This is an automated email from the ASF dual-hosted git repository. voonhous pushed a commit to branch dedupe-keybased-buffer-tests in repository https://gitbox.apache.org/repos/asf/hudi.git
commit 974c55a2d927c67f6be299d4bf9c75bbaabe36a8 Author: voon <[email protected]> AuthorDate: Fri Sep 18 14:20:49 2026 +0800 test(common): drop duplicated key-based buffer tests TestKeyBasedFileGroupRecordBuffer carried four *WithRecords tests that duplicate TestStreamingKeyBasedFileGroupRecordBuffer: readWithEventTimeOrderingWithRecords -> readWithEventTimeOrdering readWithCommitTimeOrderingWithRecords -> readWithCommitTimeOrdering readWithCustomPayloadWithRecords -> readWithCustomPayload readWithCustomMergerWithRecords -> readWithCustomMergerWithRecords Each pair is identical apart from line wrapping and one inner-class name being qualified on the key-based side. Both classes are plain siblings extending BaseTestFileGroupRecordBuffer with no lifecycle hooks, and both reach the buffer through the same static buildKeyBasedFileGroupRecordBuffer overload with a non-empty record iterator, which routes to createStreamingRecordsBufferLoader(). The same production path therefore ran twice. Delete the four copies from TestKeyBasedFileGroupRecordBuffer, along with the eight fixture fields and eight imports they alone used. The six remaining tests are byte-identical to before. --- .../buffer/TestKeyBasedFileGroupRecordBuffer.java | 172 --------------------- 1 file changed, 172 deletions(-) diff --git a/hudi-common/src/test/java/org/apache/hudi/common/table/read/buffer/TestKeyBasedFileGroupRecordBuffer.java b/hudi-common/src/test/java/org/apache/hudi/common/table/read/buffer/TestKeyBasedFileGroupRecordBuffer.java index a5dbe4c9d86f..60c0552a444a 100644 --- a/hudi-common/src/test/java/org/apache/hudi/common/table/read/buffer/TestKeyBasedFileGroupRecordBuffer.java +++ b/hudi-common/src/test/java/org/apache/hudi/common/table/read/buffer/TestKeyBasedFileGroupRecordBuffer.java @@ -21,18 +21,13 @@ package org.apache.hudi.common.table.read.buffer; import org.apache.hudi.common.avro.HoodieAvroReaderContext; import org.apache.hudi.common.config.RecordMergeMode; -import org.apache.hudi.common.config.TypedProperties; import org.apache.hudi.common.engine.HoodieReaderContext; import org.apache.hudi.common.model.DeleteRecord; import org.apache.hudi.common.model.HoodieAvroRecordMerger; -import org.apache.hudi.common.model.HoodieRecord; -import org.apache.hudi.common.model.HoodieRecordMerger; import org.apache.hudi.common.table.HoodieTableConfig; -import org.apache.hudi.common.table.HoodieTableVersion; import org.apache.hudi.common.table.log.block.HoodieDataBlock; import org.apache.hudi.common.table.log.block.HoodieDeleteBlock; import org.apache.hudi.common.table.read.BufferedRecord; -import org.apache.hudi.common.table.read.FileGroupReaderSchemaHandler; import org.apache.hudi.common.table.read.HoodieReadStats; import org.apache.hudi.common.util.Option; import org.apache.hudi.common.util.collection.ClosableIterator; @@ -49,10 +44,7 @@ import java.lang.reflect.Modifier; import java.util.Arrays; import java.util.Collections; import java.util.List; -import java.util.stream.Stream; -import static org.apache.hudi.common.model.DefaultHoodieRecordPayload.DELETE_KEY; -import static org.apache.hudi.common.model.DefaultHoodieRecordPayload.DELETE_MARKER; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.mock; @@ -65,19 +57,11 @@ class TestKeyBasedFileGroupRecordBuffer extends BaseTestFileGroupRecordBuffer { private final IndexedRecord testRecord2Update = createTestRecord("2", 1, 2L); private final IndexedRecord testRecord2EarlierUpdate = createTestRecord("2", 1, 0L); private final IndexedRecord testRecord2Delete = createTestRecord("2", 2, 3L); - private final IndexedRecord testRecord2CustomPayloadExpected = createTestRecord("2", 2, 2L); private final IndexedRecord testRecord3 = createTestRecord("3", 1, 1L); private final IndexedRecord testRecord3Update = createTestRecord("3", 1, 2L); - private final IndexedRecord testRecord3UpdateCustomPayloadExpected = createTestRecord("3", 2, 2L); private final IndexedRecord testRecord3DeleteByFieldValue = createTestRecord("3", 3, 1L); private final IndexedRecord testRecord4 = createTestRecord("4", 2, 1L); private final IndexedRecord testRecord4Update = createTestRecord("4", 1, 2L); - private final IndexedRecord testRecord4EarlierUpdate = createTestRecord("4", 1, 0L); - private final IndexedRecord testRecord5 = createTestRecord("5", 1, 1L); - private final IndexedRecord testRecord5DeleteByCustomMarker = createTestRecord("5", 3, 2L); - private final IndexedRecord testRecord6 = createTestRecord("6", 1, 5L); - private final IndexedRecord testRecord6DeleteByCustomMarker = createTestRecord("6", 3, 2L); - private final IndexedRecord testRecord7 = createTestRecord("7", 1, 5L); /** * Asserts that {@code processNextDataRecord} and {@code isPartialMergingEnabled} keep their @@ -168,46 +152,6 @@ class TestKeyBasedFileGroupRecordBuffer extends BaseTestFileGroupRecordBuffer { assertEquals(2, readStats.getNumUpdates()); } - @Test - void readWithEventTimeOrderingWithRecords() throws IOException { - HoodieReadStats readStats = new HoodieReadStats(); - TypedProperties properties = new TypedProperties(); - properties.setProperty(HoodieTableConfig.ORDERING_FIELDS.key(), "ts"); - properties.setProperty(DELETE_KEY, "counter"); - properties.setProperty(DELETE_MARKER, "3"); - HoodieTableConfig tableConfig = mock(HoodieTableConfig.class); - when(tableConfig.getRecordMergeMode()).thenReturn(RecordMergeMode.EVENT_TIME_ORDERING); - when(tableConfig.getPartialUpdateMode()).thenReturn(Option.empty()); - when(tableConfig.getTableVersion()).thenReturn(HoodieTableVersion.current()); - when(tableConfig.getRecordKeyFields()).thenReturn(Option.of(new String[] {"record_key"})); - when(tableConfig.getPartitionFields()).thenReturn(Option.empty()); - StorageConfiguration<?> storageConfiguration = mock(StorageConfiguration.class); - HoodieReaderContext<IndexedRecord> readerContext = new HoodieAvroReaderContext(storageConfiguration, tableConfig, Option.empty(), Option.empty()); - readerContext.setHasLogFiles(false); - readerContext.setHasBootstrapBaseFile(false); - readerContext.initRecordMerger(properties); - FileGroupReaderSchemaHandler schemaHandler = new FileGroupReaderSchemaHandler(readerContext, SCHEMA, SCHEMA, Option.empty(), - properties, createMockMetaClient(tableConfig)); - readerContext.setSchemaHandler(schemaHandler); - List<HoodieRecord> inputRecords = convertToHoodieRecordsList(Arrays.asList(testRecord1UpdateWithSameTime, testRecord2Update, testRecord3Update, testRecord4EarlierUpdate, testRecord7)); - inputRecords.addAll(convertToHoodieRecordsListForDeletes(Arrays.asList(testRecord5DeleteByCustomMarker, testRecord6DeleteByCustomMarker), false)); - KeyBasedFileGroupRecordBuffer<IndexedRecord> fileGroupRecordBuffer = buildKeyBasedFileGroupRecordBuffer(readerContext, tableConfig, readStats, null, - RecordMergeMode.EVENT_TIME_ORDERING, Collections.singletonList("ts"), properties, Option.of(inputRecords.iterator())); - - fileGroupRecordBuffer.setBaseFileIterator(ClosableIterator.wrap(Arrays.asList(testRecord1, testRecord2, testRecord3, testRecord4, - testRecord5, testRecord6).iterator())); - - List<IndexedRecord> actualRecords = getActualRecords(fileGroupRecordBuffer); - // update for 4 is ignored due to lower ordering value. - // record5 is deleted. - // delete for 6 is ignored due to lower ordering value. - assertEquals(Arrays.asList(getSerializableIndexedRecord(testRecord1UpdateWithSameTime), getSerializableIndexedRecord(testRecord2Update), - getSerializableIndexedRecord(testRecord3Update), testRecord4, testRecord6, getSerializableIndexedRecord(testRecord7)), actualRecords); - assertEquals(1, readStats.getNumInserts()); - assertEquals(1, readStats.getNumDeletes()); - assertEquals(3, readStats.getNumUpdates()); - } - @Test void readWithCommitTimeOrdering() throws IOException { HoodieReadStats readStats = new HoodieReadStats(); @@ -239,43 +183,6 @@ class TestKeyBasedFileGroupRecordBuffer extends BaseTestFileGroupRecordBuffer { assertEquals(2, readStats.getNumUpdates()); } - @Test - void readWithCommitTimeOrderingWithRecords() throws IOException { - HoodieReadStats readStats = new HoodieReadStats(); - TypedProperties properties = new TypedProperties(); - properties.setProperty(DELETE_KEY, "counter"); - properties.setProperty(DELETE_MARKER, "3"); - HoodieTableConfig tableConfig = mock(HoodieTableConfig.class); - when(tableConfig.getRecordMergeMode()).thenReturn(RecordMergeMode.COMMIT_TIME_ORDERING); - when(tableConfig.getPartialUpdateMode()).thenReturn(Option.empty()); - when(tableConfig.getTableVersion()).thenReturn(HoodieTableVersion.current()); - when(tableConfig.getRecordKeyFields()).thenReturn(Option.of(new String[] {"record_key"})); - when(tableConfig.getPartitionFields()).thenReturn(Option.empty()); - StorageConfiguration<?> storageConfiguration = mock(StorageConfiguration.class); - HoodieReaderContext<IndexedRecord> readerContext = new HoodieAvroReaderContext(storageConfiguration, tableConfig, Option.empty(), Option.empty()); - readerContext.setHasLogFiles(false); - readerContext.setHasBootstrapBaseFile(false); - readerContext.initRecordMerger(properties); - FileGroupReaderSchemaHandler schemaHandler = new FileGroupReaderSchemaHandler(readerContext, SCHEMA, SCHEMA, Option.empty(), - properties, createMockMetaClient(tableConfig)); - readerContext.setSchemaHandler(schemaHandler); - List<HoodieRecord> inputRecords = convertToHoodieRecordsList(Arrays.asList(testRecord1UpdateWithSameTime, testRecord2Update, testRecord3Update, - testRecord4EarlierUpdate, testRecord7)); - inputRecords.addAll(convertToHoodieRecordsListForDeletes(Arrays.asList(testRecord5DeleteByCustomMarker, testRecord6DeleteByCustomMarker), true)); - KeyBasedFileGroupRecordBuffer<IndexedRecord> fileGroupRecordBuffer = buildKeyBasedFileGroupRecordBuffer(readerContext, tableConfig, readStats, null, - RecordMergeMode.COMMIT_TIME_ORDERING, Collections.singletonList("ts"), properties, Option.of(inputRecords.iterator())); - - fileGroupRecordBuffer.setBaseFileIterator(ClosableIterator.wrap(Arrays.asList(testRecord1, testRecord2, testRecord3, testRecord4, - testRecord5, testRecord6).iterator())); - - List<IndexedRecord> actualRecords = getActualRecords(fileGroupRecordBuffer); - assertEquals(convertGenRecordsToSerializableIndexedRecords(Stream.of(testRecord1UpdateWithSameTime, testRecord2Update, - testRecord3Update, testRecord4EarlierUpdate, testRecord7)), actualRecords); - assertEquals(1, readStats.getNumInserts()); - assertEquals(2, readStats.getNumDeletes()); - assertEquals(4, readStats.getNumUpdates()); - } - @Test void readWithCustomPayload() throws IOException { HoodieReadStats readStats = new HoodieReadStats(); @@ -316,46 +223,6 @@ class TestKeyBasedFileGroupRecordBuffer extends BaseTestFileGroupRecordBuffer { assertEquals(0, readStats.getNumUpdates()); } - @Test - void readWithCustomPayloadWithRecords() throws IOException { - HoodieReadStats readStats = new HoodieReadStats(); - TypedProperties properties = new TypedProperties(); - properties.setProperty(DELETE_KEY, "counter"); - properties.setProperty(DELETE_MARKER, "3"); - properties.setProperty(HoodieTableConfig.RECORD_MERGE_MODE.key(), "CUSTOM"); - properties.setProperty(HoodieTableConfig.PAYLOAD_CLASS_NAME.key(), TestKeyBasedFileGroupRecordBuffer.CustomPayload.class.getName()); - properties.setProperty(HoodieTableConfig.RECORD_MERGE_STRATEGY_ID.key(), HoodieRecordMerger.PAYLOAD_BASED_MERGE_STRATEGY_UUID); - HoodieTableConfig tableConfig = mock(HoodieTableConfig.class); - when(tableConfig.getPayloadClass()).thenReturn(TestKeyBasedFileGroupRecordBuffer.CustomPayload.class.getName()); - when(tableConfig.getRecordKeyFields()).thenReturn(Option.of(new String[] {"record_key"})); - when(tableConfig.getPartitionFields()).thenReturn(Option.empty()); - when(tableConfig.getRecordMergeMode()).thenReturn(RecordMergeMode.CUSTOM); - when(tableConfig.getPartialUpdateMode()).thenReturn(Option.empty()); - when(tableConfig.getRecordMergeStrategyId()).thenReturn(HoodieRecordMerger.PAYLOAD_BASED_MERGE_STRATEGY_UUID); - when(tableConfig.getTableVersion()).thenReturn(HoodieTableVersion.current()); - StorageConfiguration<?> storageConfiguration = mock(StorageConfiguration.class); - HoodieReaderContext<IndexedRecord> readerContext = new HoodieAvroReaderContext(storageConfiguration, tableConfig, Option.empty(), Option.empty()); - readerContext.setHasLogFiles(false); - readerContext.setHasBootstrapBaseFile(false); - readerContext.initRecordMerger(properties); - FileGroupReaderSchemaHandler schemaHandler = new FileGroupReaderSchemaHandler(readerContext, SCHEMA, SCHEMA, Option.empty(), - properties, createMockMetaClient(tableConfig)); - readerContext.setSchemaHandler(schemaHandler); - List<HoodieRecord> inputRecords = convertToHoodieRecordsList(Arrays.asList(testRecord1UpdateWithSameTime, testRecord2Update, testRecord3Update, testRecord4EarlierUpdate)); - inputRecords.addAll(convertToHoodieRecordsListForDeletes(Arrays.asList(testRecord5DeleteByCustomMarker, testRecord6DeleteByCustomMarker), true)); - KeyBasedFileGroupRecordBuffer<IndexedRecord> fileGroupRecordBuffer = buildKeyBasedFileGroupRecordBuffer(readerContext, tableConfig, readStats, new HoodieAvroRecordMerger(), - RecordMergeMode.CUSTOM, Collections.singletonList("ts"), properties, Option.of(inputRecords.iterator())); - - fileGroupRecordBuffer.setBaseFileIterator(ClosableIterator.wrap(Arrays.asList(testRecord1, testRecord2, testRecord3, testRecord4, - testRecord5, testRecord6).iterator())); - - List<IndexedRecord> actualRecords = getActualRecords(fileGroupRecordBuffer); - assertEquals(Arrays.asList(testRecord1, testRecord2CustomPayloadExpected, testRecord3UpdateCustomPayloadExpected), actualRecords); - assertEquals(0, readStats.getNumInserts()); - assertEquals(3, readStats.getNumDeletes()); - assertEquals(2, readStats.getNumUpdates()); - } - @Test void readWithCustomMerger() throws IOException { HoodieReadStats readStats = new HoodieReadStats(); @@ -393,43 +260,4 @@ class TestKeyBasedFileGroupRecordBuffer extends BaseTestFileGroupRecordBuffer { assertEquals(3, readStats.getNumDeletes()); assertEquals(0, readStats.getNumUpdates()); } - - @Test - void readWithCustomMergerWithRecords() throws IOException { - HoodieReadStats readStats = new HoodieReadStats(); - TypedProperties properties = new TypedProperties(); - properties.setProperty(DELETE_KEY, "counter"); - properties.setProperty(DELETE_MARKER, "3"); - properties.setProperty(HoodieTableConfig.PAYLOAD_CLASS_NAME.key(), CustomPayload.class.getName()); - HoodieTableConfig tableConfig = mock(HoodieTableConfig.class); - when(tableConfig.getPayloadClass()).thenReturn(CustomPayload.class.getName()); - when(tableConfig.getRecordKeyFields()).thenReturn(Option.of(new String[] {"record_key"})); - when(tableConfig.getPartitionFields()).thenReturn(Option.empty()); - when(tableConfig.getRecordMergeMode()).thenReturn(RecordMergeMode.CUSTOM); - when(tableConfig.getPartialUpdateMode()).thenReturn(Option.empty()); - when(tableConfig.getRecordMergeStrategyId()).thenReturn(HoodieRecordMerger.PAYLOAD_BASED_MERGE_STRATEGY_UUID); - when(tableConfig.getTableVersion()).thenReturn(HoodieTableVersion.current()); - - StorageConfiguration<?> storageConfiguration = mock(StorageConfiguration.class); - HoodieReaderContext<IndexedRecord> readerContext = new HoodieAvroReaderContext(storageConfiguration, tableConfig, Option.empty(), Option.empty()); - readerContext.setHasLogFiles(false); - readerContext.setHasBootstrapBaseFile(false); - readerContext.initRecordMerger(properties); - FileGroupReaderSchemaHandler schemaHandler = new FileGroupReaderSchemaHandler(readerContext, SCHEMA, SCHEMA, Option.empty(), - properties, createMockMetaClient(tableConfig)); - readerContext.setSchemaHandler(schemaHandler); - List<HoodieRecord> inputRecords = convertToHoodieRecordsList(Arrays.asList(testRecord1UpdateWithSameTime, testRecord2Update, testRecord3Update, testRecord4EarlierUpdate)); - inputRecords.addAll(convertToHoodieRecordsListForDeletes(Arrays.asList(testRecord5DeleteByCustomMarker, testRecord6DeleteByCustomMarker), true)); - KeyBasedFileGroupRecordBuffer<IndexedRecord> fileGroupRecordBuffer = buildKeyBasedFileGroupRecordBuffer(readerContext, tableConfig, readStats, new TestKeyBasedFileGroupRecordBuffer.CustomMerger(), - RecordMergeMode.CUSTOM, Collections.singletonList("ts"), properties, Option.of(inputRecords.iterator())); - - fileGroupRecordBuffer.setBaseFileIterator(ClosableIterator.wrap(Arrays.asList(testRecord1, testRecord2, testRecord3, testRecord4, - testRecord5, testRecord6).iterator())); - - List<IndexedRecord> actualRecords = getActualRecords(fileGroupRecordBuffer); - assertEquals(Arrays.asList(testRecord1, testRecord2CustomPayloadExpected, testRecord3UpdateCustomPayloadExpected), actualRecords); - assertEquals(0, readStats.getNumInserts()); - assertEquals(3, readStats.getNumDeletes()); - assertEquals(2, readStats.getNumUpdates()); - } }
