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]