This is an automated email from the ASF dual-hosted git repository.
danny0405 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git
The following commit(s) were added to refs/heads/master by this push:
new 91510d705b60 perf(flink): scope column stats reads to candidate
partitions (#19992)
91510d705b60 is described below
commit 91510d705b60a83321e54602953512f6bbbd13d3
Author: Danny Chan <[email protected]>
AuthorDate: Tue Sep 22 10:11:20 2026 +0800
perf(flink): scope column stats reads to candidate partitions (#19992)
* perf(flink): scope column stats reads to candidate partitions
---
.../java/org/apache/hudi/source/FileIndex.java | 13 ++-
.../apache/hudi/source/stats/ColumnStatsIndex.java | 7 +-
.../apache/hudi/source/stats/FileStatsIndex.java | 27 +++--
.../hudi/source/stats/PartitionStatsIndex.java | 5 +-
.../java/org/apache/hudi/source/TestFileIndex.java | 127 +++++++++++++++++++++
.../hudi/source/stats/TestColumnStatsIndex.java | 53 ++++++++-
6 files changed, 213 insertions(+), 19 deletions(-)
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/FileIndex.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/FileIndex.java
index 23aa83b3ba40..1248e5078969 100644
---
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/FileIndex.java
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/FileIndex.java
@@ -72,6 +72,7 @@ public class FileIndex implements Serializable, AutoCloseable
{
private final ColumnStatsProbe colStatsProbe; // for
probing column stats
private final Function<String, Integer> partitionBucketIdFunc; // for
bucket pruning
private List<String> partitionPaths; // cache of
partition paths
+ private int totalPartitionCount; // partition
count before pruning
private final FileStatsIndex fileStatsIndex; // for data
skipping
private final Option<BaseRecordLevelIndex> recordLevelIndex;
private final HoodieTableMetaClient metaClient;
@@ -197,8 +198,16 @@ public class FileIndex implements Serializable,
AutoCloseable {
}
// data skipping based on column stats
+ if (colStatsProbe == null || filteredFileSlices.isEmpty()) {
+ return filteredFileSlices;
+ }
+ List<String> candidatePartitions =
filteredFileSlices.stream().map(FileSlice::getPartitionPath).distinct().collect(Collectors.toList());
+ // Avoid expanding column prefixes when the remaining files still span
every partition.
+ if (candidatePartitions.size() == totalPartitionCount) {
+ candidatePartitions = Collections.emptyList();
+ }
List<String> allFiles =
filteredFileSlices.stream().map(FileSlice::getAllFileNames).flatMap(List::stream).collect(Collectors.toList());
- Set<String> candidateFiles =
fileStatsIndex.computeCandidateFiles(colStatsProbe, allFiles);
+ Set<String> candidateFiles =
fileStatsIndex.computeCandidateFiles(colStatsProbe, allFiles,
candidatePartitions);
if (candidateFiles == null) {
// no need to filter by col stats or error occurs.
return filteredFileSlices;
@@ -232,6 +241,7 @@ public class FileIndex implements Serializable,
AutoCloseable {
@VisibleForTesting
public void reset() {
this.partitionPaths = null;
+ this.totalPartitionCount = 0;
}
// -------------------------------------------------------------------------
@@ -249,6 +259,7 @@ public class FileIndex implements Serializable,
AutoCloseable {
}
List<String> allPartitionPaths = this.tableExists ?
FSUtils.getAllPartitionPaths(new HoodieFlinkEngineContext(hadoopConf),
metaClient, metadataConfig)
: Collections.emptyList();
+ this.totalPartitionCount = allPartitionPaths.size();
this.partitionPaths = partitionPruner.map(pruner ->
pruner.filter(allPartitionPaths).stream().collect(Collectors.toList())).orElse(allPartitionPaths);
return this.partitionPaths;
}
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/stats/ColumnStatsIndex.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/stats/ColumnStatsIndex.java
index d54df6dc13ff..1c388275373c 100644
---
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/stats/ColumnStatsIndex.java
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/stats/ColumnStatsIndex.java
@@ -31,12 +31,13 @@ public interface ColumnStatsIndex extends
FlinkMetadataIndex {
/**
* Computes the filtered files with given candidates.
*
- * @param columnStatsProbe The utility to filter the column stats metadata.
- * @param allFile The file name list of the candidate files.
+ * @param columnStatsProbe The utility to filter the column stats
metadata.
+ * @param allFile The file name list of the candidate files.
+ * @param candidatePartitions The relative partition paths of the candidate
files, or an empty list to read all partitions.
*
* @return The set of filtered file names
*/
- Set<String> computeCandidateFiles(ColumnStatsProbe columnStatsProbe,
List<String> allFile);
+ Set<String> computeCandidateFiles(ColumnStatsProbe columnStatsProbe,
List<String> allFile, List<String> candidatePartitions);
/**
* Computes the filtered partition paths with given candidates.
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/stats/FileStatsIndex.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/stats/FileStatsIndex.java
index 58f6b74dcb67..99c22402a02f 100644
---
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/stats/FileStatsIndex.java
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/stats/FileStatsIndex.java
@@ -128,13 +128,13 @@ public class FileStatsIndex implements ColumnStatsIndex {
}
@Override
- public Set<String> computeCandidateFiles(ColumnStatsProbe probe,
List<String> allFiles) {
+ public Set<String> computeCandidateFiles(ColumnStatsProbe probe,
List<String> allFiles, List<String> candidatePartitions) {
if (probe == null || !isIndexAvailable()) {
return null;
}
try {
String[] targetColumns = probe.getReferencedCols();
- final List<RowData> statsRows =
readColumnStatsIndexByColumns(targetColumns);
+ final List<RowData> statsRows =
readColumnStatsIndexByColumns(targetColumns, candidatePartitions);
return candidatesInMetadataTable(probe, statsRows, allFiles);
} catch (Throwable t) {
log.error("Failed to read metadata index: {} for data skipping",
getIndexPartitionName(), t);
@@ -384,20 +384,27 @@ public class FileStatsIndex implements ColumnStatsIndex {
return converter.convert(rawVal);
}
+ /**
+ * Reads statistics for the requested columns and relative partition paths.
+ * Column prefixes limit reads to the referenced columns; partition prefixes
further restrict reads after partition pruning.
+ * An empty partition list reads all partitions using column-only prefixes.
+ */
@VisibleForTesting
- public List<RowData> readColumnStatsIndexByColumns(String[] targetColumns) {
- // NOTE: If specific columns have been provided, we can considerably trim
down amount of data fetched
- // by only fetching Column Stats Index records pertaining to the
requested columns.
- // Otherwise, we fall back to read whole Column Stats Index
+ public List<RowData> readColumnStatsIndexByColumns(String[] targetColumns,
List<String> candidatePartitions) {
ValidationUtils.checkArgument(targetColumns.length > 0,
"Column stats is only valid when push down filters have referenced
columns");
// Read Metadata Table's column stats Flink's RowData list by
- // - Fetching the records by key-prefixes (column names)
+ // - Fetching the records by key-prefixes (column names and candidate
partitions, when provided)
// - Deserializing fetched records into [[RowData]]s
- List<ColumnStatsIndexPrefixRawKey> rawKeys = Arrays.stream(targetColumns)
- .map(ColumnStatsIndexPrefixRawKey::new) // Just column name, no
partition
- .collect(Collectors.toList());
+ List<ColumnStatsIndexPrefixRawKey> rawKeys;
+ if (candidatePartitions.isEmpty()) {
+ rawKeys =
Arrays.stream(targetColumns).map(ColumnStatsIndexPrefixRawKey::new).collect(Collectors.toList());
+ } else {
+ rawKeys = candidatePartitions.stream().distinct()
+ .flatMap(partition -> Arrays.stream(targetColumns).map(column -> new
ColumnStatsIndexPrefixRawKey(column, partition)))
+ .collect(Collectors.toList());
+ }
HoodieData<HoodieRecord<HoodieMetadataPayload>> records =
getMetadataTable().getRecordsByKeyPrefixes(
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/stats/PartitionStatsIndex.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/stats/PartitionStatsIndex.java
index 4ab47c7ef28f..afc5e1d55bc9 100644
---
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/stats/PartitionStatsIndex.java
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/stats/PartitionStatsIndex.java
@@ -28,6 +28,7 @@ import org.apache.flink.table.types.logical.RowType;
import javax.annotation.Nullable;
+import java.util.Collections;
import java.util.List;
import java.util.Set;
@@ -57,7 +58,7 @@ public class PartitionStatsIndex extends FileStatsIndex {
}
@Override
- public Set<String> computeCandidateFiles(ColumnStatsProbe probe,
List<String> allFiles) {
+ public Set<String> computeCandidateFiles(ColumnStatsProbe probe,
List<String> allFiles, List<String> candidatePartitions) {
throw new UnsupportedOperationException("This method is not supported by "
+ this.getClass().getSimpleName());
}
@@ -82,6 +83,6 @@ public class PartitionStatsIndex extends FileStatsIndex {
*/
@Override
public Set<String> computeCandidatePartitions(ColumnStatsProbe probe,
List<String> allPartitions) {
- return super.computeCandidateFiles(probe, allPartitions);
+ return super.computeCandidateFiles(probe, allPartitions,
Collections.emptyList());
}
}
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestFileIndex.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestFileIndex.java
index d421f6e72098..3ad096a96870 100644
---
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestFileIndex.java
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestFileIndex.java
@@ -19,6 +19,7 @@
package org.apache.hudi.source;
import org.apache.hudi.common.config.HoodieMetadataConfig;
+import org.apache.hudi.common.fs.FSUtils;
import org.apache.hudi.common.model.FileSlice;
import org.apache.hudi.common.model.HoodieTableType;
import org.apache.hudi.common.table.HoodieTableMetaClient;
@@ -28,6 +29,7 @@ import org.apache.hudi.configuration.FlinkOptions;
import org.apache.hudi.keygen.NonpartitionedAvroKeyGenerator;
import org.apache.hudi.source.prune.ColumnStatsProbe;
import org.apache.hudi.source.prune.PartitionPruners;
+import org.apache.hudi.source.stats.FileStatsIndex;
import org.apache.hudi.storage.StoragePath;
import org.apache.hudi.storage.StoragePathInfo;
import org.apache.hudi.util.StreamerUtil;
@@ -53,6 +55,8 @@ import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.EnumSource;
import org.junit.jupiter.params.provider.MethodSource;
import org.junit.jupiter.params.provider.ValueSource;
+import org.mockito.MockedConstruction;
+import org.mockito.MockedStatic;
import java.io.File;
import java.math.BigDecimal;
@@ -64,6 +68,7 @@ import java.time.ZoneId;
import java.util.Arrays;
import java.util.Collection;
import java.util.Collections;
+import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
@@ -81,6 +86,16 @@ import static org.hamcrest.CoreMatchers.is;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyList;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
/**
* Test cases for {@link FileIndex}.
@@ -147,6 +162,18 @@ public class TestFileIndex {
assertThat(partitions.size(), is(0));
}
+ @Test
+ void testFilterFileSlicesWithoutColumnStatsProbe() {
+ Configuration conf =
TestConfigurations.getDefaultConf(tempFile.getAbsolutePath());
+ FileIndex fileIndex = FileIndex.builder().path(new
StoragePath(tempFile.getAbsolutePath())).conf(conf)
+ .rowType(TestConfigurations.ROW_TYPE).build();
+ FileSlice fileSlice = mock(FileSlice.class);
+ List<FileSlice> fileSlices = Collections.singletonList(fileSlice);
+
+ assertEquals(fileSlices, fileIndex.filterFileSlices(fileSlices));
+ verify(fileSlice, never()).getAllFileNames();
+ }
+
@Test
void testFileListingWithDataSkipping() throws Exception {
Configuration conf =
TestConfigurations.getDefaultConf(tempFile.getAbsolutePath(),
TestConfigurations.ROW_DATA_TYPE_BIGINT);
@@ -179,6 +206,106 @@ public class TestFileIndex {
assertThat(fileSlices.size(), is(2));
}
+ @ParameterizedTest
+ @MethodSource("columnStatsPartitionScopes")
+ void testColumnStatsPartitionScope(List<String> allPartitions, List<String>
selectedPartitions,
+ List<String> filePartitions, List<String>
expectedScope) throws Exception {
+ Configuration conf =
TestConfigurations.getDefaultConf(tempFile.getAbsolutePath());
+ conf.set(METADATA_ENABLED, true);
+ conf.set(READ_DATA_SKIPPING_ENABLED, true);
+ HoodieTableMetaClient metaClient = StreamerUtil.initTableIfNotExists(conf);
+ ColumnStatsProbe probe = mock(ColumnStatsProbe.class);
+
+ try (MockedStatic<FSUtils> fsUtils = mockStatic(FSUtils.class);
+ MockedConstruction<FileStatsIndex> statsIndexes =
mockConstruction(FileStatsIndex.class, (index, context) ->
+ when(index.computeCandidateFiles(any(), anyList(),
anyList())).thenReturn(null))) {
+ fsUtils.when(() -> FSUtils.getAllPartitionPaths(any(), eq(metaClient),
any(HoodieMetadataConfig.class))).thenReturn(allPartitions);
+ try (FileIndex fileIndex = FileIndex.builder().path(new
StoragePath(tempFile.getAbsolutePath())).conf(conf)
+
.rowType(TestConfigurations.ROW_TYPE).metaClient(metaClient).columnStatsProbe(probe)
+ .partitionPruner(partitions -> new
HashSet<>(selectedPartitions)).build()) {
+ assertEquals(new HashSet<>(selectedPartitions), new
HashSet<>(fileIndex.getOrBuildPartitionPaths()));
+ List<FileSlice> slices = filePartitions.stream().map(partition -> new
FileSlice(partition, "001", "file1")).collect(Collectors.toList());
+ assertEquals(slices, fileIndex.filterFileSlices(slices));
+
verify(statsIndexes.constructed().get(0)).computeCandidateFiles(eq(probe),
anyList(), eq(expectedScope));
+ // Choosing prefixes must reuse the partition listing, not request
another one.
+ fsUtils.verify(() -> FSUtils.getAllPartitionPaths(any(),
eq(metaClient), any(HoodieMetadataConfig.class)), times(1));
+ }
+ }
+ }
+
+ private static Stream<Arguments> columnStatsPartitionScopes() {
+ List<String> allPartitions = Arrays.asList("par1", "par2");
+ List<String> onePartition = Collections.singletonList("par1");
+ List<String> nonPartitioned = Collections.singletonList("");
+ return Stream.of(
+ // A partition pruner that retains every partition must still use
column-only prefixes.
+ Arguments.of(allPartitions, allPartitions, Arrays.asList("par1",
"par1", "par2"), Collections.emptyList()),
+ Arguments.of(allPartitions, onePartition, onePartition, onePartition),
+ // Incremental reads may contain files from fewer partitions than the
partition listing.
+ Arguments.of(allPartitions, allPartitions, onePartition, onePartition),
+ Arguments.of(nonPartitioned, nonPartitioned, nonPartitioned,
Collections.emptyList()));
+ }
+
+ @ParameterizedTest
+ @ValueSource(booleans = {true, false})
+ void testColumnStatsPartitionScopeAfterBucketPruning(boolean
removesPartition) throws Exception {
+ Configuration conf =
TestConfigurations.getDefaultConf(tempFile.getAbsolutePath());
+ conf.set(METADATA_ENABLED, true);
+ conf.set(READ_DATA_SKIPPING_ENABLED, true);
+ HoodieTableMetaClient metaClient = StreamerUtil.initTableIfNotExists(conf);
+ ColumnStatsProbe probe = mock(ColumnStatsProbe.class);
+
+ try (MockedStatic<FSUtils> fsUtils = mockStatic(FSUtils.class);
+ MockedConstruction<FileStatsIndex> statsIndexes =
mockConstruction(FileStatsIndex.class, (index, context) ->
+ when(index.computeCandidateFiles(any(), anyList(),
anyList())).thenReturn(null))) {
+ fsUtils.when(() -> FSUtils.getAllPartitionPaths(any(), eq(metaClient),
any(HoodieMetadataConfig.class)))
+ .thenReturn(Arrays.asList("par1", "par2"));
+ try (FileIndex fileIndex = FileIndex.builder().path(new
StoragePath(tempFile.getAbsolutePath())).conf(conf)
+
.rowType(TestConfigurations.ROW_TYPE).metaClient(metaClient).columnStatsProbe(probe)
+ .partitionBucketIdFunc(partition -> removesPartition &&
partition.equals("par2") ? 2 : 0).build()) {
+ fileIndex.getOrBuildPartitionPaths();
+ List<FileSlice> slices = Arrays.asList(
+ new FileSlice("par1", "001", "00000000-file1"), new
FileSlice("par1", "001", "00000001-file2"),
+ new FileSlice("par2", "001", "00000000-file3"), new
FileSlice("par2", "001", "00000001-file4"));
+ assertEquals(removesPartition ? 1 : 2,
fileIndex.filterFileSlices(slices).size());
+
verify(statsIndexes.constructed().get(0)).computeCandidateFiles(eq(probe),
anyList(),
+ eq(removesPartition ? Collections.singletonList("par1") :
Collections.emptyList()));
+ }
+ }
+ }
+
+ @Test
+ void testColumnStatsPartitionScopeAfterReset() throws Exception {
+ Configuration conf =
TestConfigurations.getDefaultConf(tempFile.getAbsolutePath());
+ conf.set(METADATA_ENABLED, true);
+ conf.set(READ_DATA_SKIPPING_ENABLED, true);
+ HoodieTableMetaClient metaClient = StreamerUtil.initTableIfNotExists(conf);
+ ColumnStatsProbe probe = mock(ColumnStatsProbe.class);
+ List<String> onePartition = Collections.singletonList("par1");
+ List<FileSlice> slices = Collections.singletonList(new FileSlice("par1",
"001", "file1"));
+
+ try (MockedStatic<FSUtils> fsUtils = mockStatic(FSUtils.class);
+ MockedConstruction<FileStatsIndex> statsIndexes =
mockConstruction(FileStatsIndex.class, (index, context) ->
+ when(index.computeCandidateFiles(any(), anyList(),
anyList())).thenReturn(null))) {
+ fsUtils.when(() -> FSUtils.getAllPartitionPaths(any(), eq(metaClient),
any(HoodieMetadataConfig.class)))
+ .thenReturn(onePartition, Arrays.asList("par1", "par2"));
+ try (FileIndex fileIndex = FileIndex.builder().path(new
StoragePath(tempFile.getAbsolutePath())).conf(conf)
+
.rowType(TestConfigurations.ROW_TYPE).metaClient(metaClient).columnStatsProbe(probe).build())
{
+ fileIndex.getOrBuildPartitionPaths();
+ fileIndex.filterFileSlices(slices);
+
verify(statsIndexes.constructed().get(0)).computeCandidateFiles(eq(probe),
anyList(), eq(Collections.emptyList()));
+
+ fileIndex.reset();
+ // Until the next listing, do not assume the cached partition count is
still valid.
+ fileIndex.filterFileSlices(slices);
+ fileIndex.getOrBuildPartitionPaths();
+ fileIndex.filterFileSlices(slices);
+ verify(statsIndexes.constructed().get(0),
times(2)).computeCandidateFiles(eq(probe), anyList(), eq(onePartition));
+ fsUtils.verify(() -> FSUtils.getAllPartitionPaths(any(),
eq(metaClient), any(HoodieMetadataConfig.class)), times(2));
+ }
+ }
+ }
+
@ParameterizedTest
@EnumSource(value = HoodieTableType.class)
void testFileListingWithPartitionStatsPruning(HoodieTableType tableType)
throws Exception {
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/stats/TestColumnStatsIndex.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/stats/TestColumnStatsIndex.java
index 898b72b06f08..02c6f9cea77b 100644
---
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/stats/TestColumnStatsIndex.java
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/stats/TestColumnStatsIndex.java
@@ -21,6 +21,8 @@ package org.apache.hudi.source.stats;
import org.apache.hudi.common.config.HoodieMetadataConfig;
import org.apache.hudi.common.util.collection.Pair;
import org.apache.hudi.configuration.FlinkOptions;
+import org.apache.hudi.keygen.NonpartitionedAvroKeyGenerator;
+import org.apache.hudi.source.prune.ColumnStatsProbe;
import org.apache.hudi.util.StreamerUtil;
import org.apache.hudi.utils.TestConfigurations;
import org.apache.hudi.utils.TestData;
@@ -30,9 +32,12 @@ import org.apache.flink.table.data.GenericRowData;
import org.apache.flink.table.data.RowData;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
import java.io.File;
import java.util.Arrays;
+import java.util.Collections;
import java.util.Comparator;
import java.util.List;
import java.util.stream.Collectors;
@@ -42,6 +47,9 @@ import static org.hamcrest.MatcherAssert.assertThat;
import static org.junit.jupiter.api.Assertions.assertArrayEquals;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
/**
* Test cases for {@link ColumnStatsIndex}.
@@ -61,7 +69,7 @@ public class TestColumnStatsIndex {
String[] queryColumns = {"uuid", "age"};
PartitionStatsIndex indexSupport = new PartitionStatsIndex(path,
TestConfigurations.ROW_TYPE, conf, StreamerUtil.createMetaClient(conf));
- List<RowData> indexRows =
indexSupport.readColumnStatsIndexByColumns(queryColumns);
+ List<RowData> indexRows =
indexSupport.readColumnStatsIndexByColumns(queryColumns,
Collections.emptyList());
List<String> results =
indexRows.stream().map(Object::toString).sorted(String::compareTo).collect(Collectors.toList());
List<String> expected = Arrays.asList(
"+I(par1,+I(23),+I(33),0,2,age)",
@@ -99,7 +107,8 @@ public class TestColumnStatsIndex {
// explicit query columns
String[] queryColumns1 = {"uuid", "age"};
FileStatsIndex indexSupport = new FileStatsIndex(path,
TestConfigurations.ROW_TYPE, conf, StreamerUtil.createMetaClient(conf));
- List<RowData> indexRows1 =
indexSupport.readColumnStatsIndexByColumns(queryColumns1);
+ // An empty partition list reads column statistics across all partitions.
+ List<RowData> indexRows1 =
indexSupport.readColumnStatsIndexByColumns(queryColumns1,
Collections.emptyList());
Pair<List<RowData>, String[]> transposedIndexTable1 =
indexSupport.transposeColumnStatsIndex(indexRows1, queryColumns1);
assertThat("The schema columns should sort by natural order",
Arrays.toString(transposedIndexTable1.getRight()), is("[age, uuid]"));
@@ -112,9 +121,47 @@ public class TestColumnStatsIndex {
+ "+I(2,44,56,0,id7,id8,0)]";
assertThat(transposed1.toString(), is(expected));
+ List<RowData> scopedRows =
indexSupport.readColumnStatsIndexByColumns(queryColumns1, Arrays.asList("par1",
"par3"));
+ assertEquals(4, scopedRows.size());
+ assertEquals("[+I(2,18,20,0,id5,id6,0), +I(2,23,33,0,id1,id2,0)]",
+ filterOutFileNames(indexSupport.transposeColumnStatsIndex(scopedRows,
queryColumns1).getLeft()).toString());
+
// no query columns, only for tests
assertThrows(IllegalArgumentException.class,
- () -> indexSupport.readColumnStatsIndexByColumns(new String[0]));
+ () -> indexSupport.readColumnStatsIndexByColumns(new String[0],
Collections.emptyList()));
+ }
+
+ @ParameterizedTest
+ @ValueSource(strings = {"par1", "partition=par1", ""})
+ void testReadColumnStatsForCandidatePartitions(String partition) throws
Exception {
+ final String path = tempFile.getAbsolutePath();
+ Configuration conf = TestConfigurations.getDefaultConf(path);
+ conf.set(FlinkOptions.METADATA_ENABLED, true);
+
conf.setString(HoodieMetadataConfig.ENABLE_METADATA_INDEX_COLUMN_STATS.key(),
"true");
+ conf.set(FlinkOptions.HIVE_STYLE_PARTITIONING,
partition.startsWith("partition="));
+ if (partition.isEmpty()) {
+ conf.set(FlinkOptions.PARTITION_PATH_FIELD, "");
+ conf.set(FlinkOptions.KEYGEN_CLASS_NAME,
NonpartitionedAvroKeyGenerator.class.getName());
+ }
+ TestData.writeData(TestData.DATA_SET_INSERT, conf);
+
+ String[] columns = {"uuid", "age"};
+ try (FileStatsIndex index = new FileStatsIndex(path,
TestConfigurations.ROW_TYPE, conf, StreamerUtil.createMetaClient(conf))) {
+ // Repeated candidate partitions must not duplicate statistics during
transposition.
+ List<RowData> rows = index.readColumnStatsIndexByColumns(columns,
Arrays.asList(partition, partition));
+ assertEquals(2, rows.size());
+ List<RowData> transposed =
filterOutFileNames(index.transposeColumnStatsIndex(rows, columns).getLeft());
+ assertEquals(partition.isEmpty() ? "[+I(8,18,56,0,id1,id8,0)]" :
"[+I(2,23,33,0,id1,id2,0)]", transposed.toString());
+
+ // Reject all indexed files, but retain candidates with no statistics.
+ ColumnStatsProbe probe = mock(ColumnStatsProbe.class);
+ when(probe.getReferencedCols()).thenReturn(columns);
+ List<String> files = rows.stream().map(row ->
row.getString(0).toString()).distinct().collect(Collectors.toList());
+ files.add("missing-stats.parquet");
+ assertEquals(Collections.singleton("missing-stats.parquet"),
+ index.computeCandidateFiles(probe, files,
Collections.singletonList(partition)));
+ assertTrue(index.readColumnStatsIndexByColumns(columns,
Collections.singletonList("unknown-partition")).isEmpty());
+ }
}
private static List<RowData> filterOutFileNames(List<RowData> indexRows) {