This is an automated email from the ASF dual-hosted git repository.

danny0405 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git


The following commit(s) were added to refs/heads/master by this push:
     new 3f2de46f5bd1 test(client): improve native log and CDC write coverage 
(#19353)
3f2de46f5bd1 is described below

commit 3f2de46f5bd10abb2fe749c1dc7056d432d16f17
Author: Shuo Cheng <[email protected]>
AuthorDate: Tue Jul 28 17:58:23 2026 +0800

    test(client): improve native log and CDC write coverage (#19353)
---
 ...tFileGroupReaderBasedNativeLogAppendHandle.java | 200 +++++++++++++
 .../hudi/io/TestHoodieNativeLogAppendHandle.java   | 309 +++++++++++++++++++++
 .../io/TestSortedAndChangeLogMergeHandles.java     | 264 ++++++++++++++++++
 .../hudi/io/cdc/TestHoodieNativeCDCFileWriter.java |  98 +++++++
 .../hudi/io/cdc/TestHoodieNativeCDCLoggers.java    | 219 +++++++++++++++
 5 files changed, 1090 insertions(+)

diff --git 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestFileGroupReaderBasedNativeLogAppendHandle.java
 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestFileGroupReaderBasedNativeLogAppendHandle.java
new file mode 100644
index 000000000000..54cb8db22bd3
--- /dev/null
+++ 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestFileGroupReaderBasedNativeLogAppendHandle.java
@@ -0,0 +1,200 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hudi.io;
+
+import org.apache.hudi.client.WriteStatus;
+import org.apache.hudi.common.config.RecordMergeMode;
+import org.apache.hudi.common.engine.HoodieReaderContext;
+import org.apache.hudi.common.engine.LocalTaskContextSupplier;
+import org.apache.hudi.common.model.CompactionOperation;
+import org.apache.hudi.common.model.HoodieDeltaWriteStat;
+import org.apache.hudi.common.model.HoodieFileFormat;
+import org.apache.hudi.common.model.HoodieWriteStat.RuntimeStats;
+import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.HoodieTableVersion;
+import org.apache.hudi.common.table.read.HoodieFileGroupReader;
+import org.apache.hudi.common.table.read.HoodieReadStats;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.common.util.collection.ClosableIterator;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.exception.HoodieIOException;
+import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.StoragePath;
+import org.apache.hudi.table.HoodieTable;
+
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+import java.io.IOException;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.Mockito.RETURNS_SELF;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+public class TestFileGroupReaderBasedNativeLogAppendHandle {
+
+  private static final String SCHEMA = 
"{\"type\":\"record\",\"name\":\"trip\",\"fields\":["
+      + "{\"name\":\"id\",\"type\":\"string\"}]}";
+
+  @Test
+  public void testCompactsLogRecordsAndCopiesReadStatsToFirstOutput() throws 
Exception {
+    HoodieWriteConfig config = config();
+    HoodieTable table = table(config);
+    CompactionOperation operation = operation();
+    HoodieReaderContext readerContext = mock(HoodieReaderContext.class);
+    HoodieFileGroupReader reader = mock(HoodieFileGroupReader.class);
+    ClosableIterator records = mock(ClosableIterator.class);
+    HoodieReadStats readStats = readStats();
+    when(records.hasNext()).thenReturn(false);
+    when(reader.getLogRecordsOnly()).thenReturn(records);
+    when(reader.getValidBlockInstants()).thenReturn(Arrays.asList("001", 
"002"));
+    when(reader.getReadStats()).thenReturn(readStats);
+
+    HoodieFileGroupReader.HoodieFileGroupReaderBuilder builder =
+        mock(HoodieFileGroupReader.HoodieFileGroupReaderBuilder.class, 
RETURNS_SELF);
+    when(builder.build()).thenReturn(reader);
+    try (MockedStatic<HoodieFileGroupReader> fileGroupReaders = 
mockStatic(HoodieFileGroupReader.class)) {
+      
fileGroupReaders.when(HoodieFileGroupReader::builder).thenReturn(builder);
+      TestableFileGroupReaderHandle handle = new TestableFileGroupReaderHandle(
+          config, table, operation, readerContext);
+      handle.doAppend();
+
+      WriteStatus first = status("first.log.parquet");
+      WriteStatus second = status("second.deletes.parquet");
+      handle.addStatuses(first, second);
+      List<WriteStatus> statuses = handle.close();
+
+      assertEquals(2, statuses.size());
+      for (WriteStatus status : statuses) {
+        assertEquals("partition", status.getStat().getPartitionPath());
+        assertEquals("001", status.getStat().getPrevCommit());
+      }
+      assertEquals(11L, first.getStat().getTotalLogReadTimeMs());
+      assertEquals(12L, first.getStat().getTotalUpdatedRecordsCompacted());
+      assertEquals(13L, first.getStat().getTotalLogFilesCompacted());
+      assertEquals(14L, first.getStat().getTotalLogRecords());
+      assertEquals(15L, first.getStat().getTotalLogBlocks());
+      assertEquals(16L, first.getStat().getTotalCorruptLogBlock());
+      assertEquals(17L, first.getStat().getTotalRollbackBlocks());
+      assertEquals(18L, first.getStat().getTotalLogSizeCompacted());
+      assertEquals(11L, first.getStat().getRuntimeStats().getTotalScanTime());
+      assertEquals(0L, second.getStat().getTotalLogReadTimeMs());
+      verify(reader).close();
+    }
+  }
+
+  @Test
+  public void testWrapsReaderInitializationFailure() throws Exception {
+    HoodieWriteConfig config = config();
+    HoodieTable table = table(config);
+    HoodieFileGroupReader reader = mock(HoodieFileGroupReader.class);
+    when(reader.getLogRecordsOnly()).thenThrow(new IOException("failure"));
+    HoodieFileGroupReader.HoodieFileGroupReaderBuilder builder =
+        mock(HoodieFileGroupReader.HoodieFileGroupReaderBuilder.class, 
RETURNS_SELF);
+    when(builder.build()).thenReturn(reader);
+
+    try (MockedStatic<HoodieFileGroupReader> fileGroupReaders = 
mockStatic(HoodieFileGroupReader.class)) {
+      
fileGroupReaders.when(HoodieFileGroupReader::builder).thenReturn(builder);
+      TestableFileGroupReaderHandle handle = new TestableFileGroupReaderHandle(
+          config, table, operation(), mock(HoodieReaderContext.class));
+      assertThrows(HoodieIOException.class, handle::doAppend);
+      verify(reader).close();
+    }
+  }
+
+  private static HoodieWriteConfig config() {
+    return HoodieWriteConfig.newBuilder()
+        .withPath("/tmp")
+        .withSchema(SCHEMA)
+        .withWriteTableVersion(HoodieTableVersion.TEN.versionCode())
+        .withWriteRecordPositionsEnabled(false)
+        .build();
+  }
+
+  private static HoodieTable table(HoodieWriteConfig config) {
+    HoodieTable table = mock(HoodieTable.class);
+    HoodieStorage storage = mock(HoodieStorage.class);
+    HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+    HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
+    when(table.getStorage()).thenReturn(storage);
+    when(table.getConfig()).thenReturn(config);
+    when(table.getMetaClient()).thenReturn(metaClient);
+    when(table.getBaseFileFormat()).thenReturn(HoodieFileFormat.PARQUET);
+    when(table.version()).thenReturn(HoodieTableVersion.TEN);
+    when(metaClient.getBasePath()).thenReturn(new StoragePath("/tmp"));
+    when(metaClient.getTableConfig()).thenReturn(tableConfig);
+    when(metaClient.getIndexMetadata()).thenReturn(Option.empty());
+    when(tableConfig.getTableVersion()).thenReturn(HoodieTableVersion.TEN);
+    
when(tableConfig.getRecordMergeMode()).thenReturn(RecordMergeMode.COMMIT_TIME_ORDERING);
+    when(tableConfig.getPayloadClassIfPresent()).thenReturn(Option.empty());
+    return table;
+  }
+
+  private static CompactionOperation operation() {
+    CompactionOperation operation = mock(CompactionOperation.class);
+    when(operation.getPartitionPath()).thenReturn("partition");
+    when(operation.getFileId()).thenReturn("file-1");
+    when(operation.getBaseInstantTime()).thenReturn("001");
+    
when(operation.getDeltaFileNames()).thenReturn(Collections.singletonList("file-1.log.parquet"));
+    return operation;
+  }
+
+  private static HoodieReadStats readStats() {
+    HoodieReadStats stats = new HoodieReadStats();
+    stats.setTotalLogReadTimeMs(11L);
+    stats.setTotalUpdatedRecordsCompacted(12L);
+    stats.setTotalLogFilesCompacted(13L);
+    stats.setTotalLogRecords(14L);
+    stats.setTotalLogBlocks(15L);
+    stats.setTotalCorruptLogBlock(16L);
+    stats.setTotalRollbackBlocks(17L);
+    stats.setTotalLogSizeCompacted(18L);
+    return stats;
+  }
+
+  private static WriteStatus status(String path) {
+    WriteStatus status = new WriteStatus(false, 0.0);
+    HoodieDeltaWriteStat stat = new HoodieDeltaWriteStat();
+    stat.setPath(path);
+    stat.setRuntimeStats(new RuntimeStats());
+    status.setStat(stat);
+    return status;
+  }
+
+  private static class TestableFileGroupReaderHandle extends 
FileGroupReaderBasedNativeLogAppendHandle {
+
+    private TestableFileGroupReaderHandle(
+        HoodieWriteConfig config, HoodieTable table, CompactionOperation 
operation,
+        HoodieReaderContext readerContext) {
+      super(config, "100", table, operation, new LocalTaskContextSupplier(), 
readerContext);
+    }
+
+    private void addStatuses(WriteStatus... writeStatuses) {
+      statuses.addAll(Arrays.asList(writeStatuses));
+    }
+  }
+}
diff --git 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieNativeLogAppendHandle.java
 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieNativeLogAppendHandle.java
new file mode 100644
index 000000000000..7787a2b542a5
--- /dev/null
+++ 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieNativeLogAppendHandle.java
@@ -0,0 +1,309 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hudi.io;
+
+import org.apache.hudi.client.WriteStatus;
+import org.apache.hudi.common.config.RecordMergeMode;
+import org.apache.hudi.common.engine.LocalTaskContextSupplier;
+import org.apache.hudi.common.engine.RecordContext;
+import org.apache.hudi.common.fs.FSUtils;
+import org.apache.hudi.common.model.HoodieDeltaWriteStat;
+import org.apache.hudi.common.model.HoodieFileFormat;
+import org.apache.hudi.common.model.HoodieLogFile;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.schema.HoodieSchema;
+import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.HoodieTableVersion;
+import org.apache.hudi.common.table.log.AppendResult;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.exception.HoodieAppendException;
+import org.apache.hudi.io.cdc.HoodieNativeLogFormatWriter;
+import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.StoragePath;
+import org.apache.hudi.table.HoodieTable;
+
+import org.junit.jupiter.api.Test;
+import org.mockito.InOrder;
+import org.mockito.MockedConstruction;
+
+import java.io.IOException;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.inOrder;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+public class TestHoodieNativeLogAppendHandle {
+
+  private static final String SCHEMA = 
"{\"type\":\"record\",\"name\":\"trip\",\"fields\":["
+      + 
"{\"name\":\"id\",\"type\":\"string\"},{\"name\":\"ts\",\"type\":\"long\"}]}";
+
+  @Test
+  public void testCreatesWriterAndRoutesDataAndDeleteRecords() throws 
Exception {
+    HoodieWriteConfig config = config();
+    HoodieTable table = table(config);
+    HoodieRecord inputRecord = mock(HoodieRecord.class);
+    HoodieRecord populatedRecord = mock(HoodieRecord.class);
+    when(inputRecord.prependMetaFields(any(HoodieSchema.class), 
any(HoodieSchema.class), any(), any()))
+        .thenReturn(populatedRecord);
+
+    try (MockedConstruction<HoodieNativeLogFormatWriter> writers = 
mockConstruction(
+        HoodieNativeLogFormatWriter.class, (writer, context) -> {
+          when(writer.canWriteDataFile()).thenReturn(true);
+          when(writer.canWriteDeleteFile()).thenReturn(true);
+          when(writer.hasPendingWrites()).thenReturn(true);
+          
when(writer.getLastAppendResults()).thenReturn(Collections.emptyList());
+          when(writer.getLogFile()).thenReturn(new HoodieLogFile(new 
StoragePath("/tmp/native.log.parquet")));
+        })) {
+      TestableNativeLogAppendHandle handle = new 
TestableNativeLogAppendHandle(config, table);
+      handle.createWriter();
+      HoodieNativeLogFormatWriter writer = writers.constructed().get(0);
+
+      handle.writeData(inputRecord, true);
+      verify(writer).appendRecord(eq(populatedRecord), any(HoodieSchema.class),
+          eq(HoodieRecord.RECORD_KEY_METADATA_FIELD));
+      handle.writeDeleteRecord(inputRecord);
+      verify(inputRecord).clearNewLocation();
+      verify(writer).appendDeleteRecord(eq(inputRecord), 
any(HoodieSchema.class),
+          eq(HoodieRecord.RECORD_KEY_METADATA_FIELD));
+
+      handle.flushWriter();
+      verify(writer).flushAppend(any());
+      assertEquals("/tmp/native.log.parquet", handle.logFilePath().toString());
+      assertTrue(handle.canWrite(inputRecord));
+      handle.closeWriter();
+      verify(writer).close();
+    }
+  }
+
+  @Test
+  public void testFlushesBeforeWriterRolloverAndTracksPerFileCounts() throws 
Exception {
+    HoodieWriteConfig config = config();
+    HoodieTable table = table(config);
+    HoodieRecord inputRecord = mock(HoodieRecord.class);
+    HoodieRecord populatedRecord = mock(HoodieRecord.class);
+    when(inputRecord.prependMetaFields(any(HoodieSchema.class), 
any(HoodieSchema.class), any(), any()))
+        .thenReturn(populatedRecord);
+
+    try (MockedConstruction<HoodieNativeLogFormatWriter> writers = 
mockConstruction(
+        HoodieNativeLogFormatWriter.class, (writer, context) -> {
+          when(writer.canWriteDataFile()).thenReturn(false);
+          when(writer.hasPendingWrites()).thenReturn(true);
+          when(writer.getLastAppendResults()).thenReturn(
+              Arrays.asList(appendResult(1, "log", 13L), appendResult(1, 
"deletes", 7L)),
+              Collections.singletonList(appendResult(2, "log", 11L)));
+          
when(writer.getLastDataFileFormatMetadata()).thenReturn(Option.empty());
+        })) {
+      TestableNativeLogAppendHandle handle = new 
TestableNativeLogAppendHandle(config, table);
+      handle.createWriter();
+      HoodieNativeLogFormatWriter writer = writers.constructed().get(0);
+
+      handle.setCounts(5, 2, 3, 4);
+      handle.writeData(inputRecord, false);
+      InOrder rolloverOrder = inOrder(writer);
+      rolloverOrder.verify(writer).flushAppend(any());
+      rolloverOrder.verify(writer).appendRecord(eq(populatedRecord), 
any(HoodieSchema.class), any());
+
+      handle.flushWriter();
+      verify(writer, times(2)).flushAppend(any());
+
+      List<WriteStatus> statuses = handle.getWriteStatuses();
+      assertEquals(3, statuses.size());
+      HoodieDeltaWriteStat firstDataStat = (HoodieDeltaWriteStat) 
statuses.get(0).getStat();
+      assertEquals(5, firstDataStat.getNumWrites());
+      assertEquals(2, firstDataStat.getNumUpdateWrites());
+      assertEquals(3, firstDataStat.getNumInserts());
+      assertEquals(13L, firstDataStat.getTotalWriteBytes());
+
+      HoodieDeltaWriteStat deleteStat = (HoodieDeltaWriteStat) 
statuses.get(1).getStat();
+      assertEquals(4, deleteStat.getNumDeletes());
+      assertEquals(7L, deleteStat.getTotalWriteBytes());
+
+      HoodieDeltaWriteStat secondDataStat = (HoodieDeltaWriteStat) 
statuses.get(2).getStat();
+      assertEquals(1, secondDataStat.getNumWrites());
+      assertEquals(0, secondDataStat.getNumUpdateWrites());
+      assertEquals(1, secondDataStat.getNumInserts());
+      assertEquals(11L, secondDataStat.getTotalWriteBytes());
+
+      assertThrows(HoodieAppendException.class,
+          () -> handle.accumulateWriteCounts(firstDataStat, appendResult(1, 
"log", 13L)));
+    }
+  }
+
+  @Test
+  public void 
testUsesConfiguredKeyWithoutMetadataFieldsAndSkipsIgnoredRecords() throws 
Exception {
+    HoodieWriteConfig config = HoodieWriteConfig.newBuilder()
+        .withPath("/tmp")
+        .withSchema(SCHEMA)
+        .withPopulateMetaFields(false)
+        .withWriteTableVersion(HoodieTableVersion.TEN.versionCode())
+        .withWriteRecordPositionsEnabled(false)
+        .build();
+    HoodieTable table = table(config);
+    
when(table.getMetaClient().getTableConfig().getRecordKeyFieldProp()).thenReturn("id");
+
+    try (MockedConstruction<HoodieNativeLogFormatWriter> writers = 
mockConstruction(
+        HoodieNativeLogFormatWriter.class, (writer, context) -> {
+          when(writer.canWriteDataFile()).thenReturn(false);
+          when(writer.canWriteDeleteFile()).thenReturn(false);
+          when(writer.hasPendingWrites()).thenReturn(false);
+        })) {
+      TestableNativeLogAppendHandle handle = new 
TestableNativeLogAppendHandle(config, table, new HashMap<>());
+      handle.createWriter();
+      HoodieNativeLogFormatWriter writer = writers.constructed().get(0);
+
+      HoodieRecord ignoredRecord = mock(HoodieRecord.class);
+      when(ignoredRecord.shouldIgnore(any(HoodieSchema.class), 
any())).thenReturn(true);
+      handle.writeData(ignoredRecord, false);
+      verify(writer, never()).appendRecord(eq(ignoredRecord), any(), any());
+
+      HoodieRecord inputRecord = mock(HoodieRecord.class);
+      HoodieRecord populatedRecord = mock(HoodieRecord.class);
+      when(inputRecord.prependMetaFields(any(HoodieSchema.class), 
any(HoodieSchema.class), any(), any()))
+          .thenReturn(populatedRecord);
+      handle.writeData(inputRecord, false);
+      verify(writer).appendRecord(eq(populatedRecord), 
any(HoodieSchema.class), eq("id"));
+      handle.writeDeleteWithoutMetadata(inputRecord);
+      verify(writer).appendDeleteRecord(eq(inputRecord), 
any(HoodieSchema.class), eq("id"));
+    }
+  }
+
+  @Test
+  public void testWrapsFlushFailure() throws Exception {
+    HoodieWriteConfig config = config();
+    HoodieTable table = table(config);
+
+    try (MockedConstruction<HoodieNativeLogFormatWriter> writers = 
mockConstruction(
+        HoodieNativeLogFormatWriter.class, (writer, context) -> {
+          when(writer.hasPendingWrites()).thenReturn(true);
+          doThrow(new IOException("flush 
failed")).when(writer).flushAppend(any());
+        })) {
+      TestableNativeLogAppendHandle handle = new 
TestableNativeLogAppendHandle(config, table);
+      handle.createWriter();
+
+      HoodieAppendException exception = 
assertThrows(HoodieAppendException.class, handle::flushWriter);
+      assertTrue(exception.getMessage().contains("file-1"));
+    }
+  }
+
+  private static HoodieWriteConfig config() {
+    return HoodieWriteConfig.newBuilder()
+        .withPath("/tmp")
+        .withSchema(SCHEMA)
+        .withWriteTableVersion(HoodieTableVersion.TEN.versionCode())
+        .withWriteRecordPositionsEnabled(false)
+        .build();
+  }
+
+  private static HoodieTable table(HoodieWriteConfig config) {
+    HoodieTable table = mock(HoodieTable.class);
+    HoodieStorage storage = mock(HoodieStorage.class);
+    HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+    HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
+    when(table.getStorage()).thenReturn(storage);
+    when(table.getConfig()).thenReturn(config);
+    when(table.getMetaClient()).thenReturn(metaClient);
+    when(table.getBaseFileFormat()).thenReturn(HoodieFileFormat.PARQUET);
+    
when(table.getRecordContextForWrite()).thenReturn(mock(RecordContext.class));
+    when(table.version()).thenReturn(HoodieTableVersion.TEN);
+    when(metaClient.getBasePath()).thenReturn(new StoragePath("/tmp"));
+    when(metaClient.getTableConfig()).thenReturn(tableConfig);
+    when(metaClient.getIndexMetadata()).thenReturn(Option.empty());
+    when(tableConfig.getTableVersion()).thenReturn(HoodieTableVersion.TEN);
+    
when(tableConfig.getRecordMergeMode()).thenReturn(RecordMergeMode.COMMIT_TIME_ORDERING);
+    when(tableConfig.getPayloadClassIfPresent()).thenReturn(Option.empty());
+    return table;
+  }
+
+  private static AppendResult appendResult(int version, String extension, long 
size) throws IOException {
+    StoragePath path = new StoragePath("/tmp", FSUtils.makeNativeLogFileName(
+        "file-1", "1-0-1", "100", version, extension, 
HoodieFileFormat.PARQUET));
+    return new AppendResult(new HoodieLogFile(path), 0L, size);
+  }
+
+  private static class TestableNativeLogAppendHandle extends 
HoodieNativeLogAppendHandle {
+
+    private TestableNativeLogAppendHandle(HoodieWriteConfig config, 
HoodieTable table) {
+      super(config, "100", table, "partition", "file-1", 
Collections.emptyIterator(),
+          new LocalTaskContextSupplier());
+    }
+
+    private TestableNativeLogAppendHandle(
+        HoodieWriteConfig config, HoodieTable table, HashMap header) {
+      super(config, "100", table, "partition", "file-1", 
Collections.emptyIterator(),
+          new LocalTaskContextSupplier(), header);
+    }
+
+    private void createWriter() {
+      HoodieDeltaWriteStat stat = new HoodieDeltaWriteStat();
+      stat.setPartitionPath(partitionPath);
+      stat.setFileId(fileId);
+      writeStatus.setStat(stat);
+      writeStatus.setFileId(fileId);
+      writeStatus.setPartitionPath(partitionPath);
+      createLogWriterForAppend("100", Option.empty());
+    }
+
+    private void writeData(HoodieRecord record, boolean isUpdate) throws 
IOException {
+      writeInsertAndUpdate(writeSchema, record, isUpdate);
+    }
+
+    private void writeDeleteRecord(HoodieRecord record) throws IOException {
+      writeDelete(writeSchemaWithMetaFields, record);
+    }
+
+    private void writeDeleteWithoutMetadata(HoodieRecord record) throws 
IOException {
+      writeDelete(writeSchema, record);
+    }
+
+    private void flushWriter() {
+      flushAppend();
+    }
+
+    private void closeWriter() {
+      closeLogWriter();
+    }
+
+    private StoragePath logFilePath() {
+      return getLogFilePath();
+    }
+
+    private void setCounts(long writes, long updates, long inserts, long 
deletes) {
+      recordsWritten = writes;
+      updatedRecordsWritten = updates;
+      insertRecordsWritten = inserts;
+      recordsDeleted = deletes;
+    }
+  }
+}
diff --git 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestSortedAndChangeLogMergeHandles.java
 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestSortedAndChangeLogMergeHandles.java
new file mode 100644
index 000000000000..4949a926e6ec
--- /dev/null
+++ 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestSortedAndChangeLogMergeHandles.java
@@ -0,0 +1,264 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hudi.io;
+
+import org.apache.hudi.client.WriteStatus;
+import org.apache.hudi.common.config.RecordMergeMode;
+import org.apache.hudi.common.engine.HoodieReaderContext;
+import org.apache.hudi.common.engine.LocalTaskContextSupplier;
+import org.apache.hudi.common.engine.ReaderContextFactory;
+import org.apache.hudi.common.model.HoodieAvroIndexedRecord;
+import org.apache.hudi.common.model.HoodieFileFormat;
+import org.apache.hudi.common.model.HoodieKey;
+import org.apache.hudi.common.model.HoodieOperation;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.schema.HoodieSchema;
+import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.HoodieTableVersion;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.core.io.storage.HoodieFileWriter;
+import org.apache.hudi.core.io.storage.HoodieFileWriterFactory;
+import org.apache.hudi.exception.HoodieUpsertException;
+import org.apache.hudi.io.cdc.HoodieCDCLogWriter;
+import org.apache.hudi.io.cdc.HoodieCDCLogWriterFactory;
+import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.StoragePath;
+import org.apache.hudi.storage.StoragePathInfo;
+import org.apache.hudi.table.HoodieTable;
+import org.apache.hudi.table.marker.WriteMarkers;
+import org.apache.hudi.table.marker.WriteMarkersFactory;
+
+import org.apache.avro.generic.GenericRecord;
+import org.apache.avro.generic.IndexedRecord;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Properties;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+public class TestSortedAndChangeLogMergeHandles {
+
+  private static final String SCHEMA = 
"{\"type\":\"record\",\"name\":\"trip\",\"fields\":["
+      + "{\"name\":\"id\",\"type\":\"string\"}]}";
+
+  @Test
+  public void testSortedMergeWritesInsertsAroundExistingRecordsAndOnClose() 
throws Exception {
+    HoodieWriteConfig config = config();
+    TestContext context = new TestContext(config);
+    HoodieRecord recordB = record("b");
+    HoodieRecord recordD = record("d");
+    Map<String, HoodieRecord> records = new HashMap<>();
+    records.put("d", recordD);
+    records.put("b", recordB);
+
+    try (MockedStatic<WriteMarkersFactory> markers = 
mockStatic(WriteMarkersFactory.class);
+         MockedStatic<HoodieFileWriterFactory> writers = 
mockStatic(HoodieFileWriterFactory.class)) {
+      context.stubWriters(markers, writers);
+      TestableSortedMergeHandle handle = new TestableSortedMergeHandle(config, 
context.table, records);
+      handle.write(record("c"));
+      assertEquals(Collections.singletonList("b"), handle.writtenInserts);
+
+      List<WriteStatus> statuses = handle.close();
+      assertEquals(Arrays.asList("b", "d"), handle.writtenInserts);
+      assertEquals(1, statuses.size());
+      verify(context.fileWriter).close();
+    }
+  }
+
+  @Test
+  public void testSortedMergeRejectsAlreadyWrittenInsert() throws Exception {
+    HoodieWriteConfig config = config();
+    TestContext context = new TestContext(config);
+    Map<String, HoodieRecord> records = Collections.singletonMap("b", 
record("b"));
+
+    try (MockedStatic<WriteMarkersFactory> markers = 
mockStatic(WriteMarkersFactory.class);
+         MockedStatic<HoodieFileWriterFactory> writers = 
mockStatic(HoodieFileWriterFactory.class)) {
+      context.stubWriters(markers, writers);
+      TestableSortedMergeHandle handle = new TestableSortedMergeHandle(config, 
context.table, records);
+      handle.markWritten("b");
+      assertThrows(HoodieUpsertException.class, () -> 
handle.write(record("c")));
+    }
+  }
+
+  @Test
+  public void testChangeLogHandleWritesUpdateAndInsertCDCAndPublishesStats() 
throws Exception {
+    HoodieWriteConfig config = config();
+    TestContext context = new TestContext(config);
+    HoodieCDCLogWriter<IndexedRecord> cdcWriter = 
mock(HoodieCDCLogWriter.class);
+    Map<String, Long> cdcStats = 
Collections.singletonMap("partition/file.cdc.parquet", 19L);
+    when(cdcWriter.getCDCWriteStats()).thenReturn(cdcStats);
+
+    try (MockedStatic<WriteMarkersFactory> markers = 
mockStatic(WriteMarkersFactory.class);
+         MockedStatic<HoodieFileWriterFactory> writers = 
mockStatic(HoodieFileWriterFactory.class);
+         MockedStatic<HoodieCDCLogWriterFactory> cdcWriters = 
mockStatic(HoodieCDCLogWriterFactory.class)) {
+      context.stubWriters(markers, writers);
+      cdcWriters.when(() -> HoodieCDCLogWriterFactory.createAvroCDCLogWriter(
+          anyString(), any(), any(), anyString(), any(), any(), anyString(), 
anyString(), any(), any(), any()))
+          .thenReturn(cdcWriter);
+
+      TestableChangeLogMergeHandle handle = new TestableChangeLogMergeHandle(
+          config, context.table, new HashMap<>());
+      HoodieRecord insert = record("i");
+      HoodieRecord savedInsert = record("i");
+      when(insert.newInstance()).thenReturn(savedInsert);
+      when(savedInsert.toIndexedRecord(any(HoodieSchema.class), any()))
+          .thenReturn(Option.of(mock(HoodieAvroIndexedRecord.class)));
+      handle.writeInsert(insert);
+
+      HoodieRecord newRecord = record("u");
+      HoodieRecord oldRecord = record("u");
+      HoodieRecord combinedRecord = record("u");
+      GenericRecord oldData = mock(GenericRecord.class);
+      GenericRecord combinedData = mock(GenericRecord.class);
+      when(oldRecord.getData()).thenReturn(oldData);
+      when(combinedRecord.getData()).thenReturn(combinedData);
+      when(combinedRecord.newInstance()).thenReturn(combinedRecord);
+      when(combinedRecord.toIndexedRecord(any(HoodieSchema.class), 
any())).thenReturn(Option.empty());
+      handle.writeUpdate(newRecord, oldRecord, combinedRecord);
+
+      List<WriteStatus> statuses = handle.close();
+      assertEquals(cdcStats, statuses.get(0).getStat().getCdcStats());
+      verify(cdcWriter, times(2)).put(any(HoodieRecord.class), any(), any());
+      verify(cdcWriter).close();
+    }
+  }
+
+  private static HoodieWriteConfig config() {
+    return HoodieWriteConfig.newBuilder()
+        .withPath("/tmp")
+        .withSchema(SCHEMA)
+        .withWriteTableVersion(HoodieTableVersion.TEN.versionCode())
+        .build();
+  }
+
+  private static HoodieRecord record(String key) {
+    HoodieRecord record = mock(HoodieRecord.class);
+    when(record.getRecordKey()).thenReturn(key);
+    when(record.getRecordKey(any(HoodieSchema.class), 
any(Option.class))).thenReturn(key);
+    when(record.getPartitionPath()).thenReturn("partition");
+    when(record.getKey()).thenReturn(new HoodieKey(key, "partition"));
+    when(record.getOperation()).thenReturn(HoodieOperation.INSERT);
+    when(record.newInstance()).thenReturn(record);
+    when(record.getMetadata()).thenReturn(Option.empty());
+    return record;
+  }
+
+  private static class TestContext {
+    private final HoodieTable table = mock(HoodieTable.class);
+    private final HoodieStorage storage = mock(HoodieStorage.class);
+    private final HoodieFileWriter fileWriter = mock(HoodieFileWriter.class);
+    private final WriteMarkers writeMarkers = mock(WriteMarkers.class);
+
+    private TestContext(HoodieWriteConfig config) throws IOException {
+      HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+      HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
+      ReaderContextFactory contextFactory = mock(ReaderContextFactory.class);
+      when(table.getStorage()).thenReturn(storage);
+      when(table.getConfig()).thenReturn(config);
+      when(table.getMetaClient()).thenReturn(metaClient);
+      
when(table.getBaseFileExtension()).thenReturn(HoodieFileFormat.PARQUET.getFileExtension());
+      when(table.getPartitionMetafileFormat()).thenReturn(Option.empty());
+      when(table.getReaderContextFactoryForWrite()).thenReturn(contextFactory);
+      
when(contextFactory.getContext()).thenReturn(mock(HoodieReaderContext.class));
+      when(metaClient.getBasePath()).thenReturn(new StoragePath("/tmp"));
+      when(metaClient.getTableConfig()).thenReturn(tableConfig);
+      when(metaClient.getIndexMetadata()).thenReturn(Option.empty());
+      when(tableConfig.getTableVersion()).thenReturn(HoodieTableVersion.TEN);
+      
when(tableConfig.getRecordMergeMode()).thenReturn(RecordMergeMode.COMMIT_TIME_ORDERING);
+      when(tableConfig.getPayloadClassIfPresent()).thenReturn(Option.empty());
+      
when(tableConfig.getPartitionMetafileFormat()).thenReturn(Option.empty());
+      when(storage.getPathInfo(any(StoragePath.class))).thenAnswer(invocation 
->
+          new StoragePathInfo(invocation.getArgument(0), 23L, false, (short) 
1, 1L, 1L));
+    }
+
+    private void stubWriters(
+        MockedStatic<WriteMarkersFactory> markers, 
MockedStatic<HoodieFileWriterFactory> writers) {
+      markers.when(() -> WriteMarkersFactory.get(any(), 
any(HoodieTable.class), anyString()))
+          .thenReturn(writeMarkers);
+      writers.when(() -> HoodieFileWriterFactory.getFileWriter(
+          anyString(), any(StoragePath.class), any(HoodieStorage.class), 
any(HoodieWriteConfig.class),
+          any(HoodieSchema.class), any(), 
any(HoodieRecord.HoodieRecordType.class)))
+          .thenReturn(fileWriter);
+    }
+  }
+
+  private static class TestableSortedMergeHandle extends 
HoodieSortedMergeHandle {
+    private final List<String> writtenInserts = new ArrayList<>();
+
+    private TestableSortedMergeHandle(
+        HoodieWriteConfig config, HoodieTable table, Map<String, HoodieRecord> 
records) {
+      super(config, "100", table, records, "partition", "file-1", null,
+          new LocalTaskContextSupplier(), Option.empty());
+    }
+
+    @Override
+    protected boolean writeRecord(
+        HoodieRecord newRecord, HoodieRecord insertRecord, HoodieSchema 
schema, Properties props) {
+      writtenInserts.add(newRecord.getRecordKey());
+      return true;
+    }
+
+    @Override
+    protected void writeToFile(
+        HoodieKey key, HoodieRecord record, HoodieSchema schema, Properties 
props,
+        boolean shouldPreserveRecordMetadata) {
+      // The test targets merge ordering; the physical writer is covered 
separately.
+    }
+
+    private void markWritten(String key) {
+      writtenRecordKeys.add(key);
+    }
+  }
+
+  private static class TestableChangeLogMergeHandle extends 
HoodieMergeHandleWithChangeLog {
+
+    private TestableChangeLogMergeHandle(
+        HoodieWriteConfig config, HoodieTable table, Map<String, HoodieRecord> 
records) {
+      super(config, "100", table, records, "partition", "file-1", null,
+          new LocalTaskContextSupplier(), Option.empty());
+    }
+
+    private void writeInsert(HoodieRecord record) throws IOException {
+      writeInsertRecord(record);
+    }
+
+    private boolean writeUpdate(
+        HoodieRecord newRecord, HoodieRecord oldRecord, HoodieRecord 
combinedRecord) throws IOException {
+      return writeUpdateRecord(newRecord, oldRecord, combinedRecord, 
writeSchema);
+    }
+  }
+}
diff --git 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/cdc/TestHoodieNativeCDCFileWriter.java
 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/cdc/TestHoodieNativeCDCFileWriter.java
index d7d3ae876917..341358fd6e1d 100644
--- 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/cdc/TestHoodieNativeCDCFileWriter.java
+++ 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/cdc/TestHoodieNativeCDCFileWriter.java
@@ -21,6 +21,7 @@ package org.apache.hudi.io.cdc;
 import org.apache.hudi.common.config.TypedProperties;
 import org.apache.hudi.common.engine.TaskContextSupplier;
 import org.apache.hudi.common.model.HoodieFileFormat;
+import org.apache.hudi.common.model.HoodieLogFile;
 import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.schema.HoodieSchema;
 import org.apache.hudi.common.table.log.LogFileCreationCallback;
@@ -29,8 +30,10 @@ import 
org.apache.hudi.common.table.log.block.HoodieLogBlock.HeaderMetadataType;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.core.io.storage.HoodieFileWriter;
 import org.apache.hudi.core.io.storage.HoodieFileWriterFactory;
+import org.apache.hudi.exception.HoodieUpsertException;
 import org.apache.hudi.storage.HoodieStorage;
 import org.apache.hudi.storage.StoragePath;
+import org.apache.hudi.storage.StoragePathInfo;
 
 import org.apache.avro.generic.IndexedRecord;
 import org.junit.jupiter.api.Test;
@@ -40,10 +43,14 @@ import org.mockito.MockedStatic;
 import java.util.Map;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
+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.eq;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
@@ -94,4 +101,95 @@ public class TestHoodieNativeCDCFileWriter {
     assertEquals(instantTime, header.get(HeaderMetadataType.INSTANT_TIME));
     assertEquals(schemaString, header.get(HeaderMetadataType.SCHEMA));
   }
+
+  @Test
+  public void testRollsFilesSkipsExistingVersionsAndReportsStats() throws 
Exception {
+    String instantTime = "100";
+    HoodieStorage storage = mock(HoodieStorage.class);
+    HoodieWriteConfig config = mock(HoodieWriteConfig.class);
+    HoodieSchema cdcSchema = mock(HoodieSchema.class);
+    HoodieFileWriter firstFileWriter = mock(HoodieFileWriter.class);
+    HoodieFileWriter secondFileWriter = mock(HoodieFileWriter.class);
+    LogFileCreationCallback creationCallback = 
mock(LogFileCreationCallback.class);
+    StoragePath parentPath = new StoragePath("/tmp/partition");
+
+    when(config.getProps()).thenReturn(new TypedProperties());
+    when(cdcSchema.toString()).thenReturn("cdc-schema");
+    // Version 1 already exists; the first and second writes therefore use 
versions 2 and 3.
+    when(storage.exists(any(StoragePath.class))).thenReturn(true, false, 
false);
+    when(firstFileWriter.canWrite()).thenReturn(false);
+    when(storage.getPathInfo(any(StoragePath.class))).thenAnswer(invocation ->
+        new StoragePathInfo(invocation.getArgument(0), 17L, false, (short) 1, 
1L, 1L));
+
+    try (MockedStatic<HoodieFileWriterFactory> writerFactory = 
mockStatic(HoodieFileWriterFactory.class)) {
+      writerFactory.when(() -> HoodieFileWriterFactory.getFileWriter(
+              eq(instantTime), any(StoragePath.class), eq(storage), 
eq(config), eq(cdcSchema),
+              any(TaskContextSupplier.class), 
eq(HoodieRecord.HoodieRecordType.AVRO)))
+          .thenReturn(firstFileWriter, secondFileWriter);
+
+      HoodieNativeCDCFileWriter<IndexedRecord> writer = new 
HoodieNativeCDCFileWriter<>(
+          instantTime, "partition", storage, config, cdcSchema, 
HoodieFileFormat.PARQUET,
+          parentPath, "file-1", "1-0-1", creationCallback, 
mock(TaskContextSupplier.class),
+          HoodieRecord.HoodieRecordType.AVRO);
+
+      IndexedRecord firstRecord = mock(IndexedRecord.class);
+      IndexedRecord secondRecord = mock(IndexedRecord.class);
+      writer.write("key-1", firstRecord);
+      writer.write("key-2", secondRecord);
+
+      verify(firstFileWriter).writeRow("key-1", firstRecord);
+      verify(firstFileWriter).close();
+      verify(secondFileWriter).writeRow("key-2", secondRecord);
+      verify(creationCallback, times(2)).preFileCreation(any());
+
+      ArgumentCaptor<StoragePath> pathCaptor = 
ArgumentCaptor.forClass(StoragePath.class);
+      writerFactory.verify(() -> HoodieFileWriterFactory.getFileWriter(
+          eq(instantTime), pathCaptor.capture(), eq(storage), eq(config), 
eq(cdcSchema),
+          any(TaskContextSupplier.class), 
eq(HoodieRecord.HoodieRecordType.AVRO)), times(2));
+      assertNotEquals(pathCaptor.getAllValues().get(0), 
pathCaptor.getAllValues().get(1));
+      assertEquals(2, new 
HoodieLogFile(pathCaptor.getAllValues().get(0)).getLogVersion());
+      assertEquals(3, new 
HoodieLogFile(pathCaptor.getAllValues().get(1)).getLogVersion());
+
+      Map<String, Long> stats = writer.getCDCWriteStats();
+      assertEquals(2, stats.size());
+      assertTrue(stats.keySet().stream().allMatch(path -> 
path.startsWith("partition/")));
+      assertTrue(stats.values().stream().allMatch(size -> size == 17L));
+
+      writer.close();
+      verify(secondFileWriter).close();
+    }
+  }
+
+  @Test
+  public void testUsesFileNameForNonPartitionedTableAndWrapsStatsFailure() 
throws Exception {
+    HoodieStorage storage = mock(HoodieStorage.class);
+    HoodieWriteConfig config = mock(HoodieWriteConfig.class);
+    HoodieSchema cdcSchema = mock(HoodieSchema.class);
+    HoodieFileWriter fileWriter = mock(HoodieFileWriter.class);
+    when(config.getProps()).thenReturn(new TypedProperties());
+    when(cdcSchema.toString()).thenReturn("cdc-schema");
+    when(storage.exists(any(StoragePath.class))).thenReturn(false);
+
+    try (MockedStatic<HoodieFileWriterFactory> writerFactory = 
mockStatic(HoodieFileWriterFactory.class)) {
+      writerFactory.when(() -> HoodieFileWriterFactory.getFileWriter(
+              eq("100"), any(StoragePath.class), eq(storage), eq(config), 
eq(cdcSchema),
+              any(TaskContextSupplier.class), 
eq(HoodieRecord.HoodieRecordType.AVRO)))
+          .thenReturn(fileWriter);
+
+      HoodieNativeCDCFileWriter<IndexedRecord> writer = new 
HoodieNativeCDCFileWriter<>(
+          "100", "", storage, config, cdcSchema, HoodieFileFormat.PARQUET,
+          new StoragePath("/tmp"), "file-1", "1-0-1", new 
LogFileCreationCallback() {
+          }, mock(TaskContextSupplier.class), 
HoodieRecord.HoodieRecordType.AVRO);
+      writer.write("key-1", mock(IndexedRecord.class));
+
+      when(storage.getPathInfo(any(StoragePath.class))).thenAnswer(invocation 
->
+          new StoragePathInfo(invocation.getArgument(0), 11L, false, (short) 
1, 1L, 1L));
+      Map<String, Long> stats = writer.getCDCWriteStats();
+      assertEquals(1, stats.size());
+      assertTrue(stats.keySet().stream().noneMatch(path -> 
path.contains("/")));
+
+      when(storage.getPathInfo(any(StoragePath.class))).thenThrow(new 
java.io.IOException("failure"));
+      assertThrows(HoodieUpsertException.class, writer::getCDCWriteStats);
+    }
+  }
 }
diff --git 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/cdc/TestHoodieNativeCDCLoggers.java
 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/cdc/TestHoodieNativeCDCLoggers.java
new file mode 100644
index 000000000000..294afb9e6868
--- /dev/null
+++ 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/cdc/TestHoodieNativeCDCLoggers.java
@@ -0,0 +1,219 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hudi.io.cdc;
+
+import org.apache.hudi.common.avro.AvroRecordContext;
+import org.apache.hudi.common.engine.TaskContextSupplier;
+import org.apache.hudi.common.model.HoodieFileFormat;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.schema.HoodieSchema;
+import org.apache.hudi.common.schema.HoodieSchemaUtils;
+import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.table.cdc.HoodieCDCOperation;
+import org.apache.hudi.common.table.cdc.HoodieCDCSupplementalLoggingMode;
+import org.apache.hudi.common.table.log.LogFileCreationCallback;
+import org.apache.hudi.common.table.read.BufferedRecord;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.StoragePath;
+
+import org.apache.avro.Schema;
+import org.apache.avro.generic.GenericData;
+import org.apache.avro.generic.GenericRecord;
+import org.apache.avro.generic.IndexedRecord;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.EnumSource;
+import org.mockito.ArgumentCaptor;
+import org.mockito.MockedConstruction;
+
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+
+import static org.apache.hudi.common.table.cdc.HoodieCDCUtils.CDC_AFTER_IMAGE;
+import static org.apache.hudi.common.table.cdc.HoodieCDCUtils.CDC_BEFORE_IMAGE;
+import static 
org.apache.hudi.common.table.cdc.HoodieCDCUtils.CDC_COMMIT_TIMESTAMP;
+import static 
org.apache.hudi.common.table.cdc.HoodieCDCUtils.CDC_OPERATION_TYPE;
+import static org.apache.hudi.common.table.cdc.HoodieCDCUtils.CDC_RECORD_KEY;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+public class TestHoodieNativeCDCLoggers {
+
+  private static final String COMMIT_TIME = "100";
+  private static final HoodieSchema DATA_SCHEMA = 
HoodieSchema.fromAvroSchema(new Schema.Parser().parse(
+      
"{\"type\":\"record\",\"name\":\"trip\",\"namespace\":\"org.apache.hudi\",\"fields\":["
+          + 
"{\"name\":\"id\",\"type\":\"string\"},{\"name\":\"name\",\"type\":\"string\"}]}"));
+  private static final HoodieSchema SCHEMA_WITH_METADATA = 
HoodieSchemaUtils.addMetadataFields(DATA_SCHEMA);
+
+  @ParameterizedTest
+  @EnumSource(HoodieCDCSupplementalLoggingMode.class)
+  public void testAvroLoggerWritesAllOperationsAndSupportsRetraction(
+      HoodieCDCSupplementalLoggingMode mode) throws Exception {
+    try (MockedConstruction<HoodieNativeCDCFileWriter> mockedWriters = 
mockConstruction(HoodieNativeCDCFileWriter.class)) {
+      HoodieAvroNativeCDCLogger logger = new HoodieAvroNativeCDCLogger(
+          COMMIT_TIME, mock(HoodieWriteConfig.class), tableConfig(mode), 
"partition",
+          mock(HoodieStorage.class), SCHEMA_WITH_METADATA, new 
StoragePath("/tmp/partition"),
+          "file-1", "1-0-1", new LogFileCreationCallback() {
+          }, mock(TaskContextSupplier.class));
+      HoodieNativeCDCFileWriter<IndexedRecord> writer = 
mockedWriters.constructed().get(0);
+      Map<String, Long> expectedStats = 
Collections.singletonMap("partition/file.cdc.parquet", 10L);
+      when(writer.getCDCWriteStats()).thenReturn(expectedStats);
+
+      GenericRecord oldRecord = record("id-1", "old");
+      GenericRecord newRecord = record("id-1", "new");
+      GenericRecord insertRecord = record("id-2", "insert");
+      logger.put("id-1", oldRecord, Option.of(newRecord));
+      verify(writer, never()).write(anyString(), any());
+      logger.put("id-2", null, Option.of(insertRecord));
+      logger.put("id-3", oldRecord, Option.empty());
+      logger.put("id-4", null, Option.of(record("id-4", "retracted")));
+      logger.remove("different-key");
+      logger.remove("id-4");
+      logger.close();
+
+      assertSame(expectedStats, logger.getCDCWriteStats());
+      verify(writer, times(3)).write(anyString(), any());
+      verify(writer).close();
+
+      ArgumentCaptor<IndexedRecord> recordCaptor = 
ArgumentCaptor.forClass(IndexedRecord.class);
+      verify(writer, times(3)).write(anyString(), recordCaptor.capture());
+      List<IndexedRecord> cdcRecords = recordCaptor.getAllValues();
+      assertCDCRecord((GenericRecord) cdcRecords.get(0), mode, 
HoodieCDCOperation.UPDATE, "id-1",
+          oldRecord, newRecord);
+      assertCDCRecord((GenericRecord) cdcRecords.get(1), mode, 
HoodieCDCOperation.INSERT, "id-2",
+          null, insertRecord);
+      assertCDCRecord((GenericRecord) cdcRecords.get(2), mode, 
HoodieCDCOperation.DELETE, "id-3",
+          oldRecord, null);
+    }
+  }
+
+  @ParameterizedTest
+  @EnumSource(HoodieCDCSupplementalLoggingMode.class)
+  public void testEngineNativeLoggerWritesAllOperationsAndProjectsMetadata(
+      HoodieCDCSupplementalLoggingMode mode) throws Exception {
+    AvroRecordContext recordContext = new AvroRecordContext();
+    try (MockedConstruction<HoodieNativeCDCFileWriter> mockedWriters = 
mockConstruction(HoodieNativeCDCFileWriter.class)) {
+      HoodieNativeCDCLogger<IndexedRecord> logger = new 
HoodieNativeCDCLogger<>(
+          COMMIT_TIME, mock(HoodieWriteConfig.class), tableConfig(mode), 
"partition",
+          mock(HoodieStorage.class), SCHEMA_WITH_METADATA, new 
StoragePath("/tmp/partition"),
+          "file-1", "1-0-1", new LogFileCreationCallback() {
+          }, mock(TaskContextSupplier.class), recordContext, 
HoodieRecord.HoodieRecordType.AVRO);
+      HoodieNativeCDCFileWriter<IndexedRecord> writer = 
mockedWriters.constructed().get(0);
+
+      GenericRecord oldRecord = record("id-1", "old");
+      GenericRecord newRecord = record("id-1", "new");
+      GenericRecord insertRecord = record("id-2", "insert");
+      BufferedRecord<IndexedRecord> oldBufferedRecord = 
bufferedRecord(oldRecord);
+      BufferedRecord<IndexedRecord> newBufferedRecord = 
bufferedRecord(newRecord);
+      BufferedRecord<IndexedRecord> insertBufferedRecord = 
bufferedRecord(insertRecord);
+      logger.put("id-1", oldBufferedRecord, Option.of(newBufferedRecord));
+      logger.put("id-2", null, Option.of(insertBufferedRecord));
+      logger.put("id-3", oldBufferedRecord, Option.empty());
+      logger.put("id-4", null, Option.of(bufferedRecord(record("id-4", 
"retracted"))));
+      logger.remove("id-4");
+      logger.close();
+
+      ArgumentCaptor<IndexedRecord> recordCaptor = 
ArgumentCaptor.forClass(IndexedRecord.class);
+      verify(writer, times(3)).write(anyString(), recordCaptor.capture());
+      List<IndexedRecord> cdcRecords = recordCaptor.getAllValues();
+      assertCDCRecord((GenericRecord) cdcRecords.get(0), mode, 
HoodieCDCOperation.UPDATE, "id-1",
+          oldRecord, newRecord);
+      assertCDCRecord((GenericRecord) cdcRecords.get(1), mode, 
HoodieCDCOperation.INSERT, "id-2",
+          null, insertRecord);
+      assertCDCRecord((GenericRecord) cdcRecords.get(2), mode, 
HoodieCDCOperation.DELETE, "id-3",
+          oldRecord, null);
+      verify(writer).close();
+    }
+  }
+
+  private static HoodieTableConfig 
tableConfig(HoodieCDCSupplementalLoggingMode mode) {
+    HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
+    when(tableConfig.cdcSupplementalLoggingMode()).thenReturn(mode);
+    when(tableConfig.getBaseFileFormat()).thenReturn(HoodieFileFormat.PARQUET);
+    return tableConfig;
+  }
+
+  private static GenericRecord record(String id, String name) {
+    GenericRecord record = new 
GenericData.Record(SCHEMA_WITH_METADATA.toAvroSchema());
+    record.put(HoodieRecord.COMMIT_TIME_METADATA_FIELD, "001");
+    record.put(HoodieRecord.COMMIT_SEQNO_METADATA_FIELD, "001_0_1");
+    record.put(HoodieRecord.RECORD_KEY_METADATA_FIELD, id);
+    record.put(HoodieRecord.PARTITION_PATH_METADATA_FIELD, "partition");
+    record.put(HoodieRecord.FILENAME_METADATA_FIELD, "file.parquet");
+    record.put("id", id);
+    record.put("name", name);
+    return record;
+  }
+
+  private static BufferedRecord<IndexedRecord> bufferedRecord(GenericRecord 
record) {
+    return new BufferedRecord<>(
+        record.get(HoodieRecord.RECORD_KEY_METADATA_FIELD).toString(), null, 
record, 0, null);
+  }
+
+  private static void assertCDCRecord(
+      GenericRecord record,
+      HoodieCDCSupplementalLoggingMode mode,
+      HoodieCDCOperation operation,
+      String recordKey,
+      GenericRecord expectedBefore,
+      GenericRecord expectedAfter) {
+    assertEquals(operation.getValue(), 
record.get(CDC_OPERATION_TYPE).toString());
+    if (mode == HoodieCDCSupplementalLoggingMode.DATA_BEFORE_AFTER) {
+      assertEquals(COMMIT_TIME, record.get(CDC_COMMIT_TIMESTAMP).toString());
+      assertImage(record, CDC_BEFORE_IMAGE, expectedBefore);
+      assertImage(record, CDC_AFTER_IMAGE, expectedAfter);
+    } else {
+      assertEquals(recordKey, record.get(CDC_RECORD_KEY).toString());
+      if (mode == HoodieCDCSupplementalLoggingMode.DATA_BEFORE) {
+        assertImage(record, CDC_BEFORE_IMAGE, expectedBefore);
+        assertNull(record.getSchema().getField(CDC_AFTER_IMAGE));
+      } else {
+        assertNull(record.getSchema().getField(CDC_BEFORE_IMAGE));
+        assertNull(record.getSchema().getField(CDC_AFTER_IMAGE));
+      }
+    }
+  }
+
+  private static void assertImage(GenericRecord cdcRecord, String imageField, 
GenericRecord expectedImage) {
+    assertNotNull(cdcRecord.getSchema().getField(imageField));
+    GenericRecord actualImage = (GenericRecord) cdcRecord.get(imageField);
+    if (expectedImage == null) {
+      assertNull(actualImage);
+    } else {
+      assertNotNull(actualImage);
+      assertEquals(expectedImage.get("id").toString(), 
actualImage.get("id").toString());
+      assertEquals(expectedImage.get("name").toString(), 
actualImage.get("name").toString());
+      assertFalse(actualImage.getSchema().getFields().stream()
+          .anyMatch(field -> 
HoodieRecord.HOODIE_META_COLUMNS.contains(field.name())));
+    }
+  }
+}

Reply via email to