Repository navigation
perf: Adding support for LatestBaseFilesPathFilter to Spark File Index - #18136
Conversation
95fda80 to
2dabf5d
Compare
|
@nsivabalan seems like the checks are not running on the PR. Can you please chek? |
2dabf5d to
9b8395e
Compare
|
can we consider both table types, and all query types and ensure we wire in the config only wherever applicable. |
| @@ -143,6 +146,7 @@ public abstract class BaseHoodieTableFileIndex implements AutoCloseable { | |||
| * @param configProperties unifying configuration (in the form of generic properties) | |||
| * @param queryType target query type | |||
| * @param queryPaths target DFS paths being queried | |||
There was a problem hiding this comment.
Add MOR condition.
There was a problem hiding this comment.
Added MOR table type condition.
| .fromProperties(configProperties) | ||
| .enable(configProperties.getBoolean(ENABLE.key(), DEFAULT_METADATA_ENABLE_FOR_READERS) | ||
| && HoodieTableMetadataUtil.isFilesPartitionAvailable(metaClient)) | ||
| && HoodieTableMetadataUtil.isFilesPartitionAvailable(metaClient) |
There was a problem hiding this comment.
hey @yihua : lets chat about this as well tomorrow.
|
|
||
| if (useROPathFilterForListing && !shouldIncludePendingCommits) { | ||
| // Group files by partition path, then by file group ID | ||
| Map<String, PartitionPath> partitionsMap = new HashMap<>(); |
There was a problem hiding this comment.
can we move this to a private method.
generatePartitionFileSlicesPostROTablePathFilter
| * By passing metaClient and completedTimeline, we can sync the view seen from this class against HoodieFileIndex class | ||
| */ | ||
| public HoodieROTablePathFilter(Configuration conf, | ||
| public HoodieROTablePathFilter(StorageConfiguration conf, |
There was a problem hiding this comment.
hey @yihua : can you review the changes in this patch
| " them (if possible).") | ||
|
|
||
| val FILE_INDEX_LIST_FILE_STATUSES_USING_RO_PATH_FILTER: ConfigProperty[Boolean] = | ||
| ConfigProperty.key("hoodie.datasource.read.file.index.list.file.statuses.using.ro.path.filter") |
There was a problem hiding this comment.
hoodie.datasource.read.file.index.optimize.listing.using.path.filter
There was a problem hiding this comment.
Made the change.
| properties.setProperty(DataSourceReadOptions.FILE_INDEX_LISTING_MODE_OVERRIDE.key, listingModeOverride) | ||
| } | ||
|
|
||
| var hoodieROTablePathFilterBasedFileListingEnabled = getConfigValue(options, sqlConf, |
There was a problem hiding this comment.
once we fix the config key, lets fix these vars as well
There was a problem hiding this comment.
Yes, updated.
| val result = spark.sql(s"select id, name, price, ts from $tableName order by id").collect() | ||
| // Should have deleted records where id % 3 = 0 (3, 6, 9) | ||
| // Should have doubled price for even ids (2, 4, 8, 10) | ||
| assert(result.length == 7) // 10 - 3 deleted = 7 |
There was a problem hiding this comment.
do you think below assertion would work.
we can rename one of the earlier versions of a file slice so that HoodieBaseFile parsing will fail.
so, if RO table path filter works as intended, listing files from a given partition should not fail, since we won't even try to parse the file.
but if RO table path filter did not work, it would fail.
There was a problem hiding this comment.
@suryaprasanna : did you get a chance to address this?
There was a problem hiding this comment.
HoodieROTablePathFilter also creates FSV, so renaming the older file slice will also fail the latestBaseFiles API right?
yihua
left a comment
There was a problem hiding this comment.
Nice work on adding an opt-in path filter to avoid driver OOM during file listing — the feature is well-motivated and the config is cleanly gated behind a default-off flag. The main concerns are correctness issues in the new loadFileSlicesForPartitions fast path: the partition path passed to FileSlice appears to be absolute rather than relative, and the partition map lookup can NPE if the key doesn't match. It's also worth clarifying MOR table compatibility and fixing the shared Hadoop config mutation in getPartitionPathFilter before merging.
| Map<String, PartitionPath> partitionsMap = new HashMap<>(); | ||
| partitions.forEach(p -> partitionsMap.put(p.path, p)); | ||
| Map<PartitionPath, List<FileSlice>> partitionToFileSlices = new HashMap<>(); | ||
|
|
There was a problem hiding this comment.
The partitionPathStr here is the absolute path (pathInfo.getPath().getParent().toString()), but FileSlice expects a relative partition path. The existing code path via HoodieTableFileSystemView always uses relative paths. This would cause mismatches downstream wherever FileSlice.getPartitionPath() is used. Should this be relPartitionPath instead?
There was a problem hiding this comment.
@yihua PartitionPath object stores only relative partition path str, not absolute paths. Which code path are you referring it?
| // Create FileSlice obj from StoragePathInfo. | ||
| String partitionPathStr = pathInfo.getPath().getParent().toString(); | ||
| String relPartitionPath = FSUtils.getRelativePartitionPath(basePath, pathInfo.getPath().getParent()); | ||
| HoodieBaseFile baseFile = new HoodieBaseFile(pathInfo); |
There was a problem hiding this comment.
If relPartitionPath doesn't exactly match a key in partitionsMap, partitionPathObj will be null and the computeIfAbsent call below will throw NPE. This could happen with path normalization differences (trailing slashes, scheme differences). Could you add a null check or use getRelativePartitionPath consistently with how PartitionPath.path was originally set?
| List<StoragePathInfo> allFiles = listPartitionPathFiles(partitions, activeTimeline); | ||
| log.info("On {} with query instant as {}, it took {}ms to list all files {} Hudi partitions", | ||
| metaClient.getTableConfig().getTableName(), queryInstant.map(instant -> instant).orElse("N/A"), | ||
| timer.endTimer(), partitions.size()); |
There was a problem hiding this comment.
Have you considered what happens with MOR tables here? The HoodieROTablePathFilter only returns base files (it calls fsView.getLatestBaseFiles()), so this path constructs FileSlices without log files. The !shouldIncludePendingCommits guard doesn't prevent MOR tables from reaching this code. It might be worth adding a table-type check (COW only) or documenting this limitation.
There was a problem hiding this comment.
Added table type equals COW condition.
| @@ -146,6 +147,12 @@ public List<StoragePathInfo> getAllFilesInPartition(StoragePath partitionPath) t | |||
| } | |||
|
|
|||
| @Override | |||
There was a problem hiding this comment.
It looks like the @Override annotation that belonged to getAllFilesInPartitions(Collection<String>) has been absorbed by the new method insertion. In the diff, the @Override on line 148 now applies to the new two-arg overload, while the original single-arg method (which is the actual interface abstract method) loses its @Override. Could you add @Override back to the original method?
There was a problem hiding this comment.
Yes my bad, it is a refactoring mistake, fixed it now.
| return getAllFilesInPartitions(partitions); | ||
| } | ||
|
|
||
| public Map<String, List<StoragePathInfo>> getAllFilesInPartitions(Collection<String> partitions) |
There was a problem hiding this comment.
nit: Let this call getAllFilesInPartitions(partitions, Option.empty()) to be easier to read?
There was a problem hiding this comment.
Yeah, makes sense. Refactored the code accordingly.
| Map<String, List<StoragePathInfo>> getAllFilesInPartitions(Collection<String> partitionPaths) | ||
| throws IOException; | ||
|
|
||
| default Map<String, List<StoragePathInfo>> getAllFilesInPartitions(Collection<String> partitionPaths, | ||
| Option<StoragePathFilter> pathFilterOption) | ||
| throws IOException { | ||
| return getAllFilesInPartitions(partitionPaths); | ||
| } |
There was a problem hiding this comment.
Make Map<String, List<StoragePathInfo>> getAllFilesInPartitions(Collection<String> partitionPaths) to have default implementation of getAllFilesInPartitions(partitionPaths, Option.empty()) so subclasses can avoid the repeating code? Then getAllFilesInPartitions(Collection<String> partitionPaths, Option<StoragePathFilter> pathFilterOption) becomes an abstract method.
There was a problem hiding this comment.
Good idea, made the suggested code changes.
| HoodieTimer timer = HoodieTimer.start(); | ||
| List<StoragePathInfo> allFiles = listPartitionPathFiles(partitions, activeTimeline); | ||
| log.info("On {} with query instant as {}, it took {}ms to list all files {} Hudi partitions", | ||
| metaClient.getTableConfig().getTableName(), queryInstant.map(instant -> instant).orElse("N/A"), |
There was a problem hiding this comment.
nit: queryInstant.map(instant -> instant).orElse("N/A") — the .map(instant -> instant) is a no-op. You can simplify to queryInstant.orElse("N/A").
There was a problem hiding this comment.
Made the suggested change.
There was a problem hiding this comment.
wasn't the suggestion to go w/
queryInstant.orElse("N/A")
There was a problem hiding this comment.
My bad, fixed it now, please check.
|
|
||
| public HoodieROTablePathFilter() { | ||
| this(new Configuration()); | ||
| this(HadoopFSUtils.getStorageConf()); |
There was a problem hiding this comment.
HoodieROTablePathFilter and BaseFileOnlyRelation should no longer be used based on the latest master; instead, HoodieCopyOnWriteSnapshotHadoopFsRelationFactory is used.
There was a problem hiding this comment.
+1
lets TAL at all implementations extending from HoodieBaseHadoopFsRelationFactory and we write them in
There was a problem hiding this comment.
@yihua HoodieCopyOnWriteSnapshotHadoopFsRelationFactory uses HoodieFileIndex. To optimize the file system view calls in HoodieFileIndex we are using HoodieROPathFilter. I think originally HoodieROPathFilter is used at a Relation level, now we downgraded to PathFilter level. So, with the current setup it should be fine right?
yihua
left a comment
There was a problem hiding this comment.
A better and general approach would be adding a file system view based on the latest snapshot only to limit the size of file slices in memory, which is used by the file index. That should solve the problem with better layering.
| return getAllDataFilesInPartition(storage, partitionPath, Option.empty()); | ||
| } | ||
|
|
||
| public static List<StoragePathInfo> getAllDataFilesInPartition(HoodieStorage storage, |
There was a problem hiding this comment.
we might need to change the naming, now that its not all files.
There was a problem hiding this comment.
Renamed it to getAllDataFilesInPartitionByPathFilter
nsivabalan
left a comment
There was a problem hiding this comment.
lets add tests for time travel query as well
| // Add the FileSlice to partitionToFileSlices | ||
| PartitionPath partitionPathObj = partitionsMap.get(relPartitionPath); | ||
| List<FileSlice> fileSlices = partitionToFileSlices.computeIfAbsent(partitionPathObj, k -> new ArrayList<>()); | ||
| fileSlices.add(fileSlice); |
There was a problem hiding this comment.
can we avoid this special handling.
lets route all the files into FSV.
so that we maintain one flow for all cases.
just that the the input files could have already been filtered (if path filter is applied), or could be referring to all files(if no path filter).
much simpler from maintainability standpoint.
There was a problem hiding this comment.
Creating a file system view is costly, we have actually created fsv via ropathfilter.
Adding fsv again again here means we are doing it twice so we should avoid it here.
There was a problem hiding this comment.
Yes, I see that, but building FSV here is not invoking distributed spark context and there is no metadata table involved (NoOpTableMetadata is used)). its happening in driver and we have just 1 file slice per file group. So, constructing should not add much overhead.
On the plus side, we dont' need to maintain additional custom code for, when path filter is enabled.
There was a problem hiding this comment.
I think this is a tradeoff between simplicity and performance. Also, I feel creating FSV on a bunch of files itself is very error prone.
We should skip creating another FSV that way later if any more logic added to it, we dont need to go through all that logic.
|
|
||
| public HoodieROTablePathFilter() { | ||
| this(new Configuration()); | ||
| this(HadoopFSUtils.getStorageConf()); |
There was a problem hiding this comment.
+1
lets TAL at all implementations extending from HoodieBaseHadoopFsRelationFactory and we write them in
Hey @yihua : based on latest state of the patch, I feel it nicely sits w/n HoodieTableMetadata and so, we can leverage this w/ any of FSV. |
9b8395e to
de29445
Compare
| private final boolean shouldIncludePendingCommits; | ||
| protected final boolean shouldIncludePendingCommits; | ||
| private final boolean shouldValidateInstant; | ||
| protected final boolean useROPathFilterForListing; |
There was a problem hiding this comment.
useLatestBasePathFilterForListing
There was a problem hiding this comment.
Sure, renamed the variable.
| HoodieTimer timer = HoodieTimer.start(); | ||
| List<StoragePathInfo> allFiles = listPartitionPathFiles(partitions, activeTimeline); | ||
| log.info("On {} with query instant as {}, it took {}ms to list all files {} Hudi partitions", | ||
| metaClient.getTableConfig().getTableName(), queryInstant.map(instant -> instant).orElse("N/A"), |
There was a problem hiding this comment.
wasn't the suggestion to go w/
queryInstant.orElse("N/A")
| // For MOR tables, we need log files which are not returned by HoodieROTablePathFilter | ||
| if (useROPathFilterForListing | ||
| && !shouldIncludePendingCommits | ||
| && metaClient.getTableConfig().getTableType() == HoodieTableType.COPY_ON_WRITE) { |
There was a problem hiding this comment.
for MOR, if query type is RO, we could also afford to add the filtering right?
There was a problem hiding this comment.
Sure, also adding MOR with RO as well.
| // Add the FileSlice to partitionToFileSlices | ||
| PartitionPath partitionPathObj = partitionsMap.get(relPartitionPath); | ||
| List<FileSlice> fileSlices = partitionToFileSlices.computeIfAbsent(partitionPathObj, k -> new ArrayList<>()); | ||
| fileSlices.add(fileSlice); |
There was a problem hiding this comment.
Yes, I see that, but building FSV here is not invoking distributed spark context and there is no metadata table involved (NoOpTableMetadata is used)). its happening in driver and we have just 1 file slice per file group. So, constructing should not add much overhead.
On the plus side, we dont' need to maintain additional custom code for, when path filter is enabled.
| } else { | ||
| // Also allow passing in the path filter config via Spark session conf for convenience | ||
| pathFilterOptimizedListingEnabled = getConfigValue(options, sqlConf, | ||
| "spark." + DataSourceReadOptions.FILE_INDEX_LIST_FILE_STATUSES_USING_RO_PATH_FILTER.key, null) |
There was a problem hiding this comment.
shouldn't your other patch fix this already. we should fix all config key prefix stripping at some root level and not litter across the code base.
#18205
There was a problem hiding this comment.
Yes, we also need to clean other configs as well. Let me take that as a followup. I would like to get this landed soon so the other PRs can be unblocked.
| val result = spark.sql(s"select id, name, price, ts from $tableName order by id").collect() | ||
| // Should have deleted records where id % 3 = 0 (3, 6, 9) | ||
| // Should have doubled price for even ids (2, 4, 8, 10) | ||
| assert(result.length == 7) // 10 - 3 deleted = 7 |
There was a problem hiding this comment.
@suryaprasanna : did you get a chance to address this?
| import org.apache.hudi.storage.StoragePath; | ||
| import org.apache.hudi.storage.StoragePathFilter; | ||
|
|
||
| public class HoodieROTableStoragePathFilter implements StoragePathFilter { |
There was a problem hiding this comment.
how about HoodieLatestBaseFilePathFilter
There was a problem hiding this comment.
Sure, renaming the file.
…on the driver Reviewers: O955 Project Hoodie Project Reviewer: Add blocking reviewers, pwason, jingli, meenalb, singh.sumit Reviewed By: O955 Project Hoodie Project Reviewer: Add blocking reviewers, pwason Tags: #has_java JIRA Issues: HUDI-6646 Differential Revision: https://code.uberinternal.com/D17441111 Fix build failures Fix checkstyle Refactor code Create unit tests
d34823b to
8d3be59
Compare
| val result = spark.sql(s"select id, name, price, ts from $tableName order by id").collect() | ||
| // Should have deleted records where id % 3 = 0 (3, 6, 9) | ||
| // Should have doubled price for even ids (2, 4, 8, 10) | ||
| assert(result.length == 7) // 10 - 3 deleted = 7 |
|
hey @suryaprasanna : can you fix the title and PR desc based on latest. |
|
oh. |
apache#18136) This PR adds an opt-in path filtering mechanism during file listing to prevent Spark driver OOM errors when querying large Hudi datasets with multiple file versions per partition. Problem: When file listing is performed without filtering, all file versions (including older ones) are loaded into driver memory, causing OOM on large tables. Solution: Added a new config hoodie.datasource.read.file.index.list.file.statuses.using.ro.path.filter (default: false) that enables HoodieROPathFilter during file listing to exclude older file versions. Summary and Changelog Users can now enable path filtering during file listing to avoid loading multiple file versions into memory on the driver. This is controlled by the new config hoodie.datasource.read.file.index.list.file.statuses.using.ro.path.filter. Changes: New Config: FILE_INDEX_LIST_FILE_STATUSES_USING_RO_PATH_FILTER (default: false) Enables path filtering during file listing to reduce driver memory pressure Filters out older file versions, keeping only the latest files needed for queries API Extensions: Extended HoodieTableMetadata, BaseTableMetadata, and FileSystemBackedTableMetadata to accept optional StoragePathFilter parameter Added FSUtils.getAllDataFilesInPartition overload with path filter support Created HoodieROTableStoragePathFilter wrapper to adapt Hadoop PathFilter to Hudi's StoragePathFilter interface Spark Integration: Updated BaseHoodieTableFileIndex to use path filter when enabled Modified SparkHoodieTableFileIndex to apply HoodieROTablePathFilter during partition listing
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## master #18136 +/- ##
============================================
- Coverage 57.29% 57.28% -0.02%
- Complexity 18555 18561 +6
============================================
Files 1945 1946 +1
Lines 106199 106262 +63
Branches 13128 13135 +7
============================================
+ Hits 60852 60869 +17
- Misses 39621 39662 +41
- Partials 5726 5731 +5
Flags with carried forward coverage won't be shown. Click here to find out more.
🚀 New features to boost your workflow:
|
The metadata config `BaseHoodieTableFileIndex` builds is gated on a three-way conjunction, and nothing in the repo asserted how it resolves. Both of the added conjuncts came from behavior fixes: `isFilesPartitionAvailable` from HUDI-5403 (apache#7488, a Trino listing regression) and `useLatestBaseFilesPathFilterForListing` from apache#18136. The reflective test removed earlier in this PR planted the field and asserted the getter returned it, so it discriminated nothing. Parameterize the existing `TestLocalIndex` harness, which already builds a real index over a real meta client, with the RO-path-filter flag, and assert the full truth table through it. Each row was checked against a mutant: dropping any one of the three conjuncts from the production expression makes exactly one row fail.
Describe the issue this Pull Request addresses
This PR adds an opt-in path filtering mechanism during file listing to prevent Spark driver OOM errors when querying large Hudi datasets with multiple file versions per partition.
Problem: When file listing is performed without filtering, all file versions (including older ones) are loaded into driver memory, causing OOM on large tables.
Solution: Added a new config
hoodie.datasource.read.file.index.list.file.statuses.using.ro.path.filter(default: false) that enables HoodieROPathFilter during file listing to exclude older file versions.Summary and Changelog
Users can now enable path filtering during file listing to avoid loading multiple file versions into memory on the driver. This is controlled by the new config
hoodie.datasource.read.file.index.list.file.statuses.using.ro.path.filter.Changes:
New Config:
FILE_INDEX_LIST_FILE_STATUSES_USING_RO_PATH_FILTER(default: false)API Extensions:
HoodieTableMetadata,BaseTableMetadata, andFileSystemBackedTableMetadatato accept optionalStoragePathFilterparameterFSUtils.getAllDataFilesInPartitionoverload with path filter supportHoodieROTableStoragePathFilterwrapper to adapt Hadoop PathFilter to Hudi's StoragePathFilter interfaceSpark Integration:
BaseHoodieTableFileIndexto use path filter when enabledSparkHoodieTableFileIndexto apply HoodieROTablePathFilter during partition listingImpact
Config Changes:
hoodie.datasource.read.file.index.list.file.statuses.using.ro.path.filter(default: false)Performance:
Risk Level
Low - The feature is behind a config flag (default: false) and does not change existing behavior unless explicitly enabled.
Verification:
Documentation Update
Config documentation is included in the
withDocumentationmethod of the new config property.Contributor's checklist