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 5e1b89fd490c feat: sort input for lsm write (#19079)
5e1b89fd490c is described below

commit 5e1b89fd490c70f5f23989b4736caa07a20234d4
Author: Danny Chan <[email protected]>
AuthorDate: Mon Jun 29 19:04:59 2026 +0800

    feat: sort input for lsm write (#19079)
    
    * feat: sort input for lsm write
---
 .../hudi/io/FileGroupReaderBasedMergeHandle.java   |  91 +++++++--
 .../apache/hudi/io/HoodieAvroNativeCDCLogger.java  | 187 +++++++++++++++++++
 .../org/apache/hudi/io/HoodieCDCLogWriter.java     |  44 +++++
 .../apache/hudi/io/HoodieCDCLogWriterFactory.java  |  84 +++++++++
 .../java/org/apache/hudi/io/HoodieCDCLogger.java   |  15 +-
 .../apache/hudi/io/HoodieMergeHandleFactory.java   |  27 +--
 .../hudi/io/HoodieMergeHandleWithChangeLog.java    |  32 ++--
 .../apache/hudi/io/HoodieNativeCDCFileWriter.java  | 154 +++++++++++++++
 .../org/apache/hudi/io/HoodieNativeCDCLogger.java  | 207 +++++++++++++++++++++
 .../hudi/io/HoodieNativeLogAppendHandle.java       |   1 +
 .../hudi/io/HoodieNativeLogFormatWriter.java       |   5 +-
 .../io/LsmFileGroupReaderBasedMergeHandle.java     | 100 ++++++++++
 .../java/org/apache/hudi/table/HoodieTable.java    |   3 +-
 .../hudi/io/TestHoodieMergeHandleFactory.java      |  40 +++-
 .../FlinkIncrementalMergeHandleWithChangeLog.java  |  23 ++-
 ...FileGroupReaderBasedIncrementalMergeHandle.java |  77 ++++++++
 .../FlinkLsmFileGroupReaderBasedMergeHandle.java   | 118 ++++++++++++
 .../hudi/io/FlinkMergeHandleWithChangeLog.java     |  23 ++-
 .../apache/hudi/io/FlinkWriteHandleFactory.java    |  35 +++-
 .../hudi/table/action/commit/FlinkWriteHelper.java |   4 +-
 .../table/action/commit/TestFlinkWriteHelper.java  |  65 +++++++
 .../BaseJavaDeltaCommitActionExecutor.java         |   8 +
 ...JavaUpsertPreppedDeltaCommitActionExecutor.java |   7 +-
 .../apache/hudi/common/table/read/InputSplit.java  |   4 +
 .../table/read/lsm/HoodieLsmFileGroupReader.java   |   5 +-
 .../table/read/lsm/LsmFileGroupRecordIterator.java |  51 ++++-
 .../apache/hudi/common/util/HoodieRecordUtils.java |  15 +-
 .../hudi/common/util/TestHoodieRecordUtils.java    |  28 ++-
 .../apache/hudi/configuration/OptionsResolver.java |  10 +
 .../org/apache/hudi/sink/StreamWriteFunction.java  |  54 +++++-
 .../org/apache/hudi/sink/buffer/RowDataBucket.java |   5 +
 .../org/apache/hudi/table/HoodieTableFactory.java  |   4 +
 32 files changed, 1442 insertions(+), 84 deletions(-)

diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedMergeHandle.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedMergeHandle.java
index 2893913a0d13..530734172ba6 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedMergeHandle.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedMergeHandle.java
@@ -44,6 +44,7 @@ import 
org.apache.hudi.common.table.read.BaseFileUpdateCallback;
 import org.apache.hudi.common.table.read.BufferedRecord;
 import org.apache.hudi.common.table.read.HoodieFileGroupReader;
 import org.apache.hudi.common.table.read.HoodieReadStats;
+import org.apache.hudi.common.table.read.HoodieRecordReader;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.collection.ClosableIterator;
 import org.apache.hudi.config.HoodieWriteConfig;
@@ -60,6 +61,7 @@ import 
org.apache.hudi.table.action.compact.strategy.CompactionStrategy;
 
 import lombok.extern.slf4j.Slf4j;
 import org.apache.avro.generic.GenericRecord;
+import org.apache.avro.generic.IndexedRecord;
 
 import javax.annotation.concurrent.NotThreadSafe;
 
@@ -87,10 +89,10 @@ import static 
org.apache.hudi.common.model.HoodieFileFormat.HFILE;
 public class FileGroupReaderBasedMergeHandle<T, I, K, O> extends 
HoodieWriteMergeHandle<T, I, K, O> {
 
   private final Option<CompactionOperation> compactionOperation;
-  private final String maxInstantTime;
+  protected final String maxInstantTime;
   private HoodieReadStats readStats;
   private HoodieRecord.HoodieRecordType recordType;
-  private Option<HoodieCDCLogger> cdcLogger;
+  private Option<HoodieCDCLogWriter<?>> cdcLogger;
   private final TypedProperties props;
   private final Iterator<HoodieRecord<T>> incomingRecordsItr;
 
@@ -172,18 +174,38 @@ public class FileGroupReaderBasedMergeHandle<T, I, K, O> 
extends HoodieWriteMerg
     // If the table is a metadata table or the base file is an HFile, we use 
AVRO record type, otherwise we use the engine record type.
     this.recordType = hoodieTable.isMetadataTable() || 
HFILE.getFileExtension().equals(hoodieTable.getBaseFileExtension()) ? 
HoodieRecord.HoodieRecordType.AVRO : enginRecordType;
     if (hoodieTable.getMetaClient().getTableConfig().isCDCEnabled()) {
-      this.cdcLogger = Option.of(new HoodieCDCLogger(
+      this.cdcLogger = Option.of(createCDCLogWriter());
+    } else {
+      this.cdcLogger = Option.empty();
+    }
+  }
+
+  private HoodieCDCLogWriter<?> createCDCLogWriter() {
+    if 
(HoodieCDCLogWriterFactory.shouldWriteNativeCDCLogs(hoodieTable.getMetaClient().getTableConfig()))
 {
+      return new HoodieNativeCDCLogger(
           instantTime,
           config,
           hoodieTable.getMetaClient().getTableConfig(),
           partitionPath,
           storage,
           getWriterSchema(),
-          createLogWriter(instantTime, HoodieCDCUtils.CDC_LOGFILE_SUFFIX, 
Option.empty()),
-          IOUtils.getMaxMemoryPerPartitionMerge(taskContextSupplier, config)));
-    } else {
-      this.cdcLogger = Option.empty();
+          
FSUtils.constructAbsolutePath(hoodieTable.getMetaClient().getBasePath(), 
partitionPath),
+          fileId,
+          writeToken,
+          getLogCreationCallback(),
+          taskContextSupplier,
+          readerContext.getRecordContext(),
+          recordType);
     }
+    return new HoodieCDCLogger(
+        instantTime,
+        config,
+        hoodieTable.getMetaClient().getTableConfig(),
+        partitionPath,
+        storage,
+        getWriterSchema(),
+        createLogWriter(instantTime, HoodieCDCUtils.CDC_LOGFILE_SUFFIX, 
Option.empty()),
+        IOUtils.getMaxMemoryPerPartitionMerge(taskContextSupplier, config));
   }
 
   private void init(CompactionOperation operation, String partitionPath) {
@@ -265,7 +287,7 @@ public class FileGroupReaderBasedMergeHandle<T, I, K, O> 
extends HoodieWriteMerg
         new HoodieLogFile(new StoragePath(FSUtils.constructAbsolutePath(
             config.getBasePath(), op.getPartitionPath()), logFileName))));
     // Initializes file group reader
-    try (HoodieFileGroupReader<T> fileGroupReader = 
getFileGroupReader(usePosition, internalSchemaOption, props, logFilesStreamOpt, 
incomingRecordsItr)) {
+    try (HoodieRecordReader<T> fileGroupReader = 
getFileGroupReader(usePosition, internalSchemaOption, props, logFilesStreamOpt, 
incomingRecordsItr)) {
       // Reads the records from the file slice
       try (ClosableIterator<HoodieRecord<T>> recordIterator = 
fileGroupReader.getClosableHoodieRecordIterator()) {
         while (recordIterator.hasNext()) {
@@ -316,8 +338,8 @@ public class FileGroupReaderBasedMergeHandle<T, I, K, O> 
extends HoodieWriteMerg
         : IOUtils.getMaxMemoryPerPartitionMerge(taskContextSupplier, config);
   }
 
-  private HoodieFileGroupReader<T> getFileGroupReader(boolean usePosition, 
Option<InternalSchema> internalSchemaOption, TypedProperties props,
-                                                        
Option<Stream<HoodieLogFile>> logFileStreamOpt, Iterator<HoodieRecord<T>> 
incomingRecordsItr) {
+  protected HoodieRecordReader<T> getFileGroupReader(boolean usePosition, 
Option<InternalSchema> internalSchemaOption, TypedProperties props,
+                                                     
Option<Stream<HoodieLogFile>> logFileStreamOpt, Iterator<HoodieRecord<T>> 
incomingRecordsItr) {
     HoodieFileGroupReader.HoodieFileGroupReaderBuilder<T> fileGroupBuilder = 
HoodieFileGroupReader.<T>builder().withReaderContext(readerContext).withHoodieTableMetaClient(hoodieTable.getMetaClient())
         
.withLatestCommitTime(maxInstantTime).withPartitionPath(partitionPath).withBaseFileOption(Option.ofNullable(baseFileToMerge))
         
.withDataSchema(writeSchemaWithMetaFields).withRequestedSchema(writeSchemaWithMetaFields)
@@ -366,11 +388,20 @@ public class FileGroupReaderBasedMergeHandle<T, I, K, O> 
extends HoodieWriteMerg
     }
   }
 
-  private Option<BaseFileUpdateCallback<T>> createCallback() {
+  protected Option<BaseFileUpdateCallback<T>> createCallback() {
     List<BaseFileUpdateCallback<T>> callbacks = new ArrayList<>();
     // Handle CDC workflow.
     if (cdcLogger.isPresent()) {
-      callbacks.add(new CDCCallback<>(cdcLogger.get(), readerContext));
+      HoodieCDCLogWriter<?> logger = cdcLogger.get();
+      if (logger instanceof HoodieNativeCDCLogger) {
+        @SuppressWarnings("unchecked")
+        HoodieNativeCDCLogger<T> nativeCDCLogger = (HoodieNativeCDCLogger<T>) 
logger;
+        callbacks.add(new NativeCDCCallback<>(nativeCDCLogger));
+      } else {
+        @SuppressWarnings("unchecked")
+        HoodieCDCLogWriter<IndexedRecord> inlineCDCLogger = 
(HoodieCDCLogWriter<IndexedRecord>) logger;
+        callbacks.add(new CDCCallback<>(inlineCDCLogger, readerContext));
+      }
     }
     // Indexes are not updated during compaction
     if (compactionOperation.isEmpty()) {
@@ -394,10 +425,10 @@ public class FileGroupReaderBasedMergeHandle<T, I, K, O> 
extends HoodieWriteMerg
   }
 
   private static class CDCCallback<T> implements BaseFileUpdateCallback<T> {
-    private final HoodieCDCLogger cdcLogger;
+    private final HoodieCDCLogWriter<IndexedRecord> cdcLogger;
     private final RecordContext<T> recordContext;
 
-    CDCCallback(HoodieCDCLogger cdcLogger, HoodieReaderContext<T> 
readerContext) {
+    CDCCallback(HoodieCDCLogWriter<IndexedRecord> cdcLogger, 
HoodieReaderContext<T> readerContext) {
       this.cdcLogger = cdcLogger;
       this.recordContext = readerContext.getRecordContext();
     }
@@ -434,6 +465,38 @@ public class FileGroupReaderBasedMergeHandle<T, I, K, O> 
extends HoodieWriteMerg
     }
   }
 
+  private static class NativeCDCCallback<T> implements 
BaseFileUpdateCallback<T> {
+    private final HoodieNativeCDCLogger<T> cdcLogger;
+
+    NativeCDCCallback(HoodieNativeCDCLogger<T> cdcLogger) {
+      this.cdcLogger = cdcLogger;
+    }
+
+    @Override
+    public void onUpdate(String recordKey, BufferedRecord<T> previousRecord, 
BufferedRecord<T> mergedRecord) {
+      cdcLogger.put(recordKey, previousRecord, Option.of(mergedRecord));
+    }
+
+    @Override
+    public void onInsert(String recordKey, BufferedRecord<T> newRecord) {
+      cdcLogger.put(recordKey, null, Option.of(newRecord));
+    }
+
+    @Override
+    public void onDelete(String recordKey, BufferedRecord<T> previousRecord, 
HoodieOperation hoodieOperation) {
+      // delete record from log block and update no base record from base 
file, skip generating changelog.
+      if (previousRecord == null) {
+        return;
+      }
+      cdcLogger.put(recordKey, previousRecord, Option.empty());
+    }
+
+    @Override
+    public void onFailure(String recordKey) {
+      cdcLogger.remove(recordKey);
+    }
+  }
+
   private static class RecordLevelIndexCallback<T> implements 
BaseFileUpdateCallback<T> {
     private final WriteStatus writeStatus;
     private final HoodieRecordLocation fileRecordLocation;
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAvroNativeCDCLogger.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAvroNativeCDCLogger.java
new file mode 100644
index 000000000000..be83c64a9401
--- /dev/null
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAvroNativeCDCLogger.java
@@ -0,0 +1,187 @@
+/*
+ * 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.avro.HoodieAvroUtils;
+import org.apache.hudi.common.engine.TaskContextSupplier;
+import org.apache.hudi.common.model.HoodieAvroIndexedRecord;
+import org.apache.hudi.common.model.HoodieKey;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.schema.HoodieSchema;
+import org.apache.hudi.common.schema.HoodieSchemaCache;
+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.cdc.HoodieCDCUtils;
+import org.apache.hudi.common.table.log.LogFileCreationCallback;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.exception.HoodieException;
+import org.apache.hudi.exception.HoodieIOException;
+import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.StoragePath;
+
+import org.apache.avro.generic.GenericData;
+import org.apache.avro.generic.GenericRecord;
+import org.apache.avro.generic.IndexedRecord;
+
+import java.io.IOException;
+import java.util.Map;
+
+import static 
org.apache.hudi.common.table.cdc.HoodieCDCSupplementalLoggingMode.DATA_BEFORE;
+import static 
org.apache.hudi.common.table.cdc.HoodieCDCSupplementalLoggingMode.DATA_BEFORE_AFTER;
+
+/**
+ * Writes CDC records as native CDC log files from Avro input records.
+ */
+public class HoodieAvroNativeCDCLogger implements 
HoodieCDCLogWriter<IndexedRecord> {
+
+  private final String commitTime;
+  private final String partitionPath;
+  private final HoodieSchema dataSchema;
+  private final HoodieSchema cdcSchema;
+  private final HoodieCDCSupplementalLoggingMode cdcSupplementalLoggingMode;
+  private final CDCTransformer transformer;
+  private final HoodieNativeCDCFileWriter nativeCDCFileWriter;
+  private PendingCDCRecord pendingRecord;
+
+  public HoodieAvroNativeCDCLogger(
+      String commitTime,
+      HoodieWriteConfig config,
+      HoodieTableConfig tableConfig,
+      String partitionPath,
+      HoodieStorage storage,
+      HoodieSchema schema,
+      StoragePath parentPath,
+      String fileId,
+      String writeToken,
+      LogFileCreationCallback fileCreationCallback,
+      TaskContextSupplier taskContextSupplier) {
+    this.commitTime = commitTime;
+    this.partitionPath = partitionPath;
+    this.dataSchema = 
HoodieSchemaCache.intern(HoodieSchemaUtils.removeMetadataFields(schema));
+    this.cdcSupplementalLoggingMode = tableConfig.cdcSupplementalLoggingMode();
+    this.cdcSchema = 
HoodieCDCUtils.schemaBySupplementalLoggingMode(cdcSupplementalLoggingMode, 
dataSchema);
+    this.transformer = getTransformer();
+    this.nativeCDCFileWriter = new HoodieNativeCDCFileWriter(
+        commitTime,
+        partitionPath,
+        storage,
+        config,
+        cdcSchema,
+        tableConfig.getBaseFileFormat(),
+        parentPath,
+        fileId,
+        writeToken,
+        fileCreationCallback,
+        taskContextSupplier,
+        HoodieRecord.HoodieRecordType.AVRO);
+  }
+
+  @Override
+  public void put(String recordKey, IndexedRecord oldRecord, 
Option<IndexedRecord> newRecord) {
+    GenericData.Record cdcRecord;
+    if (newRecord.isPresent()) {
+      if (oldRecord == null) {
+        cdcRecord = transformer.transform(HoodieCDCOperation.INSERT, 
recordKey, null, (GenericRecord) newRecord.get());
+      } else {
+        cdcRecord = transformer.transform(HoodieCDCOperation.UPDATE, 
recordKey, (GenericRecord) oldRecord, (GenericRecord) newRecord.get());
+      }
+    } else {
+      cdcRecord = transformer.transform(HoodieCDCOperation.DELETE, recordKey, 
(GenericRecord) oldRecord, null);
+    }
+
+    flushPendingRecord();
+    pendingRecord = new PendingCDCRecord(recordKey, cdcRecord);
+  }
+
+  @Override
+  public void remove(String recordKey) {
+    if (pendingRecord != null && pendingRecord.recordKey.equals(recordKey)) {
+      pendingRecord = null;
+    }
+  }
+
+  @Override
+  public Map<String, Long> getCDCWriteStats() {
+    return nativeCDCFileWriter.getCDCWriteStats();
+  }
+
+  @Override
+  public void close() {
+    try {
+      flushPendingRecord();
+      nativeCDCFileWriter.close();
+    } catch (IOException e) {
+      throw new HoodieIOException("Failed to close HoodieAvroNativeCDCLogger", 
e);
+    }
+  }
+
+  private void flushPendingRecord() {
+    if (pendingRecord == null) {
+      return;
+    }
+    try {
+      nativeCDCFileWriter.write(
+          pendingRecord.recordKey,
+          new HoodieAvroIndexedRecord(new HoodieKey(pendingRecord.recordKey, 
partitionPath), pendingRecord.record));
+      pendingRecord = null;
+    } catch (IOException e) {
+      throw new HoodieException("Failed to write the cdc data to native cdc 
log file", e);
+    }
+  }
+
+  private CDCTransformer getTransformer() {
+    if (cdcSupplementalLoggingMode == DATA_BEFORE_AFTER) {
+      return (operation, recordKey, oldRecord, newRecord) ->
+          HoodieCDCUtils.cdcRecord(cdcSchema, operation.getValue(), 
commitTime, removeCommitMetadata(oldRecord), removeCommitMetadata(newRecord));
+    } else if (cdcSupplementalLoggingMode == DATA_BEFORE) {
+      return (operation, recordKey, oldRecord, newRecord) ->
+          HoodieCDCUtils.cdcRecord(cdcSchema, operation.getValue(), recordKey, 
removeCommitMetadata(oldRecord));
+    } else {
+      return (operation, recordKey, oldRecord, newRecord) ->
+          HoodieCDCUtils.cdcRecord(cdcSchema, operation.getValue(), recordKey);
+    }
+  }
+
+  private GenericRecord removeCommitMetadata(GenericRecord record) {
+    return record == null ? null : 
HoodieAvroUtils.projectRecordToNewSchemaShallow(record, 
dataSchema.getAvroSchema());
+  }
+
+  private static class PendingCDCRecord {
+    private final String recordKey;
+    private final IndexedRecord record;
+
+    private PendingCDCRecord(String recordKey, IndexedRecord record) {
+      this.recordKey = recordKey;
+      this.record = record;
+    }
+  }
+
+  /**
+   * A transformer that transforms normal Avro records into CDC records.
+   */
+  private interface CDCTransformer {
+    GenericData.Record transform(HoodieCDCOperation operation,
+                                 String recordKey,
+                                 GenericRecord oldRecord,
+                                 GenericRecord newRecord);
+  }
+}
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCDCLogWriter.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCDCLogWriter.java
new file mode 100644
index 000000000000..9adf9afbd9a7
--- /dev/null
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCDCLogWriter.java
@@ -0,0 +1,44 @@
+/*
+ * 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.common.model.HoodieRecord;
+import org.apache.hudi.common.util.Option;
+
+import java.io.Closeable;
+import java.util.Map;
+
+/**
+ * Writes CDC records generated by merge handles.
+ */
+public interface HoodieCDCLogWriter<T> extends Closeable {
+
+  default void put(HoodieRecord hoodieRecord, T oldRecord, Option<T> 
newRecord) {
+    put(hoodieRecord.getRecordKey(), oldRecord, newRecord);
+  }
+
+  void put(String recordKey, T oldRecord, Option<T> newRecord);
+
+  void remove(String recordKey);
+
+  Map<String, Long> getCDCWriteStats();
+
+  @Override
+  void close();
+}
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCDCLogWriterFactory.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCDCLogWriterFactory.java
new file mode 100644
index 000000000000..ab509f84963e
--- /dev/null
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCDCLogWriterFactory.java
@@ -0,0 +1,84 @@
+/*
+ * 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.common.engine.TaskContextSupplier;
+import org.apache.hudi.common.fs.FSUtils;
+import org.apache.hudi.common.schema.HoodieSchema;
+import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.table.log.HoodieLogFormat;
+import org.apache.hudi.common.table.log.LogFileCreationCallback;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.table.HoodieTable;
+
+import org.apache.avro.generic.IndexedRecord;
+
+import java.util.function.Supplier;
+
+/**
+ * Creates CDC log writers for Avro merge handles.
+ */
+final class HoodieCDCLogWriterFactory {
+
+  private HoodieCDCLogWriterFactory() {
+  }
+
+  static <T, I, K, O> HoodieCDCLogWriter<IndexedRecord> createAvroCDCLogWriter(
+      String instantTime,
+      HoodieWriteConfig config,
+      HoodieTable<T, I, K, O> hoodieTable,
+      String partitionPath,
+      HoodieStorage storage,
+      HoodieSchema writerSchema,
+      String fileId,
+      String writeToken,
+      LogFileCreationCallback logCreationCallback,
+      TaskContextSupplier taskContextSupplier,
+      Supplier<HoodieLogFormat.Writer> logWriterSupplier) {
+    HoodieTableConfig tableConfig = 
hoodieTable.getMetaClient().getTableConfig();
+    if (shouldWriteNativeCDCLogs(tableConfig)) {
+      return new HoodieAvroNativeCDCLogger(
+          instantTime,
+          config,
+          tableConfig,
+          partitionPath,
+          storage,
+          writerSchema,
+          
FSUtils.constructAbsolutePath(hoodieTable.getMetaClient().getBasePath(), 
partitionPath),
+          fileId,
+          writeToken,
+          logCreationCallback,
+          taskContextSupplier);
+    }
+    return new HoodieCDCLogger(
+        instantTime,
+        config,
+        tableConfig,
+        partitionPath,
+        storage,
+        writerSchema,
+        logWriterSupplier.get(),
+        IOUtils.getMaxMemoryPerPartitionMerge(taskContextSupplier, config));
+  }
+
+  static boolean shouldWriteNativeCDCLogs(HoodieTableConfig tableConfig) {
+    return tableConfig.isLSMTreeStorageLayout();
+  }
+}
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCDCLogger.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCDCLogger.java
index 0d9b74a7f2dc..30de52c01c9b 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCDCLogger.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieCDCLogger.java
@@ -49,7 +49,6 @@ import org.apache.avro.generic.GenericData;
 import org.apache.avro.generic.GenericRecord;
 import org.apache.avro.generic.IndexedRecord;
 
-import java.io.Closeable;
 import java.io.IOException;
 import java.util.ArrayList;
 import java.util.Collections;
@@ -65,7 +64,7 @@ import static 
org.apache.hudi.common.table.cdc.HoodieCDCSupplementalLoggingMode.
 /**
  * This class encapsulates all the cdc-writing functions.
  */
-public class HoodieCDCLogger implements Closeable {
+public class HoodieCDCLogger implements HoodieCDCLogWriter<IndexedRecord> {
 
   private final String commitTime;
 
@@ -152,14 +151,8 @@ public class HoodieCDCLogger implements Closeable {
     }
   }
 
-  public void put(HoodieRecord hoodieRecord,
-                  GenericRecord oldRecord,
-                  Option<IndexedRecord> newRecord) {
-    put(hoodieRecord.getRecordKey(), oldRecord, newRecord);
-  }
-
   public void put(String recordKey,
-                  GenericRecord oldRecord,
+                  IndexedRecord oldRecord,
                   Option<IndexedRecord> newRecord) {
     GenericData.Record cdcRecord;
     if (newRecord.isPresent()) {
@@ -171,12 +164,12 @@ public class HoodieCDCLogger implements Closeable {
       } else {
         // UPDATE cdc record
         cdcRecord = this.transformer.transform(HoodieCDCOperation.UPDATE, 
recordKey,
-            oldRecord, record);
+            (GenericRecord) oldRecord, record);
       }
     } else {
       // DELETE cdc record
       cdcRecord = this.transformer.transform(HoodieCDCOperation.DELETE, 
recordKey,
-          oldRecord, null);
+          (GenericRecord) oldRecord, null);
     }
 
     flushIfNeeded(false);
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandleFactory.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandleFactory.java
index 01ad6a2eeaf8..6f7ef0bb8c96 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandleFactory.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandleFactory.java
@@ -118,9 +118,11 @@ public class HoodieMergeHandleFactory {
       String maxInstantTime,
       HoodieRecord.HoodieRecordType recordType) {
 
-    boolean isFallbackEnabled = config.isMergeHandleFallbackEnabled();
-
-    String mergeHandleClass = config.getCompactionMergeHandleClassName();
+    String mergeHandleClass = 
hoodieTable.getMetaClient().getTableConfig().isLSMTreeStorageLayout()
+        ? LsmFileGroupReaderBasedMergeHandle.class.getName()
+        : config.getCompactionMergeHandleClassName();
+    boolean isFallbackEnabled = config.isMergeHandleFallbackEnabled()
+        && 
!LsmFileGroupReaderBasedMergeHandle.class.getName().equals(mergeHandleClass);
     String logContext = String.format("for fileId %s and partitionPath %s at 
commit %s", operation.getFileId(), operation.getPartitionPath(), instantTime);
     log.info("Create HoodieMergeHandle implementation {} {}", 
mergeHandleClass, logContext);
 
@@ -165,26 +167,22 @@ public class HoodieMergeHandleFactory {
     String mergeHandleClass;
     String fallbackMergeHandleClass = null;
 
-    // Overwrite to a different implementation for {@link 
HoodieWriteMergeHandle} if sorting or CDC is enabled.
-    if (table.requireSortedRecords()) {
-      if (table.getMetaClient().getTableConfig().isCDCEnabled()) {
-        mergeHandleClass = 
HoodieSortedMergeHandleWithChangeLog.class.getName();
-      } else {
-        mergeHandleClass = HoodieSortedMergeHandle.class.getName();
-      }
+    // Overwrite to file-group-reader based implementations if sorted output 
is required.
+    if (table.getMetaClient().getTableConfig().isLSMTreeStorageLayout()) {
+      mergeHandleClass = LsmFileGroupReaderBasedMergeHandle.class.getName();
     } else if (!WriteOperationType.isChangingRecords(operationType) && 
writeConfig.allowDuplicateInserts()) {
       mergeHandleClass = writeConfig.getConcatHandleClassName();
       if 
(!mergeHandleClass.equals(HoodieWriteConfig.CONCAT_HANDLE_CLASS_NAME.defaultValue()))
 {
         fallbackMergeHandleClass = 
HoodieWriteConfig.CONCAT_HANDLE_CLASS_NAME.defaultValue();
       }
-    } else if (table.getMetaClient().getTableConfig().isCDCEnabled()) {
+    } else if (table.requireSortedRecords() || 
table.getMetaClient().getTableConfig().isCDCEnabled()) {
       if 
(writeConfig.getMergeHandleClassName().equals(FileGroupReaderBasedMergeHandle.class.getName()))
 {
         mergeHandleClass = writeConfig.getMergeHandleClassName();
         if 
(!mergeHandleClass.equals(HoodieWriteConfig.MERGE_HANDLE_CLASS_NAME.defaultValue()))
 {
           fallbackMergeHandleClass = 
HoodieWriteConfig.MERGE_HANDLE_CLASS_NAME.defaultValue();
         }
       } else {
-        mergeHandleClass = HoodieMergeHandleWithChangeLog.class.getName();
+        mergeHandleClass = FileGroupReaderBasedMergeHandle.class.getName();
       }
     } else {
       mergeHandleClass = writeConfig.getMergeHandleClassName();
@@ -196,6 +194,9 @@ public class HoodieMergeHandleFactory {
     return Pair.of(mergeHandleClass, fallbackMergeHandleClass);
   }
 
+  /**
+   * IMPORTANT: this is only for compaction paths without file group reader.
+   */
   @VisibleForTesting
   static Pair<String, String> 
getMergeHandleClassesCompaction(HoodieWriteConfig writeConfig, HoodieTable 
table) {
     String mergeHandleClass;
@@ -217,4 +218,4 @@ public class HoodieMergeHandleFactory {
 
     return Pair.of(mergeHandleClass, fallbackMergeHandleClass);
   }
-}
\ No newline at end of file
+}
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandleWithChangeLog.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandleWithChangeLog.java
index 10f70868cfd2..6298a11f6865 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandleWithChangeLog.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandleWithChangeLog.java
@@ -47,21 +47,13 @@ import java.util.Map;
 @Slf4j
 public class HoodieMergeHandleWithChangeLog<T, I, K, O> extends 
HoodieWriteMergeHandle<T, I, K, O> {
 
-  protected final HoodieCDCLogger cdcLogger;
+  protected final HoodieCDCLogWriter<IndexedRecord> cdcLogger;
 
   public HoodieMergeHandleWithChangeLog(HoodieWriteConfig config, String 
instantTime, HoodieTable<T, I, K, O> hoodieTable,
                                         Iterator<HoodieRecord<T>> recordItr, 
String partitionPath, String fileId,
                                         TaskContextSupplier 
taskContextSupplier, Option<BaseKeyGenerator> keyGeneratorOpt) {
     super(config, instantTime, hoodieTable, recordItr, partitionPath, fileId, 
taskContextSupplier, keyGeneratorOpt);
-    this.cdcLogger = new HoodieCDCLogger(
-        instantTime,
-        config,
-        hoodieTable.getMetaClient().getTableConfig(),
-        partitionPath,
-        storage,
-        getWriterSchema(),
-        createLogWriter(instantTime, HoodieCDCUtils.CDC_LOGFILE_SUFFIX, 
Option.empty()),
-        IOUtils.getMaxMemoryPerPartitionMerge(taskContextSupplier, config));
+    this.cdcLogger = createCDCLogWriter(instantTime, config, hoodieTable, 
partitionPath, taskContextSupplier);
   }
 
   /**
@@ -71,15 +63,27 @@ public class HoodieMergeHandleWithChangeLog<T, I, K, O> 
extends HoodieWriteMerge
                                         Map<String, HoodieRecord<T>> 
keyToNewRecords, String partitionPath, String fileId,
                                         HoodieBaseFile dataFileToBeMerged, 
TaskContextSupplier taskContextSupplier, Option<BaseKeyGenerator> 
keyGeneratorOpt) {
     super(config, instantTime, hoodieTable, keyToNewRecords, partitionPath, 
fileId, dataFileToBeMerged, taskContextSupplier, keyGeneratorOpt);
-    this.cdcLogger = new HoodieCDCLogger(
+    this.cdcLogger = createCDCLogWriter(instantTime, config, hoodieTable, 
partitionPath, taskContextSupplier);
+  }
+
+  private HoodieCDCLogWriter<IndexedRecord> createCDCLogWriter(
+      String instantTime,
+      HoodieWriteConfig config,
+      HoodieTable<T, I, K, O> hoodieTable,
+      String partitionPath,
+      TaskContextSupplier taskContextSupplier) {
+    return HoodieCDCLogWriterFactory.createAvroCDCLogWriter(
         instantTime,
         config,
-        hoodieTable.getMetaClient().getTableConfig(),
+        hoodieTable,
         partitionPath,
         storage,
         getWriterSchema(),
-        createLogWriter(instantTime, HoodieCDCUtils.CDC_LOGFILE_SUFFIX, 
Option.empty()),
-        IOUtils.getMaxMemoryPerPartitionMerge(taskContextSupplier, config));
+        fileId,
+        writeToken,
+        getLogCreationCallback(),
+        taskContextSupplier,
+        () -> createLogWriter(instantTime, HoodieCDCUtils.CDC_LOGFILE_SUFFIX, 
Option.empty()));
   }
 
   protected boolean writeUpdateRecord(HoodieRecord<T> newRecord, 
HoodieRecord<T> oldRecord, HoodieRecord combinedRecord, HoodieSchema 
writerSchema)
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeCDCFileWriter.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeCDCFileWriter.java
new file mode 100644
index 000000000000..697d5a49a069
--- /dev/null
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeCDCFileWriter.java
@@ -0,0 +1,154 @@
+/*
+ * 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.common.engine.TaskContextSupplier;
+import org.apache.hudi.common.fs.FSUtils;
+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.cdc.HoodieCDCUtils;
+import org.apache.hudi.common.table.log.LogFileCreationCallback;
+import org.apache.hudi.common.util.StringUtils;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.exception.HoodieUpsertException;
+import org.apache.hudi.io.storage.HoodieFileWriter;
+import org.apache.hudi.io.storage.HoodieFileWriterFactory;
+import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.StoragePath;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Properties;
+
+/**
+ * Manages native CDC log file creation, rolling, writes, and stats.
+ */
+class HoodieNativeCDCFileWriter {
+
+  private final String commitTime;
+  private final String partitionPath;
+  private final HoodieStorage storage;
+  private final HoodieWriteConfig config;
+  private final HoodieSchema cdcSchema;
+  private final HoodieFileFormat nativeFileFormat;
+  private final StoragePath parentPath;
+  private final String fileId;
+  private final String writeToken;
+  private final LogFileCreationCallback fileCreationCallback;
+  private final TaskContextSupplier taskContextSupplier;
+  private final HoodieRecord.HoodieRecordType recordType;
+  private final Properties recordProperties;
+  private final List<StoragePath> cdcAbsPaths;
+  private int nextLogVersion;
+  private HoodieFileWriter cdcWriter;
+
+  HoodieNativeCDCFileWriter(
+      String commitTime,
+      String partitionPath,
+      HoodieStorage storage,
+      HoodieWriteConfig config,
+      HoodieSchema cdcSchema,
+      HoodieFileFormat nativeFileFormat,
+      StoragePath parentPath,
+      String fileId,
+      String writeToken,
+      LogFileCreationCallback fileCreationCallback,
+      TaskContextSupplier taskContextSupplier,
+      HoodieRecord.HoodieRecordType recordType) {
+    this.commitTime = commitTime;
+    this.partitionPath = partitionPath;
+    this.storage = storage;
+    this.config = config;
+    this.cdcSchema = cdcSchema;
+    this.nativeFileFormat = nativeFileFormat;
+    this.parentPath = parentPath;
+    this.fileId = fileId;
+    this.writeToken = writeToken;
+    this.fileCreationCallback = fileCreationCallback;
+    this.taskContextSupplier = taskContextSupplier;
+    this.recordType = recordType;
+    this.recordProperties = new Properties();
+    this.recordProperties.putAll(config.getProps());
+    this.cdcAbsPaths = new ArrayList<>();
+    this.nextLogVersion = HoodieLogFile.LOGFILE_BASE_VERSION;
+  }
+
+  void write(String recordKey, HoodieRecord record) throws IOException {
+    ensureCDCWriter();
+    cdcWriter.write(recordKey, record, cdcSchema, recordProperties);
+  }
+
+  Map<String, Long> getCDCWriteStats() {
+    Map<String, Long> stats = new HashMap<>();
+    try {
+      for (StoragePath cdcAbsPath : cdcAbsPaths) {
+        String cdcFileName = cdcAbsPath.getName();
+        String cdcPath = StringUtils.isNullOrEmpty(partitionPath) ? 
cdcFileName : partitionPath + "/" + cdcFileName;
+        stats.put(cdcPath, storage.getPathInfo(cdcAbsPath).getLength());
+      }
+    } catch (IOException e) {
+      throw new HoodieUpsertException("Failed to get cdc write stat", e);
+    }
+    return stats;
+  }
+
+  void close() throws IOException {
+    if (cdcWriter != null) {
+      cdcWriter.close();
+      cdcWriter = null;
+    }
+  }
+
+  private void ensureCDCWriter() throws IOException {
+    if (cdcWriter != null && cdcWriter.canWrite()) {
+      return;
+    }
+    close();
+    HoodieLogFile cdcLogFile = createNativeCDCLogFile();
+    cdcWriter = HoodieFileWriterFactory.getFileWriter(
+        commitTime, cdcLogFile.getPath(), storage, config, cdcSchema, 
taskContextSupplier, recordType);
+    cdcAbsPaths.add(cdcLogFile.getPath());
+  }
+
+  private HoodieLogFile createNativeCDCLogFile() throws IOException {
+    int version = nextAvailableVersion();
+    HoodieLogFile nativeCDCLogFile = new 
HoodieLogFile(makeNativeCDCLogPath(version), 0);
+    fileCreationCallback.preFileCreation(nativeCDCLogFile);
+    nextLogVersion = version + 1;
+    return nativeCDCLogFile;
+  }
+
+  private int nextAvailableVersion() throws IOException {
+    int candidateVersion = nextLogVersion;
+    while (storage.exists(makeNativeCDCLogPath(candidateVersion))) {
+      candidateVersion++;
+    }
+    return candidateVersion;
+  }
+
+  private StoragePath makeNativeCDCLogPath(int version) {
+    return new StoragePath(parentPath, FSUtils.makeNativeLogFileName(
+        fileId, writeToken, commitTime, version, 
HoodieCDCUtils.CDC_LOGFILE_SUFFIX, nativeFileFormat));
+  }
+}
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeCDCLogger.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeCDCLogger.java
new file mode 100644
index 000000000000..cfa667d131d5
--- /dev/null
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeCDCLogger.java
@@ -0,0 +1,207 @@
+/*
+ * 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.common.engine.RecordContext;
+import org.apache.hudi.common.engine.TaskContextSupplier;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.schema.HoodieSchema;
+import org.apache.hudi.common.schema.HoodieSchemaCache;
+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.cdc.HoodieCDCUtils;
+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.exception.HoodieException;
+import org.apache.hudi.exception.HoodieIOException;
+import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.StoragePath;
+
+import java.io.IOException;
+import java.util.Map;
+
+import static 
org.apache.hudi.common.table.cdc.HoodieCDCSupplementalLoggingMode.DATA_BEFORE;
+import static 
org.apache.hudi.common.table.cdc.HoodieCDCSupplementalLoggingMode.DATA_BEFORE_AFTER;
+
+/**
+ * Writes CDC records as native CDC log files, for example {@code 
.cdc.parquet}.
+ */
+public class HoodieNativeCDCLogger<T> implements 
HoodieCDCLogWriter<BufferedRecord<T>> {
+
+  /** Instant time used for naming and writing native CDC log files. */
+  private final String commitTime;
+  /** Partition path whose records are being merged. */
+  private final String partitionPath;
+  /** Table data schema without Hudi metadata fields. */
+  private final HoodieSchema dataSchema;
+  /** CDC record schema derived from the table CDC supplemental logging mode. 
*/
+  private final HoodieSchema cdcSchema;
+  /** Cached encoded CDC schema id reused for every buffered CDC record. */
+  private final Integer cdcSchemaId;
+  /** Supplemental logging mode determining which fields are emitted in CDC 
records. */
+  private final HoodieCDCSupplementalLoggingMode cdcSupplementalLoggingMode;
+  /** Engine-specific record context for row construction, projection, and 
conversion. */
+  private final RecordContext<T> recordContext;
+  /** Manages native CDC file creation, rolling, writing, and stats. */
+  private final HoodieNativeCDCFileWriter nativeCDCFileWriter;
+  /** Last CDC record staged for write so it can be retracted if the merge 
later fails. */
+  private PendingCDCRecord<T> pendingRecord;
+
+  public HoodieNativeCDCLogger(
+      String commitTime,
+      HoodieWriteConfig config,
+      HoodieTableConfig tableConfig,
+      String partitionPath,
+      HoodieStorage storage,
+      HoodieSchema schema,
+      StoragePath parentPath,
+      String fileId,
+      String writeToken,
+      LogFileCreationCallback fileCreationCallback,
+      TaskContextSupplier taskContextSupplier,
+      RecordContext<T> recordContext,
+      HoodieRecord.HoodieRecordType recordType) {
+    this.commitTime = commitTime;
+    this.partitionPath = partitionPath;
+    this.dataSchema = 
HoodieSchemaCache.intern(HoodieSchemaUtils.removeMetadataFields(schema));
+    this.cdcSupplementalLoggingMode = tableConfig.cdcSupplementalLoggingMode();
+    this.cdcSchema = 
HoodieCDCUtils.schemaBySupplementalLoggingMode(cdcSupplementalLoggingMode, 
dataSchema);
+    this.recordContext = recordContext;
+    this.cdcSchemaId = recordContext.encodeSchema(cdcSchema);
+    this.nativeCDCFileWriter = new HoodieNativeCDCFileWriter(
+        commitTime,
+        partitionPath,
+        storage,
+        config,
+        cdcSchema,
+        tableConfig.getBaseFileFormat(),
+        parentPath,
+        fileId,
+        writeToken,
+        fileCreationCallback,
+        taskContextSupplier,
+        recordType);
+  }
+
+  @Override
+  public void put(String recordKey, BufferedRecord<T> oldRecord, 
Option<BufferedRecord<T>> newRecord) {
+    flushPendingRecord();
+    HoodieCDCOperation operation;
+    if (newRecord.isPresent()) {
+      operation = oldRecord == null ? HoodieCDCOperation.INSERT : 
HoodieCDCOperation.UPDATE;
+    } else {
+      operation = HoodieCDCOperation.DELETE;
+    }
+    this.pendingRecord = new PendingCDCRecord<>(recordKey, 
createCDCRecord(recordKey, operation, oldRecord, newRecord.orElse(null)));
+  }
+
+  @Override
+  public void remove(String recordKey) {
+    if (pendingRecord != null && pendingRecord.recordKey.equals(recordKey)) {
+      pendingRecord = null;
+    }
+  }
+
+  @Override
+  public Map<String, Long> getCDCWriteStats() {
+    return nativeCDCFileWriter.getCDCWriteStats();
+  }
+
+  @Override
+  public void close() {
+    try {
+      flushPendingRecord();
+      nativeCDCFileWriter.close();
+    } catch (IOException e) {
+      throw new HoodieIOException("Failed to close HoodieNativeCDCLogger", e);
+    }
+  }
+
+  private BufferedRecord<T> createCDCRecord(
+      String recordKey,
+      HoodieCDCOperation operation,
+      BufferedRecord<T> oldRecord,
+      BufferedRecord<T> newRecord) {
+    Object[] fieldValues = new Object[cdcSchema.getFields().size()];
+    if (cdcSupplementalLoggingMode == DATA_BEFORE_AFTER) {
+      fieldValues[0] = convertString(operation.getValue());
+      fieldValues[1] = convertString(commitTime);
+      fieldValues[2] = projectDataRecord(oldRecord);
+      fieldValues[3] = projectDataRecord(newRecord);
+    } else if (cdcSupplementalLoggingMode == DATA_BEFORE) {
+      fieldValues[0] = convertString(operation.getValue());
+      fieldValues[1] = convertString(recordKey);
+      fieldValues[2] = projectDataRecord(oldRecord);
+    } else {
+      fieldValues[0] = convertString(operation.getValue());
+      fieldValues[1] = convertString(recordKey);
+    }
+    T cdcRecord = recordContext.constructEngineRecord(cdcSchema, fieldValues);
+    return new BufferedRecord<>(recordKey, null, cdcRecord, cdcSchemaId, null);
+  }
+
+  private Object convertString(String value) {
+    return recordContext.convertValueToEngineType(value);
+  }
+
+  private T projectDataRecord(BufferedRecord<T> record) {
+    if (record == null || record.getRecord() == null) {
+      return null;
+    }
+    HoodieSchema recordSchema = 
recordContext.getSchemaFromBufferRecord(record);
+    T dataRecord = record.getRecord();
+    if (needsProjection(recordSchema)) {
+      dataRecord = recordContext.projectRecord(recordSchema, 
dataSchema).apply(dataRecord);
+    }
+    return recordContext.seal(dataRecord);
+  }
+
+  private boolean needsProjection(HoodieSchema recordSchema) {
+    return recordSchema != null
+        && (recordSchema.getFields().size() != dataSchema.getFields().size() 
|| !recordSchema.equals(dataSchema));
+  }
+
+  private void flushPendingRecord() {
+    if (pendingRecord == null) {
+      return;
+    }
+    try {
+      nativeCDCFileWriter.write(
+          pendingRecord.recordKey,
+          recordContext.constructHoodieRecord(pendingRecord.record, 
partitionPath));
+      pendingRecord = null;
+    } catch (IOException e) {
+      throw new HoodieException("Failed to write the cdc data to native cdc 
log file", e);
+    }
+  }
+
+  private static class PendingCDCRecord<T> {
+    private final String recordKey;
+    private final BufferedRecord<T> record;
+
+    private PendingCDCRecord(String recordKey, BufferedRecord<T> record) {
+      this.recordKey = recordKey;
+      this.record = record;
+    }
+  }
+}
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeLogAppendHandle.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeLogAppendHandle.java
index 7e2318aa4f1d..f245c8c6c2f0 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeLogAppendHandle.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeLogAppendHandle.java
@@ -101,6 +101,7 @@ public class HoodieNativeLogAppendHandle<T, I, K, O> 
extends HoodieAppendHandle<
           getLogCreationCallback(),
           config.getWriteVersion(),
           config,
+          hoodieTable.getBaseFileFormat(),
           writeSchemaWithMetaFields,
           taskContextSupplier,
           
hoodieTable.getReaderContextFactoryForWrite().getContext().getRecordContext(),
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeLogFormatWriter.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeLogFormatWriter.java
index f45678b9f6b7..1b35f4fbe51d 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeLogFormatWriter.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeLogFormatWriter.java
@@ -65,6 +65,7 @@ public class HoodieNativeLogFormatWriter extends 
HoodieLogFormat.Writer {
   private static final String LOG_FORMAT_METADATA_FOOTER_KEY = 
"hudi.log.format.metadata";
 
   private final HoodieWriteConfig writeConfig;
+  private final HoodieFileFormat nativeFileFormat;
   private final HoodieSchema tableSchema;
   private final TaskContextSupplier taskContextSupplier;
   private final RecordContext recordContext;
@@ -90,6 +91,7 @@ public class HoodieNativeLogFormatWriter extends 
HoodieLogFormat.Writer {
                                      LogFileCreationCallback 
fileCreationCallback,
                                      HoodieTableVersion tableVersion,
                                      HoodieWriteConfig writeConfig,
+                                     HoodieFileFormat nativeFileFormat,
                                      HoodieSchema tableSchema,
                                      TaskContextSupplier taskContextSupplier,
                                      RecordContext recordContext,
@@ -97,6 +99,7 @@ public class HoodieNativeLogFormatWriter extends 
HoodieLogFormat.Writer {
     super(bufferSize, storage, parentPath, logFileId, DATA_LOG_EXTENSION, 
instantTime, logVersion, logWriteToken,
         null, 0L, sizeThreshold, fileCreationCallback, tableVersion);
     this.writeConfig = writeConfig;
+    this.nativeFileFormat = nativeFileFormat;
     this.tableSchema = tableSchema;
     this.taskContextSupplier = taskContextSupplier;
     this.recordContext = recordContext;
@@ -307,7 +310,7 @@ public class HoodieNativeLogFormatWriter extends 
HoodieLogFormat.Writer {
 
   private StoragePath makeNativeLogPath(int version, String logExtension) {
     return new StoragePath(parentPath, FSUtils.makeNativeLogFileName(
-        logFileId, logWriteToken, instantTime, version, logExtension, 
HoodieFileFormat.PARQUET));
+        logFileId, logWriteToken, instantTime, version, logExtension, 
nativeFileFormat));
   }
 
 }
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/LsmFileGroupReaderBasedMergeHandle.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/LsmFileGroupReaderBasedMergeHandle.java
new file mode 100644
index 000000000000..e83246b5e486
--- /dev/null
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/LsmFileGroupReaderBasedMergeHandle.java
@@ -0,0 +1,100 @@
+/*
+ * 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.common.config.TypedProperties;
+import org.apache.hudi.common.engine.HoodieReaderContext;
+import org.apache.hudi.common.engine.TaskContextSupplier;
+import org.apache.hudi.common.model.CompactionOperation;
+import org.apache.hudi.common.model.HoodieBaseFile;
+import org.apache.hudi.common.model.HoodieLogFile;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.table.read.HoodieRecordReader;
+import org.apache.hudi.common.table.read.lsm.HoodieLsmFileGroupReader;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.internal.schema.InternalSchema;
+import org.apache.hudi.keygen.BaseKeyGenerator;
+import org.apache.hudi.table.HoodieTable;
+
+import java.util.Comparator;
+import java.util.Iterator;
+import java.util.Map;
+import java.util.stream.Stream;
+
+/**
+ * A merge handle that uses the LSM file-group reader to merge sorted runs.
+ *
+ * <p>The incoming records, base file records, and native parquet log records 
are expected to be
+ * sorted by record key. The LSM reader performs a k-way merge and emits 
sorted output.
+ */
+public class LsmFileGroupReaderBasedMergeHandle<T, I, K, O> extends 
FileGroupReaderBasedMergeHandle<T, I, K, O> {
+
+  public LsmFileGroupReaderBasedMergeHandle(HoodieWriteConfig config, String 
instantTime, HoodieTable<T, I, K, O> hoodieTable,
+                                            Iterator<HoodieRecord<T>> 
recordItr, String partitionPath, String fileId,
+                                            TaskContextSupplier 
taskContextSupplier, Option<BaseKeyGenerator> keyGeneratorOpt) {
+    super(config, instantTime, hoodieTable, recordItr, partitionPath, fileId, 
taskContextSupplier, keyGeneratorOpt);
+  }
+
+  public LsmFileGroupReaderBasedMergeHandle(HoodieWriteConfig config, String 
instantTime, HoodieTable<T, I, K, O> hoodieTable,
+                                            Iterator<HoodieRecord<T>> 
recordItr, String partitionPath, String fileId,
+                                            TaskContextSupplier 
taskContextSupplier, HoodieBaseFile baseFile, Option<BaseKeyGenerator> 
keyGeneratorOpt) {
+    super(config, instantTime, hoodieTable, recordItr, partitionPath, fileId, 
taskContextSupplier, baseFile, keyGeneratorOpt);
+  }
+
+  public LsmFileGroupReaderBasedMergeHandle(HoodieWriteConfig config, String 
instantTime, HoodieTable<T, I, K, O> hoodieTable,
+                                            Map<String, HoodieRecord<T>> 
keyToNewRecords, String partitionPath, String fileId,
+                                            HoodieBaseFile dataFileToBeMerged, 
TaskContextSupplier taskContextSupplier,
+                                            Option<BaseKeyGenerator> 
keyGeneratorOpt) {
+    this(config, instantTime, hoodieTable, keyToNewRecords.values().stream()
+        .sorted(Comparator.comparing(HoodieRecord::getRecordKey)).iterator(), 
partitionPath, fileId,
+        taskContextSupplier, dataFileToBeMerged, keyGeneratorOpt);
+  }
+
+  public LsmFileGroupReaderBasedMergeHandle(HoodieWriteConfig config, String 
instantTime, HoodieTable<T, I, K, O> hoodieTable,
+                                            CompactionOperation 
compactionOperation, TaskContextSupplier taskContextSupplier,
+                                            HoodieReaderContext<T> 
readerContext, String maxInstantTime,
+                                            HoodieRecord.HoodieRecordType 
enginRecordType) {
+    super(config, instantTime, hoodieTable, compactionOperation, 
taskContextSupplier, readerContext, maxInstantTime, enginRecordType);
+  }
+
+  @Override
+  protected HoodieRecordReader<T> getFileGroupReader(boolean usePosition, 
Option<InternalSchema> internalSchemaOption, TypedProperties props,
+                                                     
Option<Stream<HoodieLogFile>> logFileStreamOpt, Iterator<HoodieRecord<T>> 
incomingRecordsItr) {
+    HoodieLsmFileGroupReader.HoodieLsmFileGroupReaderBuilder<T> 
fileGroupBuilder = HoodieLsmFileGroupReader.<T>builder()
+        .withReaderContext(readerContext)
+        .withHoodieTableMetaClient(hoodieTable.getMetaClient())
+        .withLatestCommitTime(maxInstantTime)
+        .withPartitionPath(partitionPath)
+        .withBaseFileOption(Option.ofNullable(baseFileToMerge))
+        .withDataSchema(writeSchemaWithMetaFields)
+        .withRequestedSchema(writeSchemaWithMetaFields)
+        .withInternalSchemaOpt(internalSchemaOption)
+        .withProps(props)
+        .withFileGroupUpdateCallback(createCallback());
+
+    if (logFileStreamOpt.isPresent()) {
+      fileGroupBuilder.withLogFiles(logFileStreamOpt.get());
+    } else {
+      fileGroupBuilder.withRecordIterator(incomingRecordsItr);
+    }
+    return fileGroupBuilder.build();
+  }
+}
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/HoodieTable.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/HoodieTable.java
index a857ef4417d6..ac9c6ad7dcf0 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/HoodieTable.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/HoodieTable.java
@@ -1052,7 +1052,8 @@ public abstract class HoodieTable<T, I, K, O> implements 
Serializable {
   }
 
   public boolean requireSortedRecords() {
-    return getBaseFileFormat() == HoodieFileFormat.HFILE;
+    return getBaseFileFormat() == HoodieFileFormat.HFILE
+        || getMetaClient().getTableConfig().isLSMTreeStorageLayout();
   }
 
   public HoodieEngineContext getContext() {
diff --git 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieMergeHandleFactory.java
 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieMergeHandleFactory.java
index 1cc3703df22f..89b4a015fff9 100644
--- 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieMergeHandleFactory.java
+++ 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieMergeHandleFactory.java
@@ -22,7 +22,9 @@ 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.util.collection.Pair;
+import org.apache.hudi.config.HoodieIndexConfig;
 import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.index.HoodieIndex;
 import org.apache.hudi.table.HoodieTable;
 
 import org.junit.jupiter.api.Assertions;
@@ -53,6 +55,7 @@ public class TestHoodieMergeHandleFactory {
     MockitoAnnotations.initMocks(this);
     when(mockHoodieTable.getMetaClient()).thenReturn(mockMetaClient);
     when(mockMetaClient.getTableConfig()).thenReturn(mockHoodieTableConfig);
+    when(mockHoodieTableConfig.isLSMTreeStorageLayout()).thenReturn(false);
   }
 
   @Test
@@ -67,10 +70,16 @@ public class TestHoodieMergeHandleFactory {
     when(mockHoodieTable.requireSortedRecords()).thenReturn(true);
     when(mockHoodieTableConfig.isCDCEnabled()).thenReturn(true);
     mergeHandleClasses = 
HoodieMergeHandleFactory.getMergeHandleClassesWrite(WriteOperationType.UPSERT, 
getWriterConfig(properties), mockHoodieTable);
-    validateMergeClasses(mergeHandleClasses, 
HoodieSortedMergeHandleWithChangeLog.class.getName());
+    validateMergeClasses(mergeHandleClasses, 
FileGroupReaderBasedMergeHandle.class.getName());
     when(mockHoodieTableConfig.isCDCEnabled()).thenReturn(false);
     mergeHandleClasses = 
HoodieMergeHandleFactory.getMergeHandleClassesWrite(WriteOperationType.UPSERT, 
getWriterConfig(properties), mockHoodieTable);
-    validateMergeClasses(mergeHandleClasses, 
HoodieSortedMergeHandle.class.getName());
+    validateMergeClasses(mergeHandleClasses, 
FileGroupReaderBasedMergeHandle.class.getName());
+
+    // LSM layout uses the LSM file-group-reader merge handle instead of the 
generic sorted merge handle.
+    when(mockHoodieTableConfig.isLSMTreeStorageLayout()).thenReturn(true);
+    mergeHandleClasses = 
HoodieMergeHandleFactory.getMergeHandleClassesWrite(WriteOperationType.UPSERT, 
getWriterConfig(properties), mockHoodieTable);
+    validateMergeClasses(mergeHandleClasses, 
LsmFileGroupReaderBasedMergeHandle.class.getName());
+    when(mockHoodieTableConfig.isLSMTreeStorageLayout()).thenReturn(false);
 
     // non-sorted: no CDC cases
     when(mockHoodieTable.requireSortedRecords()).thenReturn(false);
@@ -104,9 +113,14 @@ public class TestHoodieMergeHandleFactory {
     mergeHandleClasses = 
HoodieMergeHandleFactory.getMergeHandleClassesWrite(WriteOperationType.UPSERT, 
getWriterConfig(properties), mockHoodieTable);
     validateMergeClasses(mergeHandleClasses, CUSTOM_MERGE_HANDLE, 
FileGroupReaderBasedMergeHandle.class.getName());
 
+    when(mockHoodieTableConfig.isCDCEnabled()).thenReturn(true);
+    mergeHandleClasses = 
HoodieMergeHandleFactory.getMergeHandleClassesWrite(WriteOperationType.UPSERT, 
getWriterConfig(properties), mockHoodieTable);
+    validateMergeClasses(mergeHandleClasses, 
FileGroupReaderBasedMergeHandle.class.getName());
+    when(mockHoodieTableConfig.isCDCEnabled()).thenReturn(false);
+
     when(mockHoodieTable.requireSortedRecords()).thenReturn(true);
     mergeHandleClasses = 
HoodieMergeHandleFactory.getMergeHandleClassesWrite(WriteOperationType.UPSERT, 
getWriterConfig(properties), mockHoodieTable);
-    validateMergeClasses(mergeHandleClasses, 
HoodieSortedMergeHandle.class.getName());
+    validateMergeClasses(mergeHandleClasses, 
FileGroupReaderBasedMergeHandle.class.getName());
 
     when(mockHoodieTable.requireSortedRecords()).thenReturn(false);
     mergeHandleClasses = 
HoodieMergeHandleFactory.getMergeHandleClassesWrite(WriteOperationType.INSERT, 
getWriterConfig(propsWithDups), mockHoodieTable);
@@ -143,12 +157,32 @@ public class TestHoodieMergeHandleFactory {
     mergeHandleClasses = 
HoodieMergeHandleFactory.getMergeHandleClassesCompaction(getWriterConfig(properties),
 mockHoodieTable);
     validateMergeClasses(mergeHandleClasses, 
HoodieSortedMergeHandle.class.getName());
 
+    // LSM layout still uses the non-reader-context merge handle selection 
here.
+    when(mockHoodieTable.requireSortedRecords()).thenReturn(false);
+    when(mockHoodieTableConfig.isLSMTreeStorageLayout()).thenReturn(true);
+    mergeHandleClasses = 
HoodieMergeHandleFactory.getMergeHandleClassesCompaction(getWriterConfig(properties),
 mockHoodieTable);
+    validateMergeClasses(mergeHandleClasses, 
FileGroupReaderBasedMergeHandle.class.getName());
+
     // custom case
+    when(mockHoodieTableConfig.isLSMTreeStorageLayout()).thenReturn(false);
     when(mockHoodieTable.requireSortedRecords()).thenReturn(false);
     properties.setProperty(HoodieWriteConfig.MERGE_HANDLE_CLASS_NAME.key(), 
CUSTOM_MERGE_HANDLE);
     mergeHandleClasses = 
HoodieMergeHandleFactory.getMergeHandleClassesCompaction(getWriterConfig(properties),
 mockHoodieTable);
     validateMergeClasses(mergeHandleClasses, CUSTOM_MERGE_HANDLE, 
FileGroupReaderBasedMergeHandle.class.getName());
 
+    Properties pureLogProps = new Properties();
+    pureLogProps.setProperty(HoodieWriteConfig.MERGE_HANDLE_CLASS_NAME.key(), 
CUSTOM_MERGE_HANDLE);
+    pureLogProps.setProperty(HoodieIndexConfig.INDEX_TYPE.key(), 
HoodieIndex.IndexType.FLINK_STATE.name());
+    when(mockHoodieTableConfig.isCDCEnabled()).thenReturn(true);
+    mergeHandleClasses = 
HoodieMergeHandleFactory.getMergeHandleClassesCompaction(getWriterConfig(pureLogProps),
 mockHoodieTable);
+    validateMergeClasses(mergeHandleClasses, 
HoodieMergeHandleWithChangeLog.class.getName());
+
+    when(mockHoodieTableConfig.isLSMTreeStorageLayout()).thenReturn(true);
+    mergeHandleClasses = 
HoodieMergeHandleFactory.getMergeHandleClassesCompaction(getWriterConfig(pureLogProps),
 mockHoodieTable);
+    validateMergeClasses(mergeHandleClasses, 
HoodieMergeHandleWithChangeLog.class.getName());
+    when(mockHoodieTableConfig.isLSMTreeStorageLayout()).thenReturn(false);
+    when(mockHoodieTableConfig.isCDCEnabled()).thenReturn(false);
+
     when(mockHoodieTable.requireSortedRecords()).thenReturn(true);
     mergeHandleClasses = 
HoodieMergeHandleFactory.getMergeHandleClassesCompaction(getWriterConfig(properties),
 mockHoodieTable);
     validateMergeClasses(mergeHandleClasses, 
HoodieSortedMergeHandle.class.getName());
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkIncrementalMergeHandleWithChangeLog.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkIncrementalMergeHandleWithChangeLog.java
index 8b6ea645905b..225902aa3acf 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkIncrementalMergeHandleWithChangeLog.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkIncrementalMergeHandleWithChangeLog.java
@@ -48,21 +48,34 @@ import java.util.List;
 public class FlinkIncrementalMergeHandleWithChangeLog<T, I, K, O>
     extends FlinkIncrementalMergeHandle<T, I, K, O> {
 
-  private final HoodieCDCLogger cdcLogger;
+  private final HoodieCDCLogWriter<IndexedRecord> cdcLogger;
 
   public FlinkIncrementalMergeHandleWithChangeLog(HoodieWriteConfig config, 
String instantTime, HoodieTable<T, I, K, O> hoodieTable,
                                                   Iterator<HoodieRecord<T>> 
recordItr, String partitionPath, String fileId,
                                                   TaskContextSupplier 
taskContextSupplier, StoragePath basePath) {
     super(config, instantTime, hoodieTable, recordItr, partitionPath, fileId, 
taskContextSupplier, basePath);
-    this.cdcLogger = new HoodieCDCLogger(
+    this.cdcLogger = createCDCLogWriter(instantTime, config, hoodieTable, 
partitionPath, fileId, taskContextSupplier);
+  }
+
+  private HoodieCDCLogWriter<IndexedRecord> createCDCLogWriter(
+      String instantTime,
+      HoodieWriteConfig config,
+      HoodieTable<T, I, K, O> hoodieTable,
+      String partitionPath,
+      String fileId,
+      TaskContextSupplier taskContextSupplier) {
+    return HoodieCDCLogWriterFactory.createAvroCDCLogWriter(
         instantTime,
         config,
-        hoodieTable.getMetaClient().getTableConfig(),
+        hoodieTable,
         partitionPath,
         getStorage(),
         getWriterSchema(),
-        createLogWriter(instantTime, HoodieCDCUtils.CDC_LOGFILE_SUFFIX, 
Option.empty()),
-        IOUtils.getMaxMemoryPerPartitionMerge(taskContextSupplier, config));
+        fileId,
+        writeToken,
+        getLogCreationCallback(),
+        taskContextSupplier,
+        () -> createLogWriter(instantTime, HoodieCDCUtils.CDC_LOGFILE_SUFFIX, 
Option.empty()));
   }
 
   @Override
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkLsmFileGroupReaderBasedIncrementalMergeHandle.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkLsmFileGroupReaderBasedIncrementalMergeHandle.java
new file mode 100644
index 000000000000..6b7bf5a655d3
--- /dev/null
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkLsmFileGroupReaderBasedIncrementalMergeHandle.java
@@ -0,0 +1,77 @@
+/*
+ * 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.engine.TaskContextSupplier;
+import org.apache.hudi.common.fs.FSUtils;
+import org.apache.hudi.common.model.HoodieBaseFile;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.exception.HoodieIOException;
+import org.apache.hudi.storage.StoragePath;
+import org.apache.hudi.table.HoodieTable;
+
+import java.io.IOException;
+import java.util.Iterator;
+import java.util.List;
+
+/**
+ * Flink incremental mini-batch merge handle backed by the LSM file-group 
reader.
+ */
+public class FlinkLsmFileGroupReaderBasedIncrementalMergeHandle<T, I, K, O>
+    extends FlinkLsmFileGroupReaderBasedMergeHandle<T, I, K, O>
+    implements MiniBatchHandle {
+
+  public FlinkLsmFileGroupReaderBasedIncrementalMergeHandle(HoodieWriteConfig 
config, String instantTime, HoodieTable<T, I, K, O> hoodieTable,
+                                                            
Iterator<HoodieRecord<T>> recordItr, String partitionPath, String fileId,
+                                                            
TaskContextSupplier taskContextSupplier, StoragePath basePath) {
+    super(config, instantTime, hoodieTable, recordItr, partitionPath, fileId, 
taskContextSupplier, new HoodieBaseFile(basePath.toString()));
+  }
+
+  @Override
+  protected String createNewFileName(String oldFileName) {
+    int rollNumber = MergeHandleUtils.calcRollNumberForBaseFile(oldFileName, 
writeToken);
+    return newFileNameWithRollover(rollNumber);
+  }
+
+  protected String newFileNameWithRollover(int rollNumber) {
+    return FSUtils.makeBaseFileName(instantTime, writeToken + "-" + rollNumber,
+        this.fileId, hoodieTable.getBaseFileExtension());
+  }
+
+  public void finalizeWrite() {
+    try {
+      storage.deleteFile(oldFilePath);
+    } catch (IOException e) {
+      throw new HoodieIOException("Error while cleaning the old base file: " + 
oldFilePath, e);
+    }
+  }
+
+  @Override
+  public List<WriteStatus> close() {
+    if (isClosed()) {
+      return getWriteStatuses();
+    }
+    List<WriteStatus> writeStatuses = super.close();
+    finalizeWrite();
+    return writeStatuses;
+  }
+}
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkLsmFileGroupReaderBasedMergeHandle.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkLsmFileGroupReaderBasedMergeHandle.java
new file mode 100644
index 000000000000..03652b0fed3e
--- /dev/null
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkLsmFileGroupReaderBasedMergeHandle.java
@@ -0,0 +1,118 @@
+/*
+ * 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.common.engine.TaskContextSupplier;
+import org.apache.hudi.common.fs.FSUtils;
+import org.apache.hudi.common.model.HoodieBaseFile;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.exception.HoodieException;
+import org.apache.hudi.storage.StoragePath;
+import org.apache.hudi.table.HoodieTable;
+import org.apache.hudi.table.marker.WriteMarkers;
+import org.apache.hudi.table.marker.WriteMarkersFactory;
+
+import lombok.extern.slf4j.Slf4j;
+
+import java.io.IOException;
+import java.util.Iterator;
+
+/**
+ * Flink mini-batch merge handle backed by the LSM file-group reader.
+ */
+@Slf4j
+public class FlinkLsmFileGroupReaderBasedMergeHandle<T, I, K, O>
+    extends LsmFileGroupReaderBasedMergeHandle<T, I, K, O>
+    implements MiniBatchHandle {
+
+  public FlinkLsmFileGroupReaderBasedMergeHandle(HoodieWriteConfig config, 
String instantTime, HoodieTable<T, I, K, O> hoodieTable,
+                                                 Iterator<HoodieRecord<T>> 
recordItr, String partitionPath, String fileId,
+                                                 TaskContextSupplier 
taskContextSupplier) {
+    this(config, instantTime, hoodieTable, recordItr, partitionPath, fileId, 
taskContextSupplier, getLatestBaseFile(hoodieTable, partitionPath, fileId));
+  }
+
+  public FlinkLsmFileGroupReaderBasedMergeHandle(HoodieWriteConfig config, 
String instantTime, HoodieTable<T, I, K, O> hoodieTable,
+                                                 Iterator<HoodieRecord<T>> 
recordItr, String partitionPath, String fileId,
+                                                 TaskContextSupplier 
taskContextSupplier, HoodieBaseFile hoodieBaseFile) {
+    super(config, instantTime, hoodieTable, recordItr, partitionPath, fileId, 
taskContextSupplier, hoodieBaseFile, Option.empty());
+    if (getAttemptId() > 0) {
+      deleteInvalidDataFile(getAttemptId() - 1);
+    }
+  }
+
+  @Override
+  protected long getMaxMemoryForMerge() {
+    return Long.MAX_VALUE;
+  }
+
+  private void deleteInvalidDataFile(long lastAttemptId) {
+    final String lastWriteToken = FSUtils.makeWriteToken(getPartitionId(), 
getStageId(), lastAttemptId);
+    final String lastDataFileName = FSUtils.makeBaseFileName(instantTime,
+        lastWriteToken, this.fileId, hoodieTable.getBaseFileExtension());
+    final StoragePath path = makeNewFilePath(partitionPath, lastDataFileName);
+    if (path.equals(oldFilePath)) {
+      return;
+    }
+    try {
+      if (storage.exists(path)) {
+        log.info("Deleting invalid MERGE base file due to task retry: {}", 
lastDataFileName);
+        storage.deleteFile(path);
+      }
+    } catch (IOException e) {
+      throw new HoodieException("Error while deleting the MERGE base file due 
to task retry: " + lastDataFileName, e);
+    }
+  }
+
+  @Override
+  protected void createMarkerFile(String partitionPath, String dataFileName) {
+    WriteMarkers writeMarkers = 
WriteMarkersFactory.get(config.getMarkersType(), hoodieTable, instantTime);
+    writeMarkers.createIfNotExists(partitionPath, dataFileName, getIOType());
+  }
+
+  @Override
+  boolean needsUpdateLocation() {
+    return false;
+  }
+
+  @Override
+  public void closeGracefully() {
+    if (isClosed()) {
+      return;
+    }
+    try {
+      close();
+    } catch (Throwable throwable) {
+      log.error("Failed to close the MERGE handle", throwable);
+      try {
+        storage.deleteFile(newFilePath);
+        log.info("Successfully deleted the intermediate MERGE data file: {}", 
newFilePath);
+      } catch (IOException e) {
+        log.warn("Failed to delete the intermediate MERGE data file: {}", 
newFilePath, e);
+      }
+    }
+  }
+
+  @Override
+  public StoragePath getWritePath() {
+    return newFilePath;
+  }
+}
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkMergeHandleWithChangeLog.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkMergeHandleWithChangeLog.java
index 17f48e9e3131..6875cf7bf5d6 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkMergeHandleWithChangeLog.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkMergeHandleWithChangeLog.java
@@ -46,21 +46,34 @@ import java.util.List;
 @Slf4j
 public class FlinkMergeHandleWithChangeLog<T, I, K, O>
     extends FlinkMergeHandle<T, I, K, O> {
-  private final HoodieCDCLogger cdcLogger;
+  private final HoodieCDCLogWriter<IndexedRecord> cdcLogger;
 
   public FlinkMergeHandleWithChangeLog(HoodieWriteConfig config, String 
instantTime, HoodieTable<T, I, K, O> hoodieTable,
                                        Iterator<HoodieRecord<T>> recordItr, 
String partitionPath, String fileId,
                                        TaskContextSupplier 
taskContextSupplier) {
     super(config, instantTime, hoodieTable, recordItr, partitionPath, fileId, 
taskContextSupplier);
-    this.cdcLogger = new HoodieCDCLogger(
+    this.cdcLogger = createCDCLogWriter(instantTime, config, hoodieTable, 
partitionPath, fileId, taskContextSupplier);
+  }
+
+  private HoodieCDCLogWriter<IndexedRecord> createCDCLogWriter(
+      String instantTime,
+      HoodieWriteConfig config,
+      HoodieTable<T, I, K, O> hoodieTable,
+      String partitionPath,
+      String fileId,
+      TaskContextSupplier taskContextSupplier) {
+    return HoodieCDCLogWriterFactory.createAvroCDCLogWriter(
         instantTime,
         config,
-        hoodieTable.getMetaClient().getTableConfig(),
+        hoodieTable,
         partitionPath,
         getStorage(),
         getWriterSchema(),
-        createLogWriter(instantTime, HoodieCDCUtils.CDC_LOGFILE_SUFFIX, 
Option.empty()),
-        IOUtils.getMaxMemoryPerPartitionMerge(taskContextSupplier, config));
+        fileId,
+        writeToken,
+        getLogCreationCallback(),
+        taskContextSupplier,
+        () -> createLogWriter(instantTime, HoodieCDCUtils.CDC_LOGFILE_SUFFIX, 
Option.empty()));
   }
 
   @Override
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkWriteHandleFactory.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkWriteHandleFactory.java
index 390645c10ba5..40ad0c8c5000 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkWriteHandleFactory.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkWriteHandleFactory.java
@@ -150,7 +150,12 @@ public class FlinkWriteHandleFactory {
 
   private static boolean isFileGroupReaderBasedHandle(HoodieWriteConfig 
writeConfig) {
     String mergeHandleClass = writeConfig.getMergeHandleClassName();
-    return 
FileGroupReaderBasedMergeHandle.class.getName().equalsIgnoreCase(mergeHandleClass);
+    return 
FileGroupReaderBasedMergeHandle.class.getName().equalsIgnoreCase(mergeHandleClass)
+        || 
LsmFileGroupReaderBasedMergeHandle.class.getName().equalsIgnoreCase(mergeHandleClass);
+  }
+
+  private static boolean isLsmTreeStorageLayout(HoodieTable<?, ?, ?, ?> table) 
{
+    return table.getMetaClient().getTableConfig().isLSMTreeStorageLayout();
   }
 
   /**
@@ -175,8 +180,11 @@ public class FlinkWriteHandleFactory {
         String fileId,
         StoragePath basePath) {
       if (isFileGroupReaderBasedHandle(config)) {
-        return new FlinkFileGroupReaderBasedIncrementalMergeHandle<>(config, 
instantTime, table, recordItr, partitionPath, fileId,
-            table.getTaskContextSupplier(), basePath);
+        return isLsmTreeStorageLayout(table)
+            ? new FlinkLsmFileGroupReaderBasedIncrementalMergeHandle<>(config, 
instantTime, table, recordItr, partitionPath, fileId,
+                table.getTaskContextSupplier(), basePath)
+            : new FlinkFileGroupReaderBasedIncrementalMergeHandle<>(config, 
instantTime, table, recordItr, partitionPath, fileId,
+                table.getTaskContextSupplier(), basePath);
       } else {
         return new FlinkIncrementalMergeHandle<>(config, instantTime, table, 
recordItr, partitionPath, fileId,
             table.getTaskContextSupplier(), basePath);
@@ -192,8 +200,11 @@ public class FlinkWriteHandleFactory {
         String partitionPath,
         String fileId) {
       if (isFileGroupReaderBasedHandle(config)) {
-        return new FlinkFileGroupReaderBasedMergeHandle<>(config, instantTime, 
table, recordItr, partitionPath,
-            fileId, table.getTaskContextSupplier());
+        return isLsmTreeStorageLayout(table)
+            ? new FlinkLsmFileGroupReaderBasedMergeHandle<>(config, 
instantTime, table, recordItr, partitionPath,
+                fileId, table.getTaskContextSupplier())
+            : new FlinkFileGroupReaderBasedMergeHandle<>(config, instantTime, 
table, recordItr, partitionPath,
+                fileId, table.getTaskContextSupplier());
       } else {
         return new FlinkMergeHandle<>(config, instantTime, table, recordItr, 
partitionPath,
             fileId, table.getTaskContextSupplier());
@@ -261,8 +272,11 @@ public class FlinkWriteHandleFactory {
         String fileId,
         StoragePath basePath) {
       if (isFileGroupReaderBasedHandle(config)) {
-        return new FlinkFileGroupReaderBasedIncrementalMergeHandle<>(config, 
instantTime, table, recordItr, partitionPath, fileId,
-            table.getTaskContextSupplier(), basePath);
+        return isLsmTreeStorageLayout(table)
+            ? new FlinkLsmFileGroupReaderBasedIncrementalMergeHandle<>(config, 
instantTime, table, recordItr, partitionPath, fileId,
+                table.getTaskContextSupplier(), basePath)
+            : new FlinkFileGroupReaderBasedIncrementalMergeHandle<>(config, 
instantTime, table, recordItr, partitionPath, fileId,
+                table.getTaskContextSupplier(), basePath);
       } else {
         return new FlinkIncrementalMergeHandleWithChangeLog<>(config, 
instantTime, table, recordItr, partitionPath, fileId,
             table.getTaskContextSupplier(), basePath);
@@ -278,8 +292,11 @@ public class FlinkWriteHandleFactory {
         String partitionPath,
         String fileId) {
       if (isFileGroupReaderBasedHandle(config)) {
-        return new FlinkFileGroupReaderBasedMergeHandle<>(config, instantTime, 
table, recordItr, partitionPath,
-            fileId, table.getTaskContextSupplier());
+        return isLsmTreeStorageLayout(table)
+            ? new FlinkLsmFileGroupReaderBasedMergeHandle<>(config, 
instantTime, table, recordItr, partitionPath,
+                fileId, table.getTaskContextSupplier())
+            : new FlinkFileGroupReaderBasedMergeHandle<>(config, instantTime, 
table, recordItr, partitionPath,
+                fileId, table.getTaskContextSupplier());
       } else {
         return new FlinkMergeHandleWithChangeLog<>(config, instantTime, table, 
recordItr, partitionPath,
             fileId, table.getTaskContextSupplier());
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/FlinkWriteHelper.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/FlinkWriteHelper.java
index 284dece174aa..57ed5c25ac17 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/FlinkWriteHelper.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/FlinkWriteHelper.java
@@ -37,6 +37,7 @@ import org.apache.hudi.table.action.HoodieWriteMetadata;
 
 import java.time.Duration;
 import java.util.Iterator;
+import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
@@ -100,7 +101,8 @@ public class FlinkWriteHelper<T, R> extends 
BaseWriteHelper<T, Iterator<HoodieRe
                                                       String[] 
orderingFieldNames) {
     // If index used is global, then records are expected to differ in their 
partitionPath
     Map<Object, List<HoodieRecord<T>>> keyedRecords = 
CollectionUtils.toStream(records)
-        .collect(Collectors.groupingBy(record -> 
record.getKey().getRecordKey()));
+        .collect(Collectors.groupingBy(
+            record -> record.getKey().getRecordKey(), LinkedHashMap::new, 
Collectors.toList()));
 
     // caution that the avro schema is not serializable
     final HoodieSchema schema = HoodieSchema.parse(schemaStr);
diff --git 
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/table/action/commit/TestFlinkWriteHelper.java
 
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/table/action/commit/TestFlinkWriteHelper.java
new file mode 100644
index 000000000000..8387dc83015a
--- /dev/null
+++ 
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/table/action/commit/TestFlinkWriteHelper.java
@@ -0,0 +1,65 @@
+/*
+ * 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.table.action.commit;
+
+import org.apache.hudi.common.config.TypedProperties;
+import org.apache.hudi.common.model.HoodieAvroPayload;
+import org.apache.hudi.common.model.HoodieAvroRecord;
+import org.apache.hudi.common.model.HoodieKey;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.util.CollectionUtils;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.Arrays;
+import java.util.List;
+import java.util.stream.Collectors;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+class TestFlinkWriteHelper {
+
+  private static final String SCHEMA = 
"{\"type\":\"record\",\"name\":\"testrec\",\"fields\":[{\"name\":\"id\",\"type\":\"string\"}]}";
+
+  @Test
+  void testDeduplicateRecordsPreservesInputKeyOrder() {
+    List<HoodieRecord<HoodieAvroPayload>> records = Arrays.asList(record("b"), 
record("a"), record("c"));
+    @SuppressWarnings("unchecked")
+    FlinkWriteHelper<HoodieAvroPayload, Object> writeHelper = 
FlinkWriteHelper.newInstance();
+
+    List<String> deduplicatedKeys = CollectionUtils.toStream(
+        writeHelper.deduplicateRecords(
+            records.iterator(),
+            null,
+            -1,
+            SCHEMA,
+            new TypedProperties(),
+            null,
+            null,
+            new String[0]))
+        .map(HoodieRecord::getRecordKey)
+        .collect(Collectors.toList());
+
+    assertEquals(Arrays.asList("b", "a", "c"), deduplicatedKeys);
+  }
+
+  private static HoodieRecord<HoodieAvroPayload> record(String recordKey) {
+    return new HoodieAvroRecord<>(new HoodieKey(recordKey, "partition"), null);
+  }
+}
diff --git 
a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/deltacommit/BaseJavaDeltaCommitActionExecutor.java
 
b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/deltacommit/BaseJavaDeltaCommitActionExecutor.java
index 595f06006839..c8e5c5bce8ce 100644
--- 
a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/deltacommit/BaseJavaDeltaCommitActionExecutor.java
+++ 
b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/deltacommit/BaseJavaDeltaCommitActionExecutor.java
@@ -42,6 +42,8 @@ import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
 
+import static 
org.apache.hudi.common.util.HoodieRecordUtils.sortRecordsByRecordKey;
+
 @Slf4j
 abstract class BaseJavaDeltaCommitActionExecutor<T> extends 
BaseJavaCommitActionExecutor<T> {
 
@@ -75,6 +77,9 @@ abstract class BaseJavaDeltaCommitActionExecutor<T> extends 
BaseJavaCommitAction
       log.info("Small file corrections for updates for commit " + instantTime 
+ " for file " + fileId);
       return super.handleUpdate(partitionPath, fileId, recordItr);
     } else {
+      if (table.requireSortedRecords()) {
+        recordItr = sortRecordsByRecordKey(recordItr);
+      }
       HoodieAppendHandle<?, ?, ?, ?> appendHandle = new AppendHandleFactory()
           .create(config, instantTime, table, partitionPath, fileId, 
recordItr, taskContextSupplier);
       appendHandle.doAppend();
@@ -86,6 +91,9 @@ abstract class BaseJavaDeltaCommitActionExecutor<T> extends 
BaseJavaCommitAction
   public Iterator<List<WriteStatus>> handleInsert(String idPfx, 
Iterator<HoodieRecord<T>> recordItr) {
     // If canIndexLogFiles, write inserts to log files else write inserts to 
base files
     if (table.getIndex().canIndexLogFiles()) {
+      if (table.requireSortedRecords()) {
+        recordItr = sortRecordsByRecordKey(recordItr);
+      }
       return new JavaLazyInsertIterable<>(recordItr, true, config, 
instantTime, table, idPfx,
           taskContextSupplier, new AppendHandleFactory<>());
     } else {
diff --git 
a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/deltacommit/JavaUpsertPreppedDeltaCommitActionExecutor.java
 
b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/deltacommit/JavaUpsertPreppedDeltaCommitActionExecutor.java
index 42a42e819a60..d4d83358b940 100644
--- 
a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/deltacommit/JavaUpsertPreppedDeltaCommitActionExecutor.java
+++ 
b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/deltacommit/JavaUpsertPreppedDeltaCommitActionExecutor.java
@@ -37,9 +37,12 @@ import lombok.extern.slf4j.Slf4j;
 
 import java.util.ArrayList;
 import java.util.HashMap;
+import java.util.Iterator;
 import java.util.LinkedList;
 import java.util.List;
 
+import static 
org.apache.hudi.common.util.HoodieRecordUtils.sortRecordsByRecordKey;
+
 @Slf4j
 public class JavaUpsertPreppedDeltaCommitActionExecutor<T> extends 
BaseJavaDeltaCommitActionExecutor<T> {
 
@@ -76,8 +79,10 @@ public class JavaUpsertPreppedDeltaCommitActionExecutor<T> 
extends BaseJavaDelta
     List<WriteStatus> allWriteStatuses = new ArrayList<>();
     try {
       recordsByFileId.forEach((k, v) -> {
+        Iterator<HoodieRecord<T>> recordItr = table.requireSortedRecords()
+            ? sortRecordsByRecordKey(v.iterator()) : v.iterator();
         HoodieAppendHandle<?, ?, ?, ?> appendHandle = new AppendHandleFactory()
-            .create(config, instantTime, table, k.getRight(), k.getLeft(), 
v.iterator(), taskContextSupplier);
+            .create(config, instantTime, table, k.getRight(), k.getLeft(), 
recordItr, taskContextSupplier);
         appendHandle.doAppend();
         allWriteStatuses.addAll(appendHandle.close());
       });
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputSplit.java 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputSplit.java
index 6b8761db5f7c..a23feac06335 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputSplit.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputSplit.java
@@ -100,6 +100,10 @@ public class InputSplit {
     return !logFiles.isEmpty();
   }
 
+  public boolean hasRecordIterator() {
+    return recordIterator.isPresent();
+  }
+
   public boolean isParquetBaseFile() {
     return baseFileOption.map(baseFile -> 
HoodieFileFormat.fromFileExtension(baseFile.getStoragePath().getFileExtension())
 == HoodieFileFormat.PARQUET).orElse(false);
   }
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/HoodieLsmFileGroupReader.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/HoodieLsmFileGroupReader.java
index 507cbe46c708..8ae074a6d946 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/HoodieLsmFileGroupReader.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/HoodieLsmFileGroupReader.java
@@ -52,6 +52,7 @@ import lombok.Builder;
 import lombok.Getter;
 
 import java.io.IOException;
+import java.util.Iterator;
 import java.util.List;
 import java.util.function.UnaryOperator;
 import java.util.stream.Stream;
@@ -94,6 +95,7 @@ public final class HoodieLsmFileGroupReader<T> implements 
HoodieRecordReader<T>
       TypedProperties props,
       Option<HoodieBaseFile> baseFileOption,
       Stream<HoodieLogFile> logFiles,
+      Iterator<? extends HoodieRecord> recordIterator,
       String partitionPath,
       Long start,
       Long length,
@@ -108,7 +110,7 @@ public final class HoodieLsmFileGroupReader<T> implements 
HoodieRecordReader<T>
     ValidationUtils.checkArgument(requestedSchema != null, "Requested schema 
is required");
     ValidationUtils.checkArgument(props != null, "Props is required");
     ValidationUtils.checkArgument(partitionPath != null, "Partition path is 
required");
-    
ValidationUtils.checkArgument(hoodieTableMetaClient.getTableConfig().getLogFileFormat()
 == HoodieFileFormat.PARQUET,
+    ValidationUtils.checkArgument(logFiles == null || 
hoodieTableMetaClient.getTableConfig().getLogFileFormat() == 
HoodieFileFormat.PARQUET,
         "LSM file group reader expects parquet log files");
 
     if (internalSchemaOpt == null) {
@@ -145,6 +147,7 @@ public final class HoodieLsmFileGroupReader<T> implements 
HoodieRecordReader<T>
     this.inputSplit = InputSplit.builder()
         .baseFileOption(baseFileOption)
         .logFileStream(logFiles)
+        .recordIterator((Iterator<HoodieRecord>) recordIterator)
         .partitionPath(partitionPath)
         .start(start)
         .length(length)
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/LsmFileGroupRecordIterator.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/LsmFileGroupRecordIterator.java
index 3bc19d50f6eb..fde5bdad4a4f 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/LsmFileGroupRecordIterator.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/LsmFileGroupRecordIterator.java
@@ -26,6 +26,8 @@ import org.apache.hudi.common.model.HoodieBaseFile;
 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.schema.HoodieSchemaCache;
+import org.apache.hudi.common.schema.HoodieSchemaUtils;
 import org.apache.hudi.common.schema.HoodieSchemas;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.table.read.BaseFileUpdateCallback;
@@ -33,6 +35,7 @@ import org.apache.hudi.common.table.read.BufferedRecord;
 import org.apache.hudi.common.table.read.BufferedRecordMerger;
 import org.apache.hudi.common.table.read.BufferedRecordMergerFactory;
 import org.apache.hudi.common.table.read.BufferedRecords;
+import org.apache.hudi.common.table.read.DeleteContext;
 import org.apache.hudi.common.table.read.HoodieReadStats;
 import org.apache.hudi.common.table.read.InputSplit;
 import org.apache.hudi.common.table.read.ReaderParameters;
@@ -52,6 +55,7 @@ import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.Comparator;
 import java.util.HashSet;
+import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
 import java.util.NoSuchElementException;
@@ -89,6 +93,7 @@ public class LsmFileGroupRecordIterator<T> implements 
ClosableIterator<BufferedR
   private final InputSplit inputSplit;
   private final HoodieSchema readerSchema;
   private final List<String> orderingFieldNames;
+  private final TypedProperties props;
   private final boolean includeBaseFile;
   private final BufferedRecordMerger<T> bufferedRecordMerger;
   private final UpdateProcessor<T> updateProcessor;
@@ -133,6 +138,7 @@ public class LsmFileGroupRecordIterator<T> implements 
ClosableIterator<BufferedR
     this.inputSplit = inputSplit;
     this.readerSchema = readerContext.getSchemaHandler().getRequiredSchema();
     this.orderingFieldNames = orderingFieldNames;
+    this.props = props;
     this.includeBaseFile = includeBaseFile;
     this.bufferedRecordMerger = BufferedRecordMergerFactory.create(
         readerContext, readerContext.getMergeMode(), false, 
readerContext.getRecordMerger(),
@@ -158,9 +164,15 @@ public class LsmFileGroupRecordIterator<T> implements 
ClosableIterator<BufferedR
       addReader(sortedRunReaders, mergeOrder++, 
createBaseFileIterator(inputSplit.getBaseFileOption().get()));
     }
 
+    if (inputSplit.hasRecordIterator()) {
+      addReader(sortedRunReaders, mergeOrder++, 
createRecordIterator(inputSplit.getRecordIterator()));
+    }
+
     List<LogReaderSpec> logReaderSpecs = new ArrayList<>();
-    for (HoodieLogFile logFile : inputSplit.getLogFiles()) {
-      logReaderSpecs.add(new LogReaderSpec(mergeOrder++, logFile));
+    if (!inputSplit.hasRecordIterator()) {
+      for (HoodieLogFile logFile : inputSplit.getLogFiles()) {
+        logReaderSpecs.add(new LogReaderSpec(mergeOrder++, logFile));
+      }
     }
     Set<Integer> directLogMergeOrders = 
selectDirectLogMergeOrders(logReaderSpecs, hasBaseFileReader);
     for (LogReaderSpec spec : logReaderSpecs) {
@@ -254,6 +266,41 @@ public class LsmFileGroupRecordIterator<T> implements 
ClosableIterator<BufferedR
     return createFileIterator(baseFile.getPathInfo(), 
baseFile.getStoragePath(), baseFile.getFileSize());
   }
 
+  /**
+   * Creates a sorted-run iterator from incoming write records.
+   */
+  private ClosableIterator<BufferedRecord<T>> 
createRecordIterator(Iterator<HoodieRecord> recordIterator) {
+    HoodieSchema recordSchema = HoodieSchemaCache.intern(getRecordSchema());
+    String[] orderingFieldsArray = orderingFieldNames.toArray(new String[0]);
+    DeleteContext deleteContext = DeleteContext.fromRecordSchema(props, 
recordSchema);
+    return new ClosableIterator<BufferedRecord<T>>() {
+      @Override
+      public boolean hasNext() {
+        return recordIterator.hasNext();
+      }
+
+      @Override
+      public BufferedRecord<T> next() {
+        return BufferedRecords.fromHoodieRecord(recordIterator.next(), 
recordSchema, readerContext.getRecordContext(),
+            props, orderingFieldsArray, deleteContext);
+      }
+
+      @Override
+      public void close() {
+        // no op.
+      }
+    };
+  }
+
+  private HoodieSchema getRecordSchema() {
+    Option<Pair<String, String>> payloadClasses = 
readerContext.getPayloadClasses(props);
+    if (payloadClasses.isPresent() && 
payloadClasses.get().getRight().equals("org.apache.spark.sql.hudi.command.payload.ExpressionPayload"))
 {
+      String schemaStr = props.getString("hoodie.payload.record.schema");
+      return HoodieSchema.parse(schemaStr);
+    }
+    return 
HoodieSchemaUtils.removeMetadataFields(readerContext.getSchemaHandler().getRequestedSchema());
+  }
+
   /**
    * Creates a sorted-run iterator for a parquet data file or a native parquet 
log file.
    *
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/util/HoodieRecordUtils.java 
b/hudi-common/src/main/java/org/apache/hudi/common/util/HoodieRecordUtils.java
index 14b5ad8efd45..1718092eba68 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/util/HoodieRecordUtils.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/util/HoodieRecordUtils.java
@@ -50,7 +50,10 @@ import org.apache.avro.generic.GenericRecord;
 
 import java.lang.reflect.Constructor;
 import java.lang.reflect.InvocationTargetException;
+import java.util.ArrayList;
 import java.util.Collections;
+import java.util.Comparator;
+import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
@@ -220,6 +223,16 @@ public class HoodieRecordUtils {
     return null;
   }
 
+  /**
+   * Returns an iterator over the input records sorted by record key.
+   */
+  public static <T> Iterator<HoodieRecord<T>> 
sortRecordsByRecordKey(Iterator<HoodieRecord<T>> records) {
+    List<HoodieRecord<T>> sortedRecords = new ArrayList<>();
+    records.forEachRemaining(sortedRecords::add);
+    sortedRecords.sort(Comparator.comparing(HoodieRecord::getRecordKey));
+    return sortedRecords.iterator();
+  }
+
   public static List<String> getOrderingFieldNames(RecordMergeMode mergeMode,
                                                    HoodieTableMetaClient 
metaClient) {
     return mergeMode == RecordMergeMode.COMMIT_TIME_ORDERING
@@ -233,4 +246,4 @@ public class HoodieRecordUtils {
         ? Collections.emptyList()
         : tableConfig.getOrderingFields();
   }
-}
\ No newline at end of file
+}
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/util/TestHoodieRecordUtils.java
 
b/hudi-common/src/test/java/org/apache/hudi/common/util/TestHoodieRecordUtils.java
index df35e7583874..2219f8043760 100644
--- 
a/hudi-common/src/test/java/org/apache/hudi/common/util/TestHoodieRecordUtils.java
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/util/TestHoodieRecordUtils.java
@@ -21,7 +21,10 @@ package org.apache.hudi.common.util;
 import org.apache.hudi.common.config.RecordMergeMode;
 import org.apache.hudi.common.config.TypedProperties;
 import org.apache.hudi.common.model.DefaultHoodieRecordPayload;
+import org.apache.hudi.common.model.HoodieAvroRecord;
 import org.apache.hudi.common.model.HoodieAvroRecordMerger;
+import org.apache.hudi.common.model.HoodieKey;
+import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.model.HoodieRecordMerger;
 import org.apache.hudi.common.model.HoodieRecordPayload;
 import org.apache.hudi.common.table.HoodieTableConfig;
@@ -30,7 +33,11 @@ import org.apache.hudi.exception.HoodieException;
 
 import org.junit.jupiter.api.Test;
 
+import java.util.ArrayList;
+import java.util.Arrays;
 import java.util.Collections;
+import java.util.Iterator;
+import java.util.List;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertThrows;
@@ -62,6 +69,21 @@ class TestHoodieRecordUtils {
     assertEquals(payload.getClass().getName(), payloadClassName);
   }
 
+  @Test
+  void sortRecordsByRecordKey() {
+    List<HoodieRecord<DefaultHoodieRecordPayload>> records = Arrays.asList(
+        record("key3"),
+        record("key1"),
+        record("key2"));
+
+    Iterator<HoodieRecord<DefaultHoodieRecordPayload>> sortedRecords =
+        HoodieRecordUtils.sortRecordsByRecordKey(records.iterator());
+
+    List<String> sortedKeys = new ArrayList<>();
+    sortedRecords.forEachRemaining(record -> 
sortedKeys.add(record.getRecordKey()));
+    assertEquals(Arrays.asList("key1", "key2", "key3"), sortedKeys);
+  }
+
   @Test
   void testGetOrderingFields() {
     HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
@@ -79,4 +101,8 @@ class TestHoodieRecordUtils {
     props.setProperty("hoodie.table.ordering.fields", "props");
     assertEquals(Collections.singletonList("tbl"), 
HoodieRecordUtils.getOrderingFieldNames(RecordMergeMode.EVENT_TIME_ORDERING, 
metaClient));
   }
-}
\ No newline at end of file
+
+  private HoodieRecord<DefaultHoodieRecordPayload> record(String recordKey) {
+    return new HoodieAvroRecord<>(new HoodieKey(recordKey, "partition"), new 
DefaultHoodieRecordPayload(Option.empty()));
+  }
+}
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/configuration/OptionsResolver.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/configuration/OptionsResolver.java
index 8853d5844984..0dd529c7031e 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/configuration/OptionsResolver.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/configuration/OptionsResolver.java
@@ -70,6 +70,7 @@ import static 
org.apache.hudi.common.config.HoodieMetadataConfig.RECORD_INDEX_GR
 import static 
org.apache.hudi.common.config.HoodieMetadataConfig.RECORD_INDEX_MAX_FILE_GROUP_SIZE_BYTES_PROP;
 import static 
org.apache.hudi.common.config.HoodieMetadataConfig.RECORD_LEVEL_INDEX_MAX_FILE_GROUP_COUNT_PROP;
 import static 
org.apache.hudi.common.config.HoodieMetadataConfig.RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_PROP;
+import static 
org.apache.hudi.common.table.HoodieTableConfig.TableStorageLayout.LSM_TREE;
 import static 
org.apache.hudi.metadata.HoodieBackedTableMetadataWriter.RECORD_INDEX_AVERAGE_RECORD_SIZE;
 
 /**
@@ -162,6 +163,15 @@ public class OptionsResolver {
         .equals(FlinkOptions.TABLE_TYPE_COPY_ON_WRITE);
   }
 
+  /**
+   * Returns whether the table uses LSM tree storage layout.
+   */
+  public static boolean isLsmTreeStorageLayout(Configuration conf) {
+    return HoodieTableConfig.TableStorageLayout.fromConfigValue(conf.getString(
+        HoodieTableConfig.TABLE_STORAGE_LAYOUT.key(),
+        HoodieTableConfig.TABLE_STORAGE_LAYOUT.defaultValue())) == LSM_TREE;
+  }
+
   /**
    * Returns whether the payload clazz is {@link DefaultHoodieRecordPayload}.
    */
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteFunction.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteFunction.java
index 4dd8f785ac6b..073b977aa8e6 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteFunction.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteFunction.java
@@ -38,6 +38,7 @@ import org.apache.hudi.sink.buffer.RowDataBucket;
 import org.apache.hudi.sink.buffer.TotalSizeTracer;
 import org.apache.hudi.sink.bulk.RowDataKeyGen;
 import org.apache.hudi.sink.bulk.RowDataKeyGens;
+import org.apache.hudi.sink.bulk.sort.SortOperatorGen;
 import org.apache.hudi.sink.common.AbstractStreamWriteFunction;
 import org.apache.hudi.sink.event.WriteMetadataEvent;
 import org.apache.hudi.sink.exception.MemoryPagesExhaustedException;
@@ -55,6 +56,10 @@ import org.apache.flink.configuration.Configuration;
 import org.apache.flink.metrics.MetricGroup;
 import org.apache.flink.table.data.RowData;
 import org.apache.flink.table.data.binary.BinaryRowData;
+import org.apache.flink.table.planner.codegen.sort.SortCodeGenerator;
+import org.apache.flink.table.runtime.generated.GeneratedNormalizedKeyComputer;
+import org.apache.flink.table.runtime.generated.GeneratedRecordComparator;
+import org.apache.flink.table.runtime.operators.sort.BinaryInMemorySortBuffer;
 import org.apache.flink.table.runtime.util.MemorySegmentPool;
 import org.apache.flink.table.types.logical.RowType;
 import org.apache.flink.util.Collector;
@@ -146,6 +151,9 @@ public class StreamWriteFunction extends 
AbstractStreamWriteFunction<HoodieFlink
 
   protected transient RecordConverter recordConverter;
 
+  private transient GeneratedNormalizedKeyComputer recordKeyComputer;
+  private transient GeneratedRecordComparator recordKeyComparator;
+
   /**
    * Constructs a StreamingSinkFunction.
    *
@@ -161,6 +169,7 @@ public class StreamWriteFunction extends 
AbstractStreamWriteFunction<HoodieFlink
   @Override
   public void open(Configuration parameters) throws IOException {
     this.tracer = new TotalSizeTracer(this.config);
+    initRecordKeySort();
     initBuffer();
     initWriteFunction();
     initIndexProcessFunction();
@@ -203,6 +212,21 @@ public class StreamWriteFunction extends 
AbstractStreamWriteFunction<HoodieFlink
     this.memorySegmentPool = 
this.memorySegmentPoolFactory.createMemorySegmentPool(config, 
OptionsResolver.getWriteBufferSizeInBytes(config));
   }
 
+  private void initRecordKeySort() {
+    if (!OptionsResolver.isLsmTreeStorageLayout(config)) {
+      return;
+    }
+    String[] recordKeyFields = OptionsResolver.getRecordKeys(config);
+    ValidationUtils.checkArgument(recordKeyFields.length > 0,
+        "Record key fields can't be empty for LSM storage layout stream 
write.");
+    SortOperatorGen sortOperatorGen = new SortOperatorGen(rowType, 
recordKeyFields);
+    SortCodeGenerator codeGenerator = 
sortOperatorGen.createSortCodeGenerator();
+    this.recordKeyComputer = 
codeGenerator.generateNormalizedKeyComputer("LsmRecordKeySortComputer");
+    this.recordKeyComparator = 
codeGenerator.generateRecordComparator("LsmRecordKeySortComparator");
+    log.info("LSM storage layout stream write will sort buffered RowData by 
record keys: {}",
+        String.join(",", recordKeyFields));
+  }
+
   private void initWriteFunction() {
     final String writeOperation = this.config.get(FlinkOptions.OPERATION);
     switch (WriteOperationType.fromValue(writeOperation)) {
@@ -287,7 +311,7 @@ public class StreamWriteFunction extends 
AbstractStreamWriteFunction<HoodieFlink
       RowDataBucket bucket = this.buckets.computeIfAbsent(bucketID,
           k -> new RowDataBucket(
               bucketID,
-              BufferUtils.createBuffer(rowType, memorySegmentPool),
+              createDataBuffer(),
               getBucketInfo(record),
               this.config.get(FlinkOptions.WRITE_BATCH_SIZE)));
 
@@ -436,6 +460,7 @@ public class StreamWriteFunction extends 
AbstractStreamWriteFunction<HoodieFlink
       RowDataBucket rowDataBucket) {
     writeMetrics.startFileFlush();
 
+    sortBucketIfNeeded(rowDataBucket);
     Iterator<BinaryRowData> rowItr =
         new MutableIteratorWrapperIterator<>(
             rowDataBucket.getDataIterator(), () -> new 
BinaryRowData(rowType.getFieldCount()));
@@ -449,6 +474,33 @@ public class StreamWriteFunction extends 
AbstractStreamWriteFunction<HoodieFlink
     return statuses;
   }
 
+  private BinaryInMemorySortBuffer createDataBuffer() {
+    if (recordKeyComputer == null) {
+      return BufferUtils.createBuffer(rowType, memorySegmentPool);
+    }
+    try {
+      ClassLoader classLoader = Thread.currentThread().getContextClassLoader();
+      return BufferUtils.createBuffer(
+          rowType,
+          memorySegmentPool,
+          recordKeyComputer.newInstance(classLoader),
+          recordKeyComparator.newInstance(classLoader));
+    } catch (Exception e) {
+      throw new HoodieException("Failed to create RowData record-key sort 
buffer for LSM storage layout.", e);
+    }
+  }
+
+  private void sortBucketIfNeeded(RowDataBucket rowDataBucket) {
+    if (recordKeyComputer == null) {
+      return;
+    }
+    try {
+      rowDataBucket.sort();
+    } catch (IOException e) {
+      throw new HoodieException("Failed to sort buffered RowData records by 
record key.", e);
+    }
+  }
+
   protected Iterator<HoodieRecord> 
deduplicateRecordsIfNeeded(Iterator<HoodieRecord> records) {
     if (config.get(FlinkOptions.PRE_COMBINE)) {
       return FlinkWriteHelper.newInstance().deduplicateRecords(
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/buffer/RowDataBucket.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/buffer/RowDataBucket.java
index 899f465a4bf5..4e35cb5e3718 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/buffer/RowDataBucket.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/buffer/RowDataBucket.java
@@ -21,6 +21,7 @@ package org.apache.hudi.sink.buffer;
 import org.apache.hudi.table.action.commit.BucketInfo;
 
 import lombok.Getter;
+import org.apache.flink.runtime.operators.sort.QuickSort;
 import org.apache.flink.table.data.RowData;
 import org.apache.flink.table.data.binary.BinaryRowData;
 import org.apache.flink.table.runtime.operators.sort.BinaryInMemorySortBuffer;
@@ -55,6 +56,10 @@ public class RowDataBucket {
     return dataBuffer.getIterator();
   }
 
+  public void sort() throws IOException {
+    new QuickSort().sort(dataBuffer);
+  }
+
   public boolean writeRow(RowData rowData) throws IOException {
     boolean success = dataBuffer.write(rowData);
     if (success) {
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/HoodieTableFactory.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/HoodieTableFactory.java
index 12a0b1092588..a68777d3fe63 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/HoodieTableFactory.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/HoodieTableFactory.java
@@ -129,6 +129,10 @@ public class HoodieTableFactory implements 
DynamicTableSourceFactory, DynamicTab
               && !conf.contains(FlinkOptions.HIVE_STYLE_PARTITIONING)) {
             conf.set(FlinkOptions.HIVE_STYLE_PARTITIONING, 
tableConfig.getBoolean(HoodieTableConfig.HIVE_STYLE_PARTITIONING_ENABLE));
           }
+          if (tableConfig.contains(HoodieTableConfig.TABLE_STORAGE_LAYOUT)
+              && 
!conf.containsKey(HoodieTableConfig.TABLE_STORAGE_LAYOUT.key())) {
+            conf.setString(HoodieTableConfig.TABLE_STORAGE_LAYOUT.key(), 
tableConfig.getString(HoodieTableConfig.TABLE_STORAGE_LAYOUT));
+          }
           if (tableConfig.contains(HoodieTableConfig.TYPE) && 
conf.contains(FlinkOptions.TABLE_TYPE)) {
             if 
(!tableConfig.getString(HoodieTableConfig.TYPE).equals(conf.get(FlinkOptions.TABLE_TYPE)))
 {
               log.error("Table type conflict : {} in {} and {} in table 
options. Update your config to match the table type in hoodie.properties.",

Reply via email to