Skip to content

Spark: Parallelize RemoveOrphanFiles prefix listing across executors - #16933

Open
arifazmidd wants to merge 1 commit into
apache:mainfrom
arifazmidd:spark/parallel-prefix-listing
Open

arifazmidd wants to merge 1 commit into
apache:mainfrom
arifazmidd:spark/parallel-prefix-listing

Conversation

@arifazmidd

@arifazmidd arifazmidd commented Jun 22, 2026 •

Copy link
Copy Markdown
Contributor

Closes #16932

Description

DeleteOrphanFilesSparkAction has two listing strategies, but only the Hadoop one parallelizes:

  • prefix_listing = false (Hadoop / listStatus) is depth-limited on the driver and fans deep sub-directories out across executors. It parallelizes, but issues ~one LIST call per directory. On tables with hundreds of thousands of partition directories that is too many round-trips to finish in a reasonable time even when distributed.
  • prefix_listing = true (FileIO / listPrefix) uses a flat recursive listing that needs ~an order of magnitude fewer LIST calls, but it is iterated serially on the driver and then parallelize(matchingFiles, 1) (a single partition). For tables with tens of millions of files the driver never finishes listing, so remove_orphan_files hangs before reaching the deletion phase.

This change gives the prefix-listing path the same executor fan-out the Hadoop path already has, so it gets both the low call count of listPrefix and cluster parallelism.

Changes

  • Parallelize the usePrefixListing branch of listedFileDS():
    • Driver discovers shallow sub-prefixes via the existing FileSystemWalker.listDirRecursivelyWithHadoop depth-limited discovery (SupportsPrefixOperations.listPrefix is recursive-only and cannot enumerate a single level, so discovery needs a delimiter-capable step).
    • Sub-prefixes are distributed with parallelize(subDirs).mapPartitions(...); each task runs the existing FileSystemWalker.listDirRecursivelyWithFileIO on its sub-prefix, with FileIO obtained from a broadcast SerializableTableWithSize.
  • Add a parallel-prefix-listing option (default true); false restores the prior serial driver-side path.
  • No FileSystemWalker changes, reuses the existing walker methods.

Testing

  • TestRemoveOrphanFilesAction is already parameterized over usePrefixListing, so the full suite exercises the parallel path (186 tests, all passing, against current main).
  • Real-world: on an S3-backed table with ~40M files across hundreds of thousands of partitions, prefix_listing => true previously hung indefinitely on the driver in listDirRecursivelyWithFileIO. With this change the listing completed (~10k sub-prefix tasks across 30 executors, finishing in minutes) and the job proceeded to delete the orphan files.

Notes

  • Targets Spark 3.5 first; happy to replicate to 3.4 / 4.0 / 4.1 once the approach looks right.
  • Highly concurrent listPrefix can trigger S3 503 throttling; raising s3.retry.num-retries / s3.retry.max-wait-ms mitigates it. Could add a docs note if useful.

The prefix_listing path enumerated the entire table on a single driver thread
via one listPrefix iterator, so remove_orphan_files could hang on large
object-store tables before ever reaching the deletion phase. Distribute the
listing across executors the way the Hadoop path already does: discover shallow
sub-prefixes on the driver, then list each sub-prefix in parallel with
listPrefix using a broadcast SerializableTable. Gated by a new
parallel-prefix-listing option (default true), with the prior serial path as
the fallback.
@github-actions github-actions Bot added the spark label Jun 22, 2026
@moomindani

moomindani commented Jul 11, 2026 •

Copy link
Copy Markdown
Contributor

Thanks for working on this — the driver-side serial listing in the usePrefixListing path is a real bottleneck and the executor fan-out direction makes sense. A few concerns, mostly about scope and the discovery approach:

  1. This PR bundles three independent changes, and two of them aren't mentioned in the description: (a) executor-side bulk deletion, which now kicks in by default for every bulk-IO table regardless of listing mode, and (b) cache() -> persist(DISK_ONLY) in findOrphanFiles. Could you split those into separate PRs? They each deserve their own discussion, and reviewing the listing change alone would be much easier.

  2. The executor-side deletion changes the result semantics. Today, non-streaming mode returns all orphan file locations in Result.orphanFileLocations; the new path always returns at most max-orphan-file-sample-size (20k) paths and silently ignores stream-results. That's a user-visible behavior change enabled by default.

  3. I'm -1 on introducing a Hadoop FileSystem dependency into the prefix-listing path, even behind a flag. Prefix listing was added in SPARK: Remove dependency on hadoop's filesystem class from remove orphan files #12254 specifically so this action works without Hadoop. The problem isn't extra configuration: with a REST catalog using vended credentials, S3FileIO gets its credentials from the catalog, so the driver's Hadoop FileSystem cannot list the same bucket no matter how it is configured. Defaulting parallel-prefix-listing to false wouldn't fix that either — the FileIO-only users this feature targets could never use the parallel path at all. IMO discovery needs to stay on FileIO: extend SupportsPrefixOperations with a single-level (delimiter-based) listing, which S3, GCS and ADLS all support natively. That's an iceberg-api change, so it likely needs its own PR (and possibly a dev@ thread) first, with this PR rebased on top. Happy to help draft that API change.

  4. Minor: the new progress logs use LOG.warn for routine progress (existing code uses LOG.info), and the parallel-prefix-listing=false fallback path has no test coverage.

Also, new functionality usually lands in the latest Spark version first (v4.1) and is then backported — worth flipping the order once the approach settles.

@github-actions

Copy link
Copy Markdown

This pull request has been marked as stale due to 30 days of inactivity. It will be closed in 1 week if no further activity occurs. If you think that’s incorrect or this pull request requires a review, please simply write any comment. If closed, you can revive the PR at any time and @mention a reviewer or discuss it on the dev@iceberg.apache.org list. Thank you for your contributions.

@github-actions github-actions Bot added the stale label Aug 28, 2026
@anuragmantri

Copy link
Copy Markdown
Collaborator

@arifazmidd, do you still have cycles to work on this? Could you move this to latest Spark version (currently 4.1)?

@github-actions github-actions Bot removed the stale label Sep 3, 2026
@anuragmantri

Copy link
Copy Markdown
Collaborator

Gentle ping @arifazmidd. Are you still able to work on this? If not, could I take over it and address the comments and add you as co-author?

@arifazmidd

Copy link
Copy Markdown
Contributor Author

Sorry for the delay. Thanks for following up @anuragmantri and for the review @moomindani. Unfortunately, I don't have the bandwidth to carry this through at the moment so yes, please do take it over.

A few notes from the review and my testing to get this moving:

  • Point 3 is the blocker and I do agree with the review there.
  • Can probably just drop the changes referenced in points 1 and 2. They came along from internal testing and aren't needed for the listing fix. The deletion change does help throughput, so it may be worth its own PR later behind an opt-in with the semantics fixed, but it shouldn't ride along here.
  • One finding from running this: the fan-out worked on a ~40M-file table. Listing completed and the action drained the table, where the serial path had hung indefinitely. On much larger tables in the same layout (write.object-storage.enabled, 100–200M data files) I couldn't get it to finish: driver discovery needed a very large heap, and once it reached the executor listing stage the job died with repeated ExecutorLostFailure on a single task after ~20h having deleted nothing. I never root-caused that. I initially suspected that maybe the entropy hash produced prefixes of uneven density and a single dense hash-prefix subtree killed every executor assigned to it but looking at the code for it a bit I think it should be uniform by construction. So it may well have been our cluster sizing or S3 retry config rather than anything architectural. Flagging it only so you know listing has been exercised at ~40M and not successfully past that. For those tables we bypassed listing entirely with file_list_view / compareToFileList fed from a daily S3 Inventory.

@moomindani

Copy link
Copy Markdown
Contributor

Thanks for the handover notes @arifazmidd, and welcome @anuragmantri.

On point 3, the part I care about is the constraint, not the remedy: discovery currently needs a Hadoop FileSystem, and a REST catalog with vended credentials has none configured for the table's bucket, so that path cannot work there. How to get around it is open. Extending SupportsPrefixOperations with a delimiter-based single-level listing was one suggestion of mine — main still has only listPrefix and deletePrefix — but I am not attached to it, and if you see a better route, take it.

The split is yours to decide as well: take the whole thing, take part of it and leave the rest as a follow-up, or keep this PR to the Spark side and let any api change land separately. If some piece would help coming from me, say so and I will pick it up; otherwise I will review.

Also worth carrying into the description: @arifazmidd's ~40M-file result, that nothing larger has completed, and that the S3 Inventory route via file_list_view / compareToFileList is the escape hatch at that size.

@anuragmantri

Copy link
Copy Markdown
Collaborator

While reiimplementing this, I found #17790, which addresses the same problem and also resolves the discovery constraint discussed above. It discovers sub-prefixes through a new delimiter-based listPrefix on SupportsPrefixOperations instead of Hadoop's FileSystem, so it works with FileIO-only setups like REST catalogs that vend credentials. I'd suggest we consolidate on #17790 and close this one.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Spark: RemoveOrphanFiles prefix_listing enumerates the whole table serially on the driver

3 participants