cshuo commented on code in PR #18372:
URL: https://github.com/apache/hudi/pull/18372#discussion_r3542103461
##########
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);
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);
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]