cshuo commented on code in PR #18372:
URL: https://github.com/apache/hudi/pull/18372#discussion_r3542079484


##########
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/utils/SparkMetadataWriterUtils.java:
##########
@@ -260,20 +261,21 @@ private static void 
setBloomFilterProps(HoodieStorageConfig storageConfig, Map<S
   /**
    * Generates expression index records
    *
-   * @param partitionFilePathAndSizeTriplet Triplet of file path, file size 
and partition name to which file belongs
-   * @param indexDefinition                 Hoodie Index Definition for the 
expression index for which records need to be generated
-   * @param metaClient                      Hoodie Table Meta Client
-   * @param parallelism                     Parallelism to use for engine 
operations
-   * @param readerSchema                    Schema of reader
-   * @param instantTime                     Instant time
-   * @param engineContext                   HoodieEngineContext
-   * @param dataWriteConfig                 Write Config for the data table
-   * @param partitionRecordsFunctionOpt     Function used to generate 
partition stat records for the EI. It takes the column range metadata generated 
for the provided partition files as input
+   * @param filesToIndex                Triplet of file path, file size and 
partition name to which file belongs

Review Comment:
   Fixed.



##########
hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metadata/index/columnstats/TestColumnStatsIndexer.java:
##########
@@ -0,0 +1,335 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.
+ */
+
+package org.apache.hudi.metadata.index.columnstats;
+
+import org.apache.hudi.avro.model.HoodieCleanMetadata;
+import org.apache.hudi.avro.model.HoodieCleanPartitionMetadata;
+import org.apache.hudi.common.config.HoodieMetadataConfig;
+import org.apache.hudi.common.data.HoodieData;
+import org.apache.hudi.common.engine.HoodieEngineContext;
+import org.apache.hudi.common.engine.HoodieLocalEngineContext;
+import org.apache.hudi.common.model.HoodieCommitMetadata;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.model.HoodieRecordMerger;
+import org.apache.hudi.common.model.HoodieWriteStat;
+import org.apache.hudi.common.schema.HoodieSchema;
+import org.apache.hudi.common.model.WriteOperationType;
+import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.view.HoodieTableFileSystemView;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.metadata.HoodieBackedTableMetadata;
+import org.apache.hudi.metadata.HoodieIndexVersion;
+import org.apache.hudi.metadata.HoodieMetadataPayload;
+import org.apache.hudi.metadata.HoodieTableMetadataUtil;
+import org.apache.hudi.metadata.MetadataPartitionType;
+import org.apache.hudi.metadata.index.model.IndexPartitionAndRecords;
+import org.apache.hudi.metadata.index.model.IndexPartitionInitialization;
+import org.apache.hudi.metadata.model.FileInfo;
+import org.apache.hudi.metadata.model.FileSliceAndPartition;
+import org.apache.hudi.stats.HoodieColumnRangeMetadata;
+import org.apache.hudi.util.Lazy;
+
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import static 
org.apache.hudi.common.testutils.HoodieTestUtils.getDefaultStorageConf;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyInt;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.when;
+
+class TestColumnStatsIndexer {
+  private static final int PARALLELISM = 4;
+  private static final int MAX_READER_BUFFER_SIZE = 1024;
+
+  private HoodieEngineContext engineContext;
+  private HoodieWriteConfig writeConfig;
+  private HoodieMetadataConfig metadataConfig;
+  private HoodieTableMetaClient metaClient;
+  private HoodieTableConfig tableConfig;
+  private HoodieRecordMerger recordMerger;
+
+  @BeforeEach
+  void setUp() {
+    engineContext = mock(HoodieEngineContext.class);
+    writeConfig = mock(HoodieWriteConfig.class);
+    metadataConfig = mock(HoodieMetadataConfig.class);
+    metaClient = mock(HoodieTableMetaClient.class);
+    tableConfig = mock(HoodieTableConfig.class);
+    recordMerger = mock(HoodieRecordMerger.class);
+
+    when(writeConfig.getMetadataConfig()).thenReturn(metadataConfig);
+    when(writeConfig.getColumnStatsIndexParallelism()).thenReturn(PARALLELISM);
+    when(writeConfig.getRecordMerger()).thenReturn(recordMerger);
+    
when(recordMerger.getRecordType()).thenReturn(HoodieRecord.HoodieRecordType.AVRO);
+    
when(metadataConfig.getColumnStatsIndexParallelism()).thenReturn(PARALLELISM);
+    
when(metadataConfig.getMaxReaderBufferSize()).thenReturn(MAX_READER_BUFFER_SIZE);
+    when(metaClient.getTableConfig()).thenReturn(tableConfig);
+  }
+
+  @Test
+  void testBuildRestoreWithEmptyInputs() {
+    ExposedColumnStatsIndexer indexer = new 
ExposedColumnStatsIndexer(engineContext, writeConfig, metaClient);
+    List<IndexPartitionAndRecords> result =
+        indexer.buildRestore("001", Collections.emptyList(), 
Collections.emptyMap(), Collections.emptyMap());
+    assertTrue(result.isEmpty());
+  }
+
+  @Test
+  void testBuildRestoreWithEmptyColumnsToIndex() {
+    try (MockedStatic<HoodieTableMetadataUtil> mockedUtil = 
mockStatic(HoodieTableMetadataUtil.class)) {
+      Map<String, Object> emptyColumnsMap = new HashMap<>();
+      mockedUtil.when(() -> HoodieTableMetadataUtil.getColumnsToIndex(
+          any(), any(), any(), eq(false), any(), 
any())).thenReturn(emptyColumnsMap);
+
+      Map<String, List<FileInfo>> filesAdded = new HashMap<>();
+      filesAdded.put("partition1", 
Collections.singletonList(FileInfo.of("file1.parquet", 1024L)));
+
+      ExposedColumnStatsIndexer indexer = new 
ExposedColumnStatsIndexer(engineContext, writeConfig, metaClient);
+      List<IndexPartitionAndRecords> result =
+          indexer.buildRestore("001", Collections.emptyList(), filesAdded, 
Collections.emptyMap());
+      assertTrue(result.isEmpty());
+    }
+  }
+
+  @Test
+  void testBuildRestoreWithValidColumns() {
+    try (MockedStatic<HoodieTableMetadataUtil> mockedUtil = 
mockStatic(HoodieTableMetadataUtil.class)) {
+      Map<String, Object> columnsMap = new HashMap<>();
+      columnsMap.put("col1", null);
+      columnsMap.put("col2", null);
+      mockedUtil.when(() -> HoodieTableMetadataUtil.getColumnsToIndex(
+          any(), any(), any(), eq(false), any(), 
any())).thenReturn(columnsMap);
+
+      HoodieData<HoodieRecord> mockHoodieData = mock(HoodieData.class);
+      mockedUtil.when(() -> 
HoodieTableMetadataUtil.convertFilesToColumnStatsRecords(
+          any(), any(), any(), any(), anyInt(), anyInt(), 
any())).thenReturn(mockHoodieData);
+
+      Map<String, List<FileInfo>> filesAdded = new HashMap<>();
+      filesAdded.put("partition1", new ArrayList<>());
+
+      ExposedColumnStatsIndexer indexer = new 
ExposedColumnStatsIndexer(engineContext, writeConfig, metaClient);
+      List<IndexPartitionAndRecords> result =
+          indexer.buildRestore("001", Collections.emptyList(), filesAdded, 
Collections.emptyMap());
+
+      assertEquals(1, result.size());
+      assertEquals(MetadataPartitionType.COLUMN_STATS.getPartitionPath(), 
result.get(0).indexPartitionName());
+      assertSame(mockHoodieData, result.get(0).indexRecords());
+
+      mockedUtil.verify(() -> 
HoodieTableMetadataUtil.convertFilesToColumnStatsRecords(
+          eq(engineContext),
+          eq(Collections.emptyMap()),
+          eq(filesAdded),
+          eq(metaClient),
+          eq(PARALLELISM),
+          eq(MAX_READER_BUFFER_SIZE),
+          any()));
+    }
+  }
+
+  @Test
+  void testBuildRestoreWithMixedInputs() {
+    try (MockedStatic<HoodieTableMetadataUtil> mockedUtil = 
mockStatic(HoodieTableMetadataUtil.class)) {
+      Map<String, Object> columnsMap = new HashMap<>();
+      columnsMap.put("col1", null);
+      columnsMap.put("col2", null);
+      columnsMap.put("col3", null);
+      mockedUtil.when(() -> HoodieTableMetadataUtil.getColumnsToIndex(
+          any(), any(), any(), eq(false), any(), 
any())).thenReturn(columnsMap);
+
+      HoodieData<HoodieRecord> mockHoodieData = mock(HoodieData.class);
+      mockedUtil.when(() -> 
HoodieTableMetadataUtil.convertFilesToColumnStatsRecords(
+          any(), any(), any(), any(), anyInt(), anyInt(), 
any())).thenReturn(mockHoodieData);
+
+      Map<String, List<FileInfo>> filesAdded = new HashMap<>();
+      List<FileInfo> filesToAdd = new ArrayList<>();
+      filesToAdd.add(FileInfo.of("file1.parquet", 1024L));
+      filesToAdd.add(FileInfo.of("file2.parquet", 2048L));
+      filesAdded.put("partition1", filesToAdd);
+
+      Map<String, List<String>> filesDeleted = new HashMap<>();
+      filesDeleted.put("partition1", List.of("old_file1.parquet", 
"old_file2.parquet"));
+
+      ExposedColumnStatsIndexer indexer = new 
ExposedColumnStatsIndexer(engineContext, writeConfig, metaClient);
+      List<IndexPartitionAndRecords> result =
+          indexer.buildRestore("001", Collections.emptyList(), filesAdded, 
filesDeleted);
+
+      assertEquals(1, result.size());
+      assertEquals(MetadataPartitionType.COLUMN_STATS.getPartitionPath(), 
result.get(0).indexPartitionName());
+      assertSame(mockHoodieData, result.get(0).indexRecords());
+
+      mockedUtil.verify(() -> 
HoodieTableMetadataUtil.convertFilesToColumnStatsRecords(
+          eq(engineContext),
+          eq(filesDeleted),
+          eq(filesAdded),
+          eq(metaClient),
+          eq(PARALLELISM),
+          eq(MAX_READER_BUFFER_SIZE),
+          any()));
+    }
+  }
+
+  @Test
+  void testInitializeDataWithEmptyInputUsesEmptyHoodieData() throws 
IOException {
+    HoodieEngineContext engineContext = mock(HoodieEngineContext.class);
+    HoodieWriteConfig writeConfig = mock(HoodieWriteConfig.class);
+    HoodieMetadataConfig metadataConfig = mock(HoodieMetadataConfig.class);
+    HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+    HoodieData<HoodieRecord> emptyData = mock(HoodieData.class);
+
+    when(writeConfig.getMetadataConfig()).thenReturn(metadataConfig);
+    when(metadataConfig.getColumnStatsIndexFileGroupCount()).thenReturn(2);
+    when(engineContext.emptyHoodieData()).thenReturn((HoodieData) emptyData);
+
+    ExposedColumnStatsIndexer indexer = new 
ExposedColumnStatsIndexer(engineContext, writeConfig, metaClient);
+    List<IndexPartitionInitialization> initializationList = 
indexer.callGetData("001", "002", Collections.emptyMap(), 
Lazy.lazily(Collections::emptyList));
+    assertEquals(1, initializationList.size());
+
+    assertEquals(MetadataPartitionType.COLUMN_STATS.getPartitionPath(), 
initializationList.get(0).indexPartitionName());
+    assertSame(emptyData, 
initializationList.get(0).dataPartitionAndRecords().get(0).indexRecords());
+  }
+
+  @SuppressWarnings("unchecked")
+  @Test
+  void testInitializeDataWithRealEngineContextAndIndexDataContent() throws 
IOException {
+    HoodieEngineContext engineContext = new 
HoodieLocalEngineContext(getDefaultStorageConf());
+    HoodieWriteConfig writeConfig = mock(HoodieWriteConfig.class);
+    HoodieMetadataConfig metadataConfig = mock(HoodieMetadataConfig.class);
+    HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+    HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
+    HoodieRecordMerger recordMerger = mock(HoodieRecordMerger.class);
+
+    when(writeConfig.getMetadataConfig()).thenReturn(metadataConfig);
+    when(writeConfig.getColumnStatsIndexParallelism()).thenReturn(4);
+    when(writeConfig.getRecordMerger()).thenReturn(recordMerger);
+    
when(recordMerger.getRecordType()).thenReturn(HoodieRecord.HoodieRecordType.AVRO);
+    when(metadataConfig.getColumnStatsIndexFileGroupCount()).thenReturn(5);
+    when(metadataConfig.getMaxReaderBufferSize()).thenReturn(4096);
+    when(metaClient.getTableConfig()).thenReturn(tableConfig);
+
+    Map<String, List<FileInfo>> files = new HashMap<>();
+    files.put("p1", Collections.singletonList(FileInfo.of("f1.parquet", 1L)));
+
+    HoodieData<HoodieRecord> records = (HoodieData<HoodieRecord>) 
(HoodieData<?>) engineContext.parallelize(
+        
Collections.singletonList(HoodieMetadataPayload.createPartitionFilesRecord("p_col",
+            Collections.singletonMap("f_col.parquet", 22L), 
Collections.emptyList())),
+        1);
+
+    try (MockedStatic<HoodieTableMetadataUtil> mockedUtil = 
mockStatic(HoodieTableMetadataUtil.class)) {
+      Map<String, Object> columns = new HashMap<>();
+      columns.put("c1", new Object());
+      mockedUtil.when(() -> HoodieTableMetadataUtil.getColumnsToIndex(any(), 
any(), any(), eq(true), eq(Option.of(HoodieRecord.HoodieRecordType.AVRO)), 
any()))
+          .thenReturn(columns);
+      mockedUtil.when(() -> 
HoodieTableMetadataUtil.convertFilesToColumnStatsRecords(any(), any(), any(), 
any(), anyInt(), anyInt(), any()))
+          .thenReturn(records);
+
+      ExposedColumnStatsIndexer indexer = new 
ExposedColumnStatsIndexer(engineContext, writeConfig, metaClient);
+      List<IndexPartitionInitialization> initializationList = 
indexer.callGetData("001", "002", files, Lazy.lazily(Collections::emptyList));
+      assertEquals(1, initializationList.size());
+
+      assertEquals(5, initializationList.get(0).totalFileGroups());
+      List<HoodieRecord> collected = 
initializationList.get(0).dataPartitionAndRecords().get(0).indexRecords().collectAsList();
+      assertEquals(1, collected.size());
+      assertEquals("p_col", collected.get(0).getRecordKey());
+    }
+  }
+
+  @Test
+  @SuppressWarnings("unchecked")
+  void testBuildUpdateWithNonEmptyCommitMetadataProducesPartitionEntry() {
+    HoodieEngineContext mockedEngineContext = mock(HoodieEngineContext.class);
+    HoodieEngineContext realEngineContext = new 
HoodieLocalEngineContext(getDefaultStorageConf());
+    ExposedColumnStatsIndexer indexer = new 
ExposedColumnStatsIndexer(mockedEngineContext, writeConfig, metaClient);
+
+    HoodieCommitMetadata commitMetadata = new HoodieCommitMetadata();
+    commitMetadata.setOperationType(WriteOperationType.UPSERT);
+    String filePath = "p1/fileid-1_1-0-1_014.parquet";
+    HoodieWriteStat writeStat = new HoodieWriteStat();
+    writeStat.setPartitionPath("p1");
+    writeStat.setPath(filePath);
+    writeStat.setRecordsStats(Collections.singletonMap("c1",
+        HoodieColumnRangeMetadata.stub(filePath, "c1", 
HoodieIndexVersion.V1)));
+    commitMetadata.getPartitionToWriteStats().put("p1", 
Collections.singletonList(writeStat));
+
+    HoodieData<HoodieWriteStat> writeStatsData = (HoodieData<HoodieWriteStat>) 
mock(HoodieData.class);
+    HoodieData<HoodieRecord> columnStatsData = (HoodieData<HoodieRecord>) 
(HoodieData<?>) realEngineContext.parallelize(
+        HoodieMetadataPayload.createColumnStatsRecords("p1",
+            Collections.singletonList(HoodieColumnRangeMetadata.stub(filePath, 
"c1", HoodieIndexVersion.V1)),
+            false).collect(java.util.stream.Collectors.toList()),
+        1);
+    when(mockedEngineContext.parallelize(any(List.class), 
anyInt())).thenReturn((HoodieData) writeStatsData);
+    when(writeStatsData.flatMap(any())).thenReturn((HoodieData) 
columnStatsData);
+
+    List<IndexPartitionAndRecords> result;
+    try (MockedStatic<HoodieTableMetadataUtil> mockedUtil = 
mockStatic(HoodieTableMetadataUtil.class)) {
+      Map<String, HoodieSchema> columnsToIndex = new HashMap<>();
+      columnsToIndex.put("c1", mock(HoodieSchema.class));
+      mockedUtil.when(() -> 
HoodieTableMetadataUtil.getColumnsToIndex(any(HoodieCommitMetadata.class), 
any(HoodieTableMetaClient.class), any(HoodieMetadataConfig.class), any()))
+          .thenReturn(columnsToIndex);
+
+      result = indexer.buildUpdate(
+          "014",
+          mock(HoodieBackedTableMetadata.class),
+          Lazy.lazily(() -> mock(HoodieTableFileSystemView.class)),
+          commitMetadata);
+    }
+
+    assertEquals(1, result.size());
+    assertEquals(MetadataPartitionType.COLUMN_STATS.getPartitionPath(), 
result.get(0).indexPartitionName());
+    List<HoodieRecord> indexRecords = 
result.get(0).indexRecords().collectAsList();
+    assertEquals(1, indexRecords.size());
+    HoodieMetadataPayload payload = (HoodieMetadataPayload) 
indexRecords.get(0).getData();
+    assertTrue(payload.getColumnStatMetadata().isPresent());
+    assertEquals("c1", payload.getColumnStatMetadata().get().getColumnName());
+    assertFalse(payload.getColumnStatMetadata().get().getIsDeleted());
+  }
+
+  @Test
+  void testBuildCleanWithNoDeletedFilesProducesEmptyRecords() {
+    Map<String, HoodieCleanPartitionMetadata> partitionMetadata = new 
HashMap<>();
+    partitionMetadata.put("p1", new HoodieCleanPartitionMetadata(
+        "p1", "KEEP_LATEST_COMMITS", Collections.emptyList(),
+        Collections.emptyList(), Collections.emptyList(), false));
+    HoodieCleanMetadata cleanMetadata = new HoodieCleanMetadata(
+        "014", 100L, 0, "013", "013", partitionMetadata, 2, 
Collections.emptyMap(), Collections.emptyMap());
+
+    HoodieEngineContext localEngineContext = new 
HoodieLocalEngineContext(getDefaultStorageConf());
+    ExposedColumnStatsIndexer indexer = new 
ExposedColumnStatsIndexer(localEngineContext, writeConfig, metaClient);
+    List<IndexPartitionAndRecords> result = indexer.buildClean("015", 
cleanMetadata);
+
+    assertEquals(1, result.size());
+    assertEquals(MetadataPartitionType.COLUMN_STATS.getPartitionPath(), 
result.get(0).indexPartitionName());
+    assertEquals(0, result.get(0).indexRecords().collectAsList().size());
+  }
+
+  private static class ExposedColumnStatsIndexer extends ColumnStatsIndexer {
+    ExposedColumnStatsIndexer(HoodieEngineContext engineContext, 
HoodieWriteConfig dataTableWriteConfig, HoodieTableMetaClient 
dataTableMetaClient) {
+      super(engineContext, dataTableWriteConfig, dataTableMetaClient);

Review Comment:
   Fixed.



##########
hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metadata/index/partitionstats/TestPartitionStatsIndexer.java:
##########
@@ -0,0 +1,223 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.
+ */
+
+package org.apache.hudi.metadata.index.partitionstats;
+
+import org.apache.hudi.avro.model.HoodieCleanMetadata;
+import org.apache.hudi.common.config.HoodieMetadataConfig;
+import org.apache.hudi.common.data.HoodieData;
+import org.apache.hudi.common.engine.HoodieEngineContext;
+import org.apache.hudi.common.engine.HoodieLocalEngineContext;
+import org.apache.hudi.common.model.HoodieCommitMetadata;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.model.HoodieRecordMerger;
+import org.apache.hudi.common.model.HoodieWriteStat;
+import org.apache.hudi.common.model.WriteOperationType;
+import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.view.HoodieTableFileSystemView;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.metadata.HoodieBackedTableMetadata;
+import org.apache.hudi.metadata.HoodieIndexVersion;
+import org.apache.hudi.metadata.HoodieMetadataPayload;
+import org.apache.hudi.metadata.HoodieTableMetadataUtil;
+import org.apache.hudi.metadata.MetadataPartitionType;
+import org.apache.hudi.metadata.index.model.IndexPartitionAndRecords;
+import org.apache.hudi.metadata.index.model.IndexPartitionInitialization;
+import org.apache.hudi.metadata.model.FileInfo;
+import org.apache.hudi.metadata.model.FileSliceAndPartition;
+import org.apache.hudi.stats.HoodieColumnRangeMetadata;
+import org.apache.hudi.util.Lazy;
+
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+import java.io.IOException;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+
+import static 
org.apache.hudi.common.testutils.HoodieTestUtils.getDefaultStorageConf;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyBoolean;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.when;
+
+class TestPartitionStatsIndexer {
+
+  @Test
+  void testSkipForNonPartitionedTable() throws IOException {
+    HoodieEngineContext engineContext = mock(HoodieEngineContext.class);
+    HoodieWriteConfig writeConfig = mock(HoodieWriteConfig.class);
+    HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+    HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
+
+    when(metaClient.getTableConfig()).thenReturn(tableConfig);
+    when(tableConfig.isTablePartitioned()).thenReturn(false);
+
+    ExposedPartitionStatsIndexer indexer = new 
ExposedPartitionStatsIndexer(engineContext, writeConfig, metaClient);
+    List<IndexPartitionInitialization> result = indexer.callGetData("001", 
"002", Collections.emptyMap(), Lazy.lazily(Collections::emptyList));
+    assertTrue(result.isEmpty());
+  }
+
+  @Test
+  void testSkipWhenColumnStatsDisabled() throws IOException {
+    HoodieEngineContext engineContext = mock(HoodieEngineContext.class);
+    HoodieWriteConfig writeConfig = mock(HoodieWriteConfig.class);
+    HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+    HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
+
+    when(metaClient.getTableConfig()).thenReturn(tableConfig);
+    when(tableConfig.isTablePartitioned()).thenReturn(true);
+    when(writeConfig.isMetadataColumnStatsIndexEnabled()).thenReturn(false);
+
+    ExposedPartitionStatsIndexer indexer = new 
ExposedPartitionStatsIndexer(engineContext, writeConfig, metaClient);
+    List<IndexPartitionInitialization> result = indexer.callGetData("001", 
"002", Collections.emptyMap(), Lazy.lazily(Collections::emptyList));
+    assertTrue(result.isEmpty());
+  }
+
+  @SuppressWarnings("unchecked")
+  @Test
+  void testInitializeWithRealEngineContextAndIndexDataContent() throws 
IOException {
+    HoodieEngineContext engineContext = new 
HoodieLocalEngineContext(getDefaultStorageConf());
+    HoodieWriteConfig writeConfig = mock(HoodieWriteConfig.class);
+    HoodieMetadataConfig metadataConfig = mock(HoodieMetadataConfig.class);
+    HoodieRecordMerger recordMerger = mock(HoodieRecordMerger.class);
+    HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+    HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
+
+    when(metaClient.getTableConfig()).thenReturn(tableConfig);
+    when(tableConfig.isTablePartitioned()).thenReturn(true);
+    when(writeConfig.isMetadataColumnStatsIndexEnabled()).thenReturn(true);
+    when(writeConfig.getMetadataConfig()).thenReturn(metadataConfig);
+    when(writeConfig.getRecordMerger()).thenReturn(recordMerger);
+    
when(recordMerger.getRecordType()).thenReturn(HoodieRecord.HoodieRecordType.AVRO);
+    when(metadataConfig.getPartitionStatsIndexFileGroupCount()).thenReturn(4);
+
+    HoodieData<HoodieRecord> records = (HoodieData<HoodieRecord>) 
(HoodieData<?>) engineContext.parallelize(
+        
Collections.singletonList(HoodieMetadataPayload.createPartitionFilesRecord("p_part",
+            Collections.singletonMap("f_part.parquet", 33L), 
Collections.emptyList())),
+        1);
+
+    try (MockedStatic<HoodieTableMetadataUtil> mockedUtil = 
mockStatic(HoodieTableMetadataUtil.class)) {
+      mockedUtil.when(() -> 
HoodieTableMetadataUtil.convertFilesToPartitionStatsRecords(any(), any(), 
any(), any(), any(), any()))
+          .thenReturn(records);
+
+      ExposedPartitionStatsIndexer indexer = new 
ExposedPartitionStatsIndexer(engineContext, writeConfig, metaClient);
+      List<IndexPartitionInitialization> initializationList = 
indexer.callGetData("001", "002", Collections.emptyMap(), 
Lazy.lazily(Collections::emptyList));
+      assertEquals(1, initializationList.size());
+
+      assertEquals(4, initializationList.get(0).totalFileGroups());
+      List<HoodieRecord> collected = 
initializationList.get(0).dataPartitionAndRecords().get(0).indexRecords().collectAsList();
+      assertEquals(1, collected.size());
+      assertEquals("p_part", collected.get(0).getRecordKey());
+    }
+  }
+
+  @Test
+  void testBuildUpdateThrowsWhenColumnStatsPartitionNotAvailable() {
+    HoodieEngineContext engineContext = mock(HoodieEngineContext.class);
+    HoodieWriteConfig writeConfig = mock(HoodieWriteConfig.class);
+    HoodieMetadataConfig metadataConfig = mock(HoodieMetadataConfig.class);
+    HoodieRecordMerger recordMerger = mock(HoodieRecordMerger.class);
+    HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+    HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
+
+    when(metaClient.getTableConfig()).thenReturn(tableConfig);
+    
when(tableConfig.isMetadataPartitionAvailable(any(MetadataPartitionType.class))).thenReturn(false);
+    when(writeConfig.getMetadataConfig()).thenReturn(metadataConfig);
+    when(writeConfig.getRecordMerger()).thenReturn(recordMerger);
+    
when(recordMerger.getRecordType()).thenReturn(HoodieRecord.HoodieRecordType.AVRO);
+
+    ExposedPartitionStatsIndexer indexer = new 
ExposedPartitionStatsIndexer(engineContext, writeConfig, metaClient);
+    assertThrows(IllegalStateException.class, () -> indexer.buildUpdate(
+        "010",
+        mock(HoodieBackedTableMetadata.class),
+        Lazy.lazily(() -> mock(HoodieTableFileSystemView.class)),
+        new HoodieCommitMetadata()));
+  }
+
+  @Test
+  @SuppressWarnings("unchecked")
+  void testBuildUpdateWithNonEmptyCommitMetadataProducesPartitionEntry() {
+    HoodieEngineContext engineContext = new 
HoodieLocalEngineContext(getDefaultStorageConf());
+    HoodieEngineContext realEngineContext = new 
HoodieLocalEngineContext(getDefaultStorageConf());

Review Comment:
   Fixed.



##########
hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestHoodieBackedMetadata.java:
##########
@@ -4089,20 +4089,20 @@ private void validateMetadata(SparkRDDWriteClient 
testClient, Option<String> ign
     List<String> metadataTablePartitions = 
FSUtils.getAllPartitionPaths(engineContext, metadataMetaClient, false);
     // Secondary index is enabled by default but no MDT partition 
corresponding to it is available
     final boolean isPartitionStatsEnabled;
-    if (!metadataWriter.getEnabledPartitionTypes().contains(COLUMN_STATS)) {
+    if (!metadataWriter.getEnabledIndexerMap().containsKey(COLUMN_STATS)) {

Review Comment:
   Fixed.



##########
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestGlobalRecordLevelIndexTableVersionSix.scala:
##########
@@ -18,15 +18,49 @@
 
 package org.apache.hudi.functional
 
-import org.apache.hudi.common.table.HoodieTableConfig
+import org.apache.hudi.DataSourceWriteOptions
+import org.apache.hudi.common.config.HoodieMetadataConfig
+import org.apache.hudi.common.model.HoodieTableType
+import org.apache.hudi.common.table.{HoodieTableConfig, HoodieTableMetaClient}
+import org.apache.hudi.common.testutils.HoodieTestDataGenerator
 import org.apache.hudi.config.HoodieWriteConfig
+import org.apache.hudi.metadata.MetadataPartitionType
 
+import org.apache.spark.sql.SaveMode
+import org.junit.jupiter.api.Assertions.assertEquals
 import org.junit.jupiter.api.Tag
+import org.junit.jupiter.params.ParameterizedTest
+import org.junit.jupiter.params.provider.EnumSource
+
+import java.util.Collections
 
 @Tag("functional-b")
 class TestGlobalRecordLevelIndexTableVersionSix extends 
TestGlobalRecordLevelIndex {
   override def commonOpts: Map[String, String] = super.commonOpts ++ Map(
     HoodieTableConfig.VERSION.key() -> "6",
     HoodieWriteConfig.WRITE_TABLE_VERSION.key() -> "6"
   )
+
+  @ParameterizedTest
+  @EnumSource(classOf[HoodieTableType])
+  override def testRLIUpsertAndDropIndex(tableType: HoodieTableType): Unit = {
+    val hudiOpts = commonOpts ++ Map(DataSourceWriteOptions.TABLE_TYPE.key -> 
tableType.name(),
+      HoodieMetadataConfig.ENABLE_METADATA_INDEX_COLUMN_STATS.key -> "true")
+    doWriteAndValidateDataAndRecordIndex(hudiOpts,
+      operation = DataSourceWriteOptions.INSERT_OPERATION_OPT_VAL,
+      saveMode = SaveMode.Overwrite)
+

Review Comment:
   Fixed.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to