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) {

Reply via email to