Repository navigation
Spark: Parallelize RemoveOrphanFiles prefix listing across executors - #16933
arifazmidd wants to merge 1 commit into
Conversation
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.
|
Thanks for working on this — the driver-side serial listing in the
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. |
|
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. |
|
@arifazmidd, do you still have cycles to work on this? Could you move this to latest Spark version (currently 4.1)? |
|
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? |
|
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:
|
|
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 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 |
|
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 |
Closes #16932
Description
DeleteOrphanFilesSparkActionhas 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 thenparallelize(matchingFiles, 1)(a single partition). For tables with tens of millions of files the driver never finishes listing, soremove_orphan_fileshangs 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
listPrefixand cluster parallelism.Changes
usePrefixListingbranch oflistedFileDS():FileSystemWalker.listDirRecursivelyWithHadoopdepth-limited discovery (SupportsPrefixOperations.listPrefixis recursive-only and cannot enumerate a single level, so discovery needs a delimiter-capable step).parallelize(subDirs).mapPartitions(...); each task runs the existingFileSystemWalker.listDirRecursivelyWithFileIOon its sub-prefix, withFileIOobtained from a broadcastSerializableTableWithSize.parallel-prefix-listingoption (defaulttrue);falserestores the prior serial driver-side path.FileSystemWalkerchanges, reuses the existing walker methods.Testing
TestRemoveOrphanFilesActionis already parameterized overusePrefixListing, so the full suite exercises the parallel path (186 tests, all passing, against currentmain).prefix_listing => truepreviously hung indefinitely on the driver inlistDirRecursivelyWithFileIO. 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
listPrefixcan trigger S3 503 throttling; raisings3.retry.num-retries/s3.retry.max-wait-msmitigates it. Could add a docs note if useful.