Repository navigation
perf: Adding support for LatestBaseFilesPathFilter to Spark File Index #18136
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
c329058
2d7a6be
0dc9842
85ae0c4
c501757
958491d
8d3be59
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -20,13 +20,16 @@ | |
|
|
||
| import org.apache.hudi.common.config.HoodieMemoryConfig; | ||
| import org.apache.hudi.common.config.HoodieMetadataConfig; | ||
| import org.apache.hudi.common.model.HoodieBaseFile; | ||
| import org.apache.hudi.storage.StoragePathFilter; | ||
| import org.apache.hudi.common.config.TypedProperties; | ||
| import org.apache.hudi.common.engine.HoodieEngineContext; | ||
| import org.apache.hudi.common.fs.FSUtils; | ||
| import org.apache.hudi.common.model.BaseFile; | ||
| import org.apache.hudi.common.model.FileSlice; | ||
| import org.apache.hudi.common.model.HoodieLogFile; | ||
| import org.apache.hudi.common.model.HoodieTableQueryType; | ||
| import org.apache.hudi.common.model.HoodieTableType; | ||
| import org.apache.hudi.common.serialization.HoodieFileSliceSerializer; | ||
| import org.apache.hudi.common.table.HoodieTableMetaClient; | ||
| import org.apache.hudi.common.table.HoodieTableVersion; | ||
|
|
@@ -107,8 +110,9 @@ public abstract class BaseHoodieTableFileIndex implements AutoCloseable { | |
| @Getter(AccessLevel.PROTECTED) | ||
| private final List<StoragePath> queryPaths; | ||
|
|
||
| private final boolean shouldIncludePendingCommits; | ||
| protected final boolean shouldIncludePendingCommits; | ||
| private final boolean shouldValidateInstant; | ||
| protected final boolean useLatestBaseFilesPathFilterForListing; | ||
|
|
||
| // The `shouldListLazily` variable controls how we initialize/refresh the TableFileIndex: | ||
| // - non-lazy/eager listing (shouldListLazily=false): all partitions and file slices will be loaded eagerly during initialization. | ||
|
|
@@ -138,24 +142,26 @@ public abstract class BaseHoodieTableFileIndex implements AutoCloseable { | |
| private transient HoodieTableMetadata tableMetadata = null; | ||
|
|
||
| /** | ||
| * @param engineContext Hudi engine-specific context | ||
| * @param metaClient Hudi table's meta-client | ||
| * @param configProperties unifying configuration (in the form of generic properties) | ||
| * @param queryType target query type | ||
| * @param queryPaths target DFS paths being queried | ||
| * @param specifiedQueryInstant instant as of which table is being queried | ||
| * @param shouldIncludePendingCommits flags whether file-index should exclude any pending operations | ||
| * @param shouldValidateInstant flags to validate whether query instant is present in the timeline | ||
| * @param fileStatusCache transient cache of fetched [[FileStatus]]es | ||
| * @param incrementalQueryStartTime start completion time for incremental query (optional) | ||
| * @param incrementalQueryEndTime end completion time for incremental query (optional) | ||
| * @param engineContext Hudi engine-specific context | ||
| * @param metaClient Hudi table's meta-client | ||
| * @param configProperties unifying configuration (in the form of generic properties) | ||
| * @param queryType target query type | ||
| * @param queryPaths target DFS paths being queried | ||
| * @param useLatestBaseFilesPathFilterForListing memory optimization on the driver while fetching read optimized results | ||
| * @param specifiedQueryInstant instant as of which table is being queried | ||
| * @param shouldIncludePendingCommits flags whether file-index should exclude any pending operations | ||
| * @param shouldValidateInstant flags to validate whether query instant is present in the timeline | ||
| * @param fileStatusCache transient cache of fetched [[FileStatus]]es | ||
| * @param incrementalQueryStartTime start completion time for incremental query (optional) | ||
| * @param incrementalQueryEndTime end completion time for incremental query (optional) | ||
| */ | ||
| public BaseHoodieTableFileIndex(HoodieEngineContext engineContext, | ||
| HoodieTableMetaClient metaClient, | ||
| TypedProperties configProperties, | ||
| HoodieTableQueryType queryType, | ||
| List<StoragePath> queryPaths, | ||
| Option<String> specifiedQueryInstant, | ||
| boolean useLatestBaseFilesPathFilterForListing, | ||
| boolean shouldIncludePendingCommits, | ||
| boolean shouldValidateInstant, | ||
| FileStatusCache fileStatusCache, | ||
|
|
@@ -165,14 +171,17 @@ public BaseHoodieTableFileIndex(HoodieEngineContext engineContext, | |
| this.partitionColumns = metaClient.getTableConfig().getPartitionFields() | ||
| .orElseGet(() -> new String[0]); | ||
|
|
||
| // Disable metadata when ro_path_filter is enabled. | ||
| this.metadataConfig = HoodieMetadataConfig.newBuilder() | ||
| .fromProperties(configProperties) | ||
| .enable(configProperties.getBoolean(ENABLE.key(), DEFAULT_METADATA_ENABLE_FOR_READERS) | ||
| && HoodieTableMetadataUtil.isFilesPartitionAvailable(metaClient)) | ||
| && HoodieTableMetadataUtil.isFilesPartitionAvailable(metaClient) | ||
| && !useLatestBaseFilesPathFilterForListing) | ||
| .build(); | ||
|
|
||
| this.queryType = queryType; | ||
| this.queryPaths = queryPaths; | ||
| this.useLatestBaseFilesPathFilterForListing = useLatestBaseFilesPathFilterForListing; | ||
| this.specifiedQueryInstant = specifiedQueryInstant; | ||
| this.shouldIncludePendingCommits = shouldIncludePendingCommits; | ||
| this.shouldValidateInstant = shouldValidateInstant; | ||
|
|
@@ -267,14 +276,70 @@ private Map<PartitionPath, List<FileSlice>> loadFileSlicesForPartitions(List<Par | |
| validateTimestampAsOf(metaClient, specifiedQueryInstant.get()); | ||
| } | ||
|
|
||
| List<StoragePathInfo> allFiles = listPartitionPathFiles(partitions); | ||
| HoodieTimeline activeTimeline = getActiveTimeline(); | ||
| Option<HoodieInstant> latestInstant = activeTimeline.lastInstant(); | ||
| Option<String> queryInstant = specifiedQueryInstant.or(() -> latestInstant.map(HoodieInstant::requestedTime)); | ||
| validate(activeTimeline, queryInstant); | ||
|
|
||
| try (HoodieTableFileSystemView fileSystemView = new HoodieTableFileSystemView(metaClient, activeTimeline, allFiles)) { | ||
| Option<String> queryInstant = specifiedQueryInstant.or(() -> latestInstant.map(HoodieInstant::requestedTime)); | ||
| validate(activeTimeline, queryInstant); | ||
| 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.orElse("N/A"), | ||
| timer.endTimer(), partitions.size()); | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Have you considered what happens with MOR tables here? The
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Added table type equals COW condition. |
||
|
|
||
| // ROPathFilter optimization is only applicable for COW tables with snapshot queries | ||
| // For MOR tables with READ_OPTIMIZED queries, we also only need base files | ||
| if (useLatestBaseFilesPathFilterForListing | ||
| && !shouldIncludePendingCommits | ||
| && (metaClient.getTableConfig().getTableType() == HoodieTableType.COPY_ON_WRITE | ||
| || queryType == HoodieTableQueryType.READ_OPTIMIZED)) { | ||
| return generatePartitionFileSlicesPostROTablePathFilter(partitions, allFiles); | ||
| } | ||
| return filterFiles(partitions, activeTimeline, allFiles, queryInstant); | ||
| } | ||
|
|
||
| /** | ||
| * Generates FileSlices from the filtered files returned by ROPathFilter. | ||
| * This is a fast path that avoids constructing a full HoodieTableFileSystemView. | ||
| * Only applicable for COW tables since ROPathFilter only returns base files. | ||
| * | ||
| * @param partitions List of partitions to process | ||
| * @param allFiles Files already filtered by ROPathFilter | ||
| * @return Map of PartitionPath to list of FileSlices | ||
| */ | ||
| private Map<PartitionPath, List<FileSlice>> generatePartitionFileSlicesPostROTablePathFilter( | ||
| List<PartitionPath> partitions, List<StoragePathInfo> allFiles) { | ||
| // Group files by partition path, then by file group ID | ||
| Map<String, PartitionPath> partitionsMap = new HashMap<>(); | ||
| partitions.forEach(p -> partitionsMap.put(p.path, p)); | ||
| Map<PartitionPath, List<FileSlice>> partitionToFileSlices = new HashMap<>(); | ||
|
|
||
| for (StoragePathInfo pathInfo : allFiles) { | ||
| // Create FileSlice obj from StoragePathInfo. | ||
| String relPartitionPath = FSUtils.getRelativePartitionPath(basePath, pathInfo.getPath().getParent()); | ||
| HoodieBaseFile baseFile = new HoodieBaseFile(pathInfo); | ||
| // Use relative partition path for FileSlice - consistent with HoodieTableFileSystemView | ||
| FileSlice fileSlice = new FileSlice(relPartitionPath, baseFile.getCommitTime(), baseFile.getFileId()); | ||
| fileSlice.setBaseFile(baseFile); | ||
|
|
||
| // Add the FileSlice to partitionToFileSlices | ||
| PartitionPath partitionPathObj = partitionsMap.get(relPartitionPath); | ||
| if (partitionPathObj != null) { | ||
| List<FileSlice> fileSlices = partitionToFileSlices.computeIfAbsent(partitionPathObj, k -> new ArrayList<>()); | ||
| fileSlices.add(fileSlice); | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. can we avoid this special handling. 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).
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Creating a file system view is costly, we have actually created fsv via ropathfilter.
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 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.
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 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. |
||
| } else { | ||
| log.warn("Could not find partition path object for relative path: {}. Skipping file: {}", | ||
| relPartitionPath, pathInfo.getPath()); | ||
| } | ||
| } | ||
| return partitionToFileSlices; | ||
| } | ||
|
|
||
| private Map<PartitionPath, List<FileSlice>> filterFiles(List<PartitionPath> partitions, | ||
| HoodieTimeline activeTimeline, | ||
| List<StoragePathInfo> allFiles, | ||
| Option<String> queryInstant) { | ||
| try (HoodieTableFileSystemView fileSystemView = new HoodieTableFileSystemView(metaClient, activeTimeline, allFiles)) { | ||
| // NOTE: For MOR table, when the compaction is inflight, we need to not only fetch the | ||
| // latest slices, but also include the base and log files of the second-last version of | ||
| // the file slice in the same file group as the latest file slice that is under compaction. | ||
|
|
@@ -391,7 +456,8 @@ private Object[] getPartitionColumnValues(String[] partitionColumns, String part | |
| /** | ||
| * Load partition paths and it's files under the query table path. | ||
| */ | ||
| private List<StoragePathInfo> listPartitionPathFiles(List<PartitionPath> partitions) { | ||
| private List<StoragePathInfo> listPartitionPathFiles(List<PartitionPath> partitions, | ||
| HoodieTimeline activeTimeline) { | ||
| List<StoragePath> partitionPaths = partitions.stream() | ||
| // NOTE: We're using [[createPathUnsafe]] to create Hadoop's [[Path]] objects | ||
| // instances more efficiently, provided that | ||
|
|
@@ -420,7 +486,7 @@ private List<StoragePathInfo> listPartitionPathFiles(List<PartitionPath> partiti | |
|
|
||
| try { | ||
| Map<String, List<StoragePathInfo>> fetchedPartitionsMap = | ||
| tableMetadata.getAllFilesInPartitions(missingPartitionPathsMap.keySet()); | ||
| tableMetadata.getAllFilesInPartitions(missingPartitionPathsMap.keySet(), getPartitionPathFilter(activeTimeline)); | ||
|
|
||
| // Ingest newly fetched partitions into cache | ||
| fetchedPartitionsMap.forEach((absolutePath, files) -> { | ||
|
|
@@ -440,6 +506,10 @@ private List<StoragePathInfo> listPartitionPathFiles(List<PartitionPath> partiti | |
| } | ||
| } | ||
|
|
||
| protected Option<StoragePathFilter> getPartitionPathFilter(HoodieTimeline activeTimeline) { | ||
| return Option.empty(); | ||
| } | ||
|
|
||
| private void doRefresh() { | ||
| HoodieTimer timer = HoodieTimer.start(); | ||
|
|
||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -43,6 +43,7 @@ | |
| import org.apache.hudi.storage.StorageConfiguration; | ||
| import org.apache.hudi.storage.StoragePath; | ||
| import org.apache.hudi.storage.StoragePathInfo; | ||
| import org.apache.hudi.storage.StoragePathFilter; | ||
|
|
||
| import lombok.Getter; | ||
| import lombok.extern.slf4j.Slf4j; | ||
|
|
@@ -146,7 +147,8 @@ public List<StoragePathInfo> getAllFilesInPartition(StoragePath partitionPath) t | |
| } | ||
|
|
||
| @Override | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. It looks like the
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Yes my bad, it is a refactoring mistake, fixed it now. |
||
| public Map<String, List<StoragePathInfo>> getAllFilesInPartitions(Collection<String> partitions) | ||
| public Map<String, List<StoragePathInfo>> getAllFilesInPartitions(Collection<String> partitions, | ||
| Option<StoragePathFilter> unused) | ||
| throws IOException { | ||
| ValidationUtils.checkArgument(isMetadataTableInitialized); | ||
| if (partitions.isEmpty()) { | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
hey @yihua : lets chat about this as well tomorrow.