Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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,
Expand All @@ -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)

Copy link
Copy Markdown
Contributor

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.

&& !useLatestBaseFilesPathFilterForListing)
.build();

this.queryType = queryType;
this.queryPaths = queryPaths;
this.useLatestBaseFilesPathFilterForListing = useLatestBaseFilesPathFilterForListing;
this.specifiedQueryInstant = specifiedQueryInstant;
this.shouldIncludePendingCommits = shouldIncludePendingCommits;
this.shouldValidateInstant = shouldValidateInstant;
Expand Down Expand Up @@ -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());

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The 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);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The 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.
Adding fsv again again here means we are doing it twice so we should avoid it here.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The 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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The 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.
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.

} 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.
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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) -> {
Expand All @@ -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();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -481,14 +481,22 @@ public static boolean isDataFile(StoragePath path) {
public static List<StoragePathInfo> getAllDataFilesInPartition(HoodieStorage storage,
StoragePath partitionPath)
throws IOException {
return getAllDataFilesInPartitionByPathFilter(storage, partitionPath, Option.empty());
}

public static List<StoragePathInfo> getAllDataFilesInPartitionByPathFilter(HoodieStorage storage,
StoragePath partitionPath,
Option<StoragePathFilter> pathFilterOption)
throws IOException {
final Set<String> validFileExtensions = Arrays.stream(HoodieFileFormat.values())
.map(HoodieFileFormat::getFileExtension).collect(Collectors.toCollection(HashSet::new));
final String logFileExtension = HoodieFileFormat.HOODIE_LOG.getFileExtension();

try {
return storage.listDirectEntries(partitionPath, path -> {
String extension = FSUtils.getFileExtension(path.getName());
return validFileExtensions.contains(extension) || path.getName().contains(logFileExtension);
return (validFileExtensions.contains(extension) || path.getName().contains(logFileExtension))
&& pathFilterOption.map(filter -> filter.accept(path)).orElse(true);
}).stream().filter(StoragePathInfo::isFile).collect(Collectors.toList());
} catch (FileNotFoundException ex) {
// return empty FileStatus if partition does not exist already
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
import org.apache.hudi.metadata.MetadataPartitionType;
import org.apache.hudi.metadata.RawKey;
import org.apache.hudi.storage.StoragePath;
import org.apache.hudi.storage.StoragePathFilter;
import org.apache.hudi.storage.StoragePathInfo;

import java.io.IOException;
Expand Down Expand Up @@ -75,6 +76,11 @@ public Map<String, List<StoragePathInfo>> getAllFilesInPartitions(Collection<Str
throw new HoodieMetadataException("Unsupported operation: getAllFilesInPartitions!");
}

@Override
public Map<String, List<StoragePathInfo>> getAllFilesInPartitions(Collection<String> partitionPaths, Option<StoragePathFilter> pathFilterOption) throws IOException {
throw new HoodieMetadataException("Unsupported operation: getAllFilesInPartitions!");
}

@Override
public Option<BloomFilter> getBloomFilter(String partitionName, String fileName) throws HoodieMetadataException {
throw new HoodieMetadataException("Unsupported operation: getBloomFilter!");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -146,7 +147,8 @@ public List<StoragePathInfo> getAllFilesInPartition(StoragePath partitionPath) t
}

@Override

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The 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()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@
import org.apache.hudi.internal.schema.Types;
import org.apache.hudi.storage.HoodieStorage;
import org.apache.hudi.storage.StoragePath;
import org.apache.hudi.storage.StoragePathFilter;
import org.apache.hudi.storage.StoragePathInfo;

import java.io.FileNotFoundException;
Expand Down Expand Up @@ -243,7 +244,8 @@ private List<String> getPartitionPathWithPathPrefixUsingFilterExpression(String
}

@Override
public Map<String, List<StoragePathInfo>> getAllFilesInPartitions(Collection<String> partitionPaths)
public Map<String, List<StoragePathInfo>> getAllFilesInPartitions(Collection<String> partitionPaths,
Option<StoragePathFilter> pathFilterOption)
throws IOException {
if (partitionPaths == null || partitionPaths.isEmpty()) {
return Collections.emptyMap();
Expand All @@ -260,7 +262,7 @@ public Map<String, List<StoragePathInfo>> getAllFilesInPartitions(Collection<Str
partitionPathStr -> {
StoragePath partitionPath = new StoragePath(partitionPathStr);
return Pair.of(partitionPathStr,
FSUtils.getAllDataFilesInPartition(getStorage(), partitionPath));
FSUtils.getAllDataFilesInPartitionByPathFilter(getStorage(), partitionPath, pathFilterOption));
}, parallelism);
engineContext.clearJobStatus();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
import org.apache.hudi.expression.Expression;
import org.apache.hudi.internal.schema.Types;
import org.apache.hudi.storage.StoragePath;
import org.apache.hudi.storage.StoragePathFilter;
import org.apache.hudi.storage.StoragePathInfo;

import org.slf4j.Logger;
Expand Down Expand Up @@ -153,8 +154,13 @@ List<String> getPartitionPathWithPathPrefixUsingFilterExpression(List<String> re
*
* NOTE: Absolute partition paths are expected here
*/
Map<String, List<StoragePathInfo>> getAllFilesInPartitions(Collection<String> partitionPaths)
throws IOException;
default Map<String, List<StoragePathInfo>> getAllFilesInPartitions(Collection<String> partitionPaths)
throws IOException {
return getAllFilesInPartitions(partitionPaths, Option.empty());
}

Map<String, List<StoragePathInfo>> getAllFilesInPartitions(Collection<String> partitionPaths,
Option<StoragePathFilter> pathFilterOption) throws IOException;

/**
* Get the bloom filter for the FileID from the metadata table.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,7 @@ private static class TestLocalIndex extends BaseHoodieTableFileIndex {
public TestLocalIndex(HoodieEngineContext engineContext, HoodieTableMetaClient metaClient, TypedProperties configProperties, HoodieTableQueryType queryType,
List<StoragePath> queryPaths, Option<String> specifiedQueryInstant, boolean shouldIncludePendingCommits, boolean shouldValidateInstant,
FileStatusCache fileStatusCache, boolean shouldListLazily, Option<String> startCompletionTime, Option<String> endCompletionTime) {
super(engineContext, metaClient, configProperties, queryType, queryPaths, specifiedQueryInstant, shouldIncludePendingCommits, shouldValidateInstant, fileStatusCache, shouldListLazily,
super(engineContext, metaClient, configProperties, queryType, queryPaths, specifiedQueryInstant, false, shouldIncludePendingCommits, shouldValidateInstant, fileStatusCache, shouldListLazily,
startCompletionTime, endCompletionTime);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,7 @@ public HiveHoodieTableFileIndex(HoodieEngineContext engineContext,
queryType,
queryPaths,
specifiedQueryInstant,
false,
shouldIncludePendingCommits,
true,
new NoopCache(),
Expand Down
Loading
Loading