anuragmantri commented on code in PR #17790:
URL: https://github.com/apache/iceberg/pull/17790#discussion_r4161502836
##########
core/src/main/java/org/apache/iceberg/util/FileSystemWalker.java:
##########
@@ -63,19 +68,139 @@ public static void listDirRecursivelyWithFileIO(
Map<Integer, PartitionSpec> specs,
Predicate<FileInfo> filter,
Consumer<String> fileConsumer) {
+ listDirRecursivelyWithFileIO(io, dir, specs,
filter).forEachRemaining(fileConsumer);
+ }
+
+ public static Iterator<String> listDirRecursivelyWithFileIO(
+ SupportsPrefixOperations io,
+ String dir,
+ Map<Integer, PartitionSpec> specs,
+ Predicate<FileInfo> filter) {
PathFilter pathFilter = PartitionAwareHiddenPathFilter.forSpecs(specs);
String listPath = dir;
if (!dir.endsWith("/")) {
listPath = dir + "/";
}
- Iterable<FileInfo> files = io.listPrefix(listPath);
- for (FileInfo file : files) {
- Path path = new Path(file.location());
- if (!isHiddenPath(dir, path, pathFilter) && filter.test(file)) {
- fileConsumer.accept(file.location());
+ Iterator<FileInfo> files = io.listPrefix(listPath).iterator();
+ return Iterators.transform(
+ Iterators.filter(
+ files,
+ file -> {
+ Path path = new Path(file.location());
+ return !isHiddenPath(dir, path, pathFilter) && filter.test(file);
+ }),
+ FileInfo::location);
+ }
+
+ /**
+ * Recursively lists files in the specified directory that satisfy the given
conditions. Use
+ * {@link PartitionAwareHiddenPathFilter} to filter out hidden paths.
+ *
+ * <p>Provides the same depth-and-fan-out controls as {@link
#listDirRecursivelyWithHadoop}:
+ *
+ * <ul>
+ * <li>Stops traversal when the maximum recursion depth is reached and
adds the current location
+ * to the pending list via {@code directoryConsumer}.
+ * <li>Stops traversal when the number of direct sub-prefixes at a level
exceeds the threshold
+ * and adds those sub-prefixes to the pending list.
+ * </ul>
+ *
+ * @param io FileIO implementation that supports prefix listing
+ * @param dir the starting prefix to traverse
+ * @param specs partition specs used to preserve partition-name-based hidden
paths
+ * @param filter file filter; only files satisfying this condition will be
collected
+ * @param maxDepth maximum recursion depth
+ * @param maxDirectSubDirs upper limit of sub-prefixes that can be processed
directly
+ * @param directoryConsumer consumer for sub-prefixes that were not expanded
further
+ * @param fileConsumer consumer for qualifying file locations
+ */
+ public static void listDirRecursivelyWithFileIO(
+ SupportsPrefixOperations io,
+ String dir,
+ Map<Integer, PartitionSpec> specs,
+ Predicate<FileInfo> filter,
+ int maxDepth,
+ int maxDirectSubDirs,
+ Consumer<String> directoryConsumer,
+ Consumer<String> fileConsumer) {
+ PathFilter pathFilter = PartitionAwareHiddenPathFilter.forSpecs(specs);
+ listDirRecursivelyWithFileIO(
+ io,
+ dir,
+ dir,
+ pathFilter,
+ filter,
+ maxDepth,
+ maxDirectSubDirs,
+ directoryConsumer,
+ fileConsumer);
+ }
+
+ private static void listDirRecursivelyWithFileIO(
+ SupportsPrefixOperations io,
+ String baseDir,
+ String dir,
+ PathFilter pathFilter,
+ Predicate<FileInfo> filter,
+ int maxDepth,
+ int maxDirectSubDirs,
+ Consumer<String> directoryConsumer,
+ Consumer<String> fileConsumer) {
+ if (maxDepth <= 0) {
+ directoryConsumer.accept(dir);
+ return;
+ }
+
+ String listPath = dir.endsWith("/") ? dir : dir + "/";
+ Preconditions.checkArgument(
+ io.supportsPrefixListingWithDelimiter(listPath, "/"),
Review Comment:
A `FileIO` without delimiter support only fails after the action starts,
from inside the recursion, and the message prints the `FileIO` object. Would it
work to check once in `listedFileDS()` when `prefixListingMaxSeedDepth > 0`? It
could either fail with a message `prefix_listing_max_seed_depth` and the
`FileIO` class, or fall back to depth 0.
##########
api/src/main/java/org/apache/iceberg/io/SupportsPrefixOperations.java:
##########
@@ -35,6 +35,39 @@ public interface SupportsPrefixOperations extends FileIO {
*/
Iterable<FileInfo> listPrefix(String prefix);
+ /**
+ * Lists files and common prefixes under a prefix, grouped by a delimiter.
+ *
+ * <p>A file is returned in {@link PrefixListingPage#files()} when the part
of its location after
+ * {@code prefix} does not contain {@code delimiter}. When the remaining
part contains the
+ * delimiter, the file is not returned directly. Instead, {@link
PrefixListingPage#subPrefixes()}
+ * contains the common prefix through the first occurrence of the delimiter.
Common prefixes are
+ * unique, include the delimiter, and are suitable for use in a subsequent
listing operation.
+ *
+ * <p>Implementations can restrict the supported delimiters. Callers must
use {@link
+ * #supportsPrefixListingWithDelimiter(String, String)} before calling this
method.
+ *
+ * @param prefix prefix to list
+ * @param delimiter non-empty delimiter used to group matching locations
+ * @return files and common prefixes directly below the prefix
+ * @throws UnsupportedOperationException if prefix listing with the
delimiter is not supported
+ */
+ default PrefixListing listPrefix(String prefix, String delimiter) {
Review Comment:
The two `listPrefix` overloads do different things: one is recursive and
returns `Iterable<FileInfo>`, the other lists one level and returns
`PrefixListing`. Would a distinct method name be clearer? Does `PrefixListing`
need to be its own public type, or could this return
`Iterable<PrefixListingPage>`? The per-prefix probe makes sense to me for
`ResolvingFileIO` and S3 directory buckets. Could the javadoc say why support
can vary by prefix?
##########
spark/v4.1/spark/src/main/java/org/apache/iceberg/spark/actions/DeleteOrphanFilesSparkAction.java:
##########
@@ -422,17 +432,36 @@ private Dataset<String> listedFileDS() {
"Cannot use prefix listing with FileIO %s which does not support
prefix operations.",
table.io());
- Predicate<org.apache.iceberg.io.FileInfo> predicate =
- fileInfo -> fileInfo.createdAtMillis() < olderThanTimestamp;
- FileSystemWalker.listDirRecursivelyWithFileIO(
- (SupportsPrefixOperations) table.io(),
- location,
- table.specs(),
- predicate,
- matchingFiles::add);
+ List<String> seedPrefixes = Lists.newArrayList();
+ if (prefixListingMaxSeedDepth == 0) {
+ seedPrefixes.add(location);
+ } else {
+ Predicate<org.apache.iceberg.io.FileInfo> predicate =
+ fileInfo -> fileInfo.createdAtMillis() < olderThanTimestamp;
+ FileSystemWalker.listDirRecursivelyWithFileIO(
+ (SupportsPrefixOperations) table.io(),
+ location,
+ table.specs(),
+ predicate,
+ prefixListingMaxSeedDepth,
+ MAX_DRIVER_LISTING_DIRECT_SUB_DIRS,
+ seedPrefixes::add,
+ matchingFiles::add);
+ }
- JavaRDD<String> matchingFileRDD =
sparkContext().parallelize(matchingFiles, 1);
- return spark().createDataset(matchingFileRDD.rdd(), Encoders.STRING());
+ int parallelism = Math.min(Math.max(seedPrefixes.size(), 1),
listingParallelism);
+ JavaRDD<String> seedPrefixRDD = sparkContext().parallelize(seedPrefixes,
parallelism);
+ ListPrefixes listPrefixes =
+ new ListPrefixes(
Review Comment:
Other actions get the table to executors through a broadcast
`SerializableTableWithSize` (`BaseSparkAction.contentFileDS`,
`RewriteManifestsSparkAction`, `RewriteTablePathSparkAction`). Could this do
the same and read `io()` and `specs()` from the broadcast copy?
##########
azure/src/main/java/org/apache/iceberg/azure/adlsv2/ADLSLocation.java:
##########
@@ -46,6 +46,7 @@
class ADLSLocation {
private static final Pattern URI_PATTERN =
Pattern.compile("^(abfss?|wasbs?)://([^/?#]+)(.*)?$");
+ private final String scheme;
Review Comment:
Thanks, #17594 seems to be closed. Do you want to reopen it?
##########
spark/v4.1/spark/src/test/java/org/apache/iceberg/spark/actions/TestRemoveOrphanFilesAction.java:
##########
@@ -1194,6 +1194,46 @@ public void testRemoveOrphanFileActionWithDeleteMode() {
DeleteOrphanFiles.PrefixMismatchMode.DELETE);
}
+ @TestTemplate
+ public void testPrefixListingMaxSeedDepthDiscoversOrphans() throws
IOException {
Review Comment:
This would still pass if `prefixListingMaxSeedDepth` were ignored, since the
single-seed path finds the same files. Could it assert the exact count of 2,
and check that seeding actually happened? `mockStatic(FileSystemWalker.class,
CALLS_REAL_METHODS)` is already used at line 1253 of this file and would work
here. AGENTS.md also asks that new test methods drop the `test` prefix.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]