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]

Reply via email to