Repository navigation
AWS, Azure, Core, GCP, Spark: Add support for delimiter-based prefix listing - #17790
tom-s-powell wants to merge 10 commits into
Conversation
| class ADLSLocation { | ||
| private static final Pattern URI_PATTERN = Pattern.compile("^(abfss?|wasbs?)://([^/?#]+)(.*)?$"); | ||
|
|
||
| private final String scheme; |
There was a problem hiding this comment.
Raised #17594 separately to address this specific change
There was a problem hiding this comment.
Thanks, #17594 seems to be closed. Do you want to reopen it?
anuragmantri
left a comment
There was a problem hiding this comment.
Thanks for PR @tom-s-powell . It's the FileIO-only direction suggested on #16933, and it'd be great to consolidate that discussion here. cc @moomindani @arifazmidd, since this covers the discovery constraint discussed on #16933.
Two high-level suggestions before doing a deeper review:
- Decouple the
remove_orphan_fileschanges from the FileIO API. The delimiter-based listPrefix in SupportsPrefixOperations, its S3/GCS/ADLS/Hadoop implementations, and the FileSystemWalker support stand on their own. - The ADLS fully-qualified-URI fix could go first on its own
- (In a new PR) Make the Spark change on spark/v4.2 only. Recent Spark changes land in the latest version first and backport v3.5/v4.0/v4.1 in one follow-up which keeps the diff small.
- Fix the CI failures.
|
|
||
| String listPath = dir.endsWith("/") ? dir : dir + "/"; | ||
| Preconditions.checkArgument( | ||
| io.supportsPrefixListingWithDelimiter(listPath, "/"), |
There was a problem hiding this 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.
| * @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) { |
There was a problem hiding this 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?
| int parallelism = Math.min(Math.max(seedPrefixes.size(), 1), listingParallelism); | ||
| JavaRDD<String> seedPrefixRDD = sparkContext().parallelize(seedPrefixes, parallelism); | ||
| ListPrefixes listPrefixes = | ||
| new ListPrefixes( |
There was a problem hiding this 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?
| } | ||
|
|
||
| @TestTemplate | ||
| public void testPrefixListingMaxSeedDepthDiscoversOrphans() throws IOException { |
There was a problem hiding this 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.
| class ADLSLocation { | ||
| private static final Pattern URI_PATTERN = Pattern.compile("^(abfss?|wasbs?)://([^/?#]+)(.*)?$"); | ||
|
|
||
| private final String scheme; |
There was a problem hiding this comment.
Thanks, #17594 seems to be closed. Do you want to reopen it?
|
Hi @tom-s-powell. Do you still have bandwidth to work on this PR? |
|
@anuragmantri I've pushed up changes in response to the comments so far. Happy to continue iterating |
Summary
This PR extends
SupportsPrefixOperationswith delimited prefix listing. The newlistPrefix(prefix, delimiter)method returns files directly below a prefix and common sub-prefixes grouped by the delimiter. Results are returned lazily as pages so that implementations can retain the pagination behaviour of the underlying storage service.The
supportsPrefixListingWithDelimiter(prefix, delimiter)method lets aFileIOreport whether it supports a specific delimiter. The new operation has default implementations, so existingSupportsPrefixOperationsimplementations remain compatible.Delimited prefix listing is implemented for S3, GCS, Azure Data Lake Storage, and Hadoop
FileIO. Wrapper implementations, includingResolvingFileIOandEncryptingFileIO, delegate the capability to their underlyingFileIO.Motivation
When
prefix_listing=true,remove_orphan_filescurrently lists all files from the Spark driver on one thread. For tables with many files, this can create a driver bottleneck and increase driver memory use.The Hadoop listing path already limits driver-side recursive listing. It discovers directories to a configured depth and then distributes the remaining directory listings across Spark executors. This PR applies the same approach to prefix-based listing.
I've be happy to split the PRs if preferred but figured it would be useful to understand the full scope of change.
Spark behaviour
The new
prefix_listing_max_seed_depthargument controls how many levels of common sub-prefixes are discovered on the driver before listing is distributed to Spark executors.When
prefix_listing_max_seed_depth=0, which is the default, the table location is used as a single seed prefix. Listing therefore uses one Spark partition, but it runs on an executor instead of the driver.When the value is greater than zero, the driver uses delimited prefix listing to discover seed prefixes up to that depth. Spark then distributes those prefixes, subject to the configured listing parallelism, so that executors can list them concurrently. Files found during seed discovery are retained and combined with the executor-side results.
This reduces driver-side work and memory use. For tables with a suitable directory or object-key layout, it also allows listing to run in parallel.
Testing
Unit tests added for all
FileIOimplementations and SparkRemoveOrphanFilesProcedure(all versions).Alternatives considered
I considered introducing a
SupportsDirectoryListingOperationsinterface or some such that extendsSupportsPrefixOperationsand deals specifically with listing directories thatremove_orphan_filescares about. This avoided thesupportsPrefixListingWithDelimiterbut felt a bit more of an invasive change.