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

voonhous 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 49fade8653c2 feat(spark): enable format-aware sort ordering and LSM 
reading for Spark (#19502)
49fade8653c2 is described below

commit 49fade8653c2fc6079c0457890223696d130c08f
Author: Shuo Cheng <[email protected]>
AuthorDate: Fri Aug 7 17:11:21 2026 +0800

    feat(spark): enable format-aware sort ordering and LSM reading for Spark 
(#19502)
    
    * feat(spark): enable UTF-8 record key ordering and LSM reading in Spark
    
    * fix comment
    
    * fix comments
---
 .../io/LsmFileGroupReaderBasedMergeHandle.java     |  11 -
 .../table/read/lsm/LsmFileGroupRecordIterator.java |   4 +-
 .../hudi/common/table/read/lsm/LsmReaderUtils.java |  41 +++
 .../hudi/metadata/HoodieBackedTableMetadata.java   |  12 +-
 .../read/lsm/TestHoodieLsmFileGroupReader.java     |  23 ++
 .../read/lsm/TestLsmFileGroupRecordIterator.java   |  15 +
 .../common/table/read/lsm/TestLsmReaderUtils.java  |  43 +++
 .../org/apache/hudi/table/format/FormatUtils.java  |  21 +-
 .../apache/hudi/table/format/TestFormatUtils.java  |  40 ---
 .../org/apache/hudi/HoodieMergeOnReadRDDV2.scala   |  43 ++-
 .../org/apache/hudi/HoodieSparkSqlWriter.scala     |   2 +
 .../HoodieFileGroupReaderBasedFileFormat.scala     |  57 ++--
 .../apache/hudi/functional/TestMORDataSource.scala | 326 ++++++++++++++++++++-
 13 files changed, 530 insertions(+), 108 deletions(-)

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
index 57fda150e0fd..01c327de6374 100644
--- 
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
@@ -34,9 +34,7 @@ import org.apache.hudi.config.HoodieWriteConfig;
 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;
 
 /**
@@ -59,15 +57,6 @@ public class LsmFileGroupReaderBasedMergeHandle<T, I, K, O> 
extends FileGroupRea
     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,
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 b8a1a243e5d4..4f3805405061 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
@@ -39,6 +39,7 @@ import org.apache.hudi.common.table.read.InputSplit;
 import org.apache.hudi.common.table.read.ReaderParameters;
 import org.apache.hudi.common.table.read.UpdateProcessor;
 import org.apache.hudi.common.util.Option;
+import org.apache.hudi.common.util.StringUtils;
 import org.apache.hudi.common.util.VisibleForTesting;
 import org.apache.hudi.common.util.collection.ClosableIterator;
 import org.apache.hudi.common.util.collection.Pair;
@@ -473,7 +474,8 @@ public class LsmFileGroupRecordIterator<T> implements 
ClosableIterator<BufferedR
     private int compare(int leftIndex, int rightIndex) {
       SortedRunReader<T> left = leaves.get(leftIndex);
       SortedRunReader<T> right = leaves.get(rightIndex);
-      int keyCompare = 
left.current.getRecordKey().compareTo(right.current.getRecordKey());
+      int keyCompare = StringUtils.compareUtf8Bytes(
+          left.current.getRecordKey(), right.current.getRecordKey());
       if (keyCompare != 0) {
         return keyCompare;
       }
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/LsmReaderUtils.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/LsmReaderUtils.java
new file mode 100644
index 000000000000..4052c817df20
--- /dev/null
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/LsmReaderUtils.java
@@ -0,0 +1,41 @@
+/*
+ * 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.common.table.read.lsm;
+
+import org.apache.hudi.common.config.HoodieReaderConfig;
+import org.apache.hudi.common.table.HoodieTableConfig;
+
+/**
+ * Utilities for selecting the LSM file group reader.
+ */
+public final class LsmReaderUtils {
+
+  private LsmReaderUtils() {
+  }
+
+  /**
+   * Returns whether the file group can be read with the LSM reader for the 
configured merge type.
+   */
+  public static boolean shouldUseLsmReader(HoodieTableConfig tableConfig, 
String mergeType) {
+    // The LSM reader collapses all sorted versions of a key. Skip-merge 
queries intentionally
+    // expose those versions independently, so retain the classic unmerged 
reader for that mode.
+    return !HoodieReaderConfig.REALTIME_SKIP_MERGE.equalsIgnoreCase(mergeType)
+        && tableConfig.isLSMTreeStorageLayout();
+  }
+}
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadata.java
 
b/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadata.java
index 3d31a2ccfed3..43fd598e4f35 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadata.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadata.java
@@ -22,6 +22,7 @@ import org.apache.hudi.avro.model.HoodieMetadataRecord;
 import org.apache.hudi.common.avro.HoodieAvroReaderContext;
 import org.apache.hudi.common.config.HoodieConfig;
 import org.apache.hudi.common.config.HoodieMetadataConfig;
+import org.apache.hudi.common.config.HoodieReaderConfig;
 import org.apache.hudi.common.config.TypedProperties;
 import org.apache.hudi.common.data.HoodieData;
 import org.apache.hudi.common.data.HoodieListData;
@@ -34,7 +35,6 @@ import org.apache.hudi.common.expression.Expression;
 import org.apache.hudi.common.expression.Literal;
 import org.apache.hudi.common.expression.Predicate;
 import org.apache.hudi.common.expression.Predicates;
-import org.apache.hudi.common.fs.FSUtils;
 import org.apache.hudi.common.function.SerializableBiFunction;
 import org.apache.hudi.common.function.SerializableFunction;
 import org.apache.hudi.common.function.SerializableFunctionUnchecked;
@@ -54,6 +54,7 @@ import 
org.apache.hudi.common.table.read.HoodieFileGroupReader;
 import org.apache.hudi.common.table.read.buffer.FileGroupRecordBufferLoader;
 import 
org.apache.hudi.common.table.read.buffer.ReusableFileGroupRecordBufferLoader;
 import org.apache.hudi.common.table.read.lsm.HoodieLsmFileGroupReader;
+import org.apache.hudi.common.table.read.lsm.LsmReaderUtils;
 import org.apache.hudi.common.table.timeline.HoodieInstant;
 import org.apache.hudi.common.table.view.HoodieTableFileSystemView;
 import org.apache.hudi.common.util.ConfigUtils;
@@ -574,7 +575,9 @@ public class HoodieBackedTableMetadata extends 
BaseTableMetadata {
 
     // If reuse is enabled and full scan is allowed for the partition, we can 
reuse the file readers for base files and the reader context for the log files.
     boolean shouldReuse = reuse && 
isFullScanAllowedForPartition(fileSlice.getPartitionPath());
-    boolean useLsmReader = !shouldReuse && 
shouldUseLsmReader(metadataMetaClient, fileSlice);
+    boolean useLsmReader = !shouldReuse
+        && LsmReaderUtils.shouldUseLsmReader(
+            metadataMetaClient.getTableConfig(), 
HoodieReaderConfig.REALTIME_PAYLOAD_COMBINE);
     Map<StoragePath, HoodieAvroFileReader> baseFileReaders = 
Collections.emptyMap();
     ReusableFileGroupRecordBufferLoader<IndexedRecord> recordBufferLoader = 
null;
     TypedProperties fileGroupReaderProps = 
ConfigUtils.buildFileGroupReaderProperties(metadataConfig, shouldReuse);
@@ -639,11 +642,6 @@ public class HoodieBackedTableMetadata extends 
BaseTableMetadata {
     }
   }
 
-  private static boolean shouldUseLsmReader(HoodieTableMetaClient metaClient, 
FileSlice fileSlice) {
-    return metaClient.getTableConfig().isLSMTreeStorageLayout()
-        && fileSlice.getLogFiles().allMatch(logFile -> 
FSUtils.isNativeLogFile(logFile.getFileName()));
-  }
-
   private ReusableFileGroupRecordBufferLoader<IndexedRecord> 
buildReusableRecordBufferLoader(FileSlice fileSlice, String 
latestMetadataInstantTime,
                                                                                
              Option<InstantRange> instantRangeOption) {
     // initialize without any filters
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestHoodieLsmFileGroupReader.java
 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestHoodieLsmFileGroupReader.java
index 67d003d36ed4..73ad2c1b3a8a 100644
--- 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestHoodieLsmFileGroupReader.java
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestHoodieLsmFileGroupReader.java
@@ -227,6 +227,29 @@ class TestHoodieLsmFileGroupReader {
     }
   }
 
+  @Test
+  void testBaseFileOnlyPathPreservesDuplicateRecordKeys() throws IOException {
+    HoodieReaderContext<IndexedRecord> readerContext = spy(context());
+    StoragePathInfo baseFilePathInfo = 
pathInfo("/tmp/file1_1-0-1_001.parquet");
+    doReturn(ClosableIterator.wrap(Arrays.asList(
+        recordWithCommitTime("001", "a", "first", 1),
+        recordWithCommitTime("001", "a", "second", 2)).iterator()))
+        .when(readerContext).getFileRecordIterator(
+            eq(baseFilePathInfo), anyLong(), anyLong(), 
any(HoodieSchema.class),
+            any(HoodieSchema.class), any(HoodieStorage.class));
+
+    try (HoodieLsmFileGroupReader<IndexedRecord> reader = reader(
+        readerContext, Option.of(new HoodieBaseFile(baseFilePathInfo)), 
Collections.emptyList(), 0L);
+         ClosableIterator<IndexedRecord> iterator = 
reader.getClosableIterator()) {
+      List<IndexedRecord> records = drain(iterator);
+      assertEquals(2, records.size());
+      assertEquals("a", records.get(0).get(1).toString());
+      assertEquals("first", records.get(0).get(2).toString());
+      assertEquals("a", records.get(1).get(1).toString());
+      assertEquals("second", records.get(1).get(2).toString());
+    }
+  }
+
   @Test
   void testMetadataTableBaseFileIsNotFilteredByInstantRange() throws 
IOException {
     when(metaClient.getBasePath()).thenReturn(
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestLsmFileGroupRecordIterator.java
 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestLsmFileGroupRecordIterator.java
index e60d5511d1ac..edd5a83284b4 100644
--- 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestLsmFileGroupRecordIterator.java
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestLsmFileGroupRecordIterator.java
@@ -196,6 +196,21 @@ class TestLsmFileGroupRecordIterator {
         "key3:log2-key3"), drain(loserTree));
   }
 
+  @Test
+  void testLoserTreeUsesUtf8Ordering() {
+    String bmpPrivateUseKey = new String(Character.toChars(0xE000));
+    String supplementaryKey = new String(Character.toChars(0x20000));
+
+    LsmFileGroupRecordIterator.LoserTree<String> loserTree =
+        new LsmFileGroupRecordIterator.LoserTree<>(
+            Arrays.asList(
+                sortedRunReader(0, record(bmpPrivateUseKey, "bmp")),
+                sortedRunReader(1, record(supplementaryKey, 
"supplementary"))));
+    assertEquals(Arrays.asList(
+        bmpPrivateUseKey + ":bmp",
+        supplementaryKey + ":supplementary"), drain(loserTree));
+  }
+
   @Test
   void testSelectDirectLogReadersPrioritizesDeletesThenSmallFiles() {
     List<LsmFileGroupRecordIterator.LogReaderSpec> logReaderSpecs = 
Arrays.asList(
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestLsmReaderUtils.java
 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestLsmReaderUtils.java
new file mode 100644
index 000000000000..6f1949902eed
--- /dev/null
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestLsmReaderUtils.java
@@ -0,0 +1,43 @@
+/*
+ * 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.common.table.read.lsm;
+
+import org.apache.hudi.common.config.HoodieReaderConfig;
+import org.apache.hudi.common.table.HoodieTableConfig;
+
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class TestLsmReaderUtils {
+
+  @Test
+  void testShouldUseLsmReader() {
+    HoodieTableConfig tableConfig = new HoodieTableConfig();
+    assertFalse(LsmReaderUtils.shouldUseLsmReader(
+        tableConfig, HoodieReaderConfig.REALTIME_PAYLOAD_COMBINE));
+
+    tableConfig.setValue(HoodieTableConfig.TABLE_STORAGE_LAYOUT, 
HoodieTableConfig.TableStorageLayout.LSM_TREE.configValue());
+    assertTrue(LsmReaderUtils.shouldUseLsmReader(
+        tableConfig, HoodieReaderConfig.REALTIME_PAYLOAD_COMBINE));
+    assertFalse(LsmReaderUtils.shouldUseLsmReader(
+        tableConfig, HoodieReaderConfig.REALTIME_SKIP_MERGE));
+  }
+}
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FormatUtils.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FormatUtils.java
index d451d61a8894..ac067d36f412 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FormatUtils.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FormatUtils.java
@@ -21,7 +21,6 @@ package org.apache.hudi.table.format;
 import org.apache.hudi.common.config.ConfigProperty;
 import org.apache.hudi.common.config.HoodieReaderConfig;
 import org.apache.hudi.common.config.TypedProperties;
-import org.apache.hudi.common.fs.FSUtils;
 import org.apache.hudi.common.model.FileSlice;
 import org.apache.hudi.common.schema.HoodieSchema;
 import org.apache.hudi.common.schema.HoodieSchemaField;
@@ -31,6 +30,7 @@ import org.apache.hudi.common.table.log.InstantRange;
 import org.apache.hudi.common.table.read.HoodieFileGroupReader;
 import org.apache.hudi.common.table.read.HoodieRecordReader;
 import org.apache.hudi.common.table.read.lsm.HoodieLsmFileGroupReader;
+import org.apache.hudi.common.table.read.lsm.LsmReaderUtils;
 import org.apache.hudi.common.util.DefaultSizeEstimator;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.collection.ClosableIterator;
@@ -125,9 +125,8 @@ public class FormatUtils {
   /**
    * Creates the record reader matching the physical layout of the file slice.
    *
-   * <p>Pure native-log file slices in an LSM-tree table use the sorted LSM 
reader. Mixed file
-   * slices containing legacy inline logs continue to use the general 
file-group reader so table
-   * upgrades remain readable.</p>
+   * <p>LSM-tree tables use the sorted LSM reader, except for skip-merge 
queries which use the
+   * general file-group reader to expose record versions independently.</p>
    */
   public static HoodieRecordReader<RowData> createRecordReader(
       HoodieTableMetaClient metaClient,
@@ -141,7 +140,7 @@ public class FormatUtils {
       boolean emitDelete,
       List<ExpressionPredicates.Predicate> predicates,
       Option<InstantRange> instantRangeOption) {
-    if (!shouldUseLsmReader(metaClient, fileSlice, mergeType)) {
+    if (!LsmReaderUtils.shouldUseLsmReader(metaClient.getTableConfig(), 
mergeType)) {
       return createFileGroupReader(metaClient, writeConfig, 
internalSchemaManager, fileSlice,
           tableSchema, requiredSchema, latestInstant, mergeType, emitDelete, 
predicates, instantRangeOption);
     }
@@ -172,18 +171,6 @@ public class FormatUtils {
         .build();
   }
 
-  static boolean shouldUseLsmReader(HoodieTableMetaClient metaClient, 
FileSlice fileSlice) {
-    return metaClient.getTableConfig().isLSMTreeStorageLayout()
-        && fileSlice.getLogFiles().allMatch(logFile -> 
FSUtils.isNativeLogFile(logFile.getFileName()));
-  }
-
-  static boolean shouldUseLsmReader(HoodieTableMetaClient metaClient, 
FileSlice fileSlice, String mergeType) {
-    // The LSM reader collapses all sorted versions of a key. Skip-merge 
queries intentionally
-    // expose those versions independently, so retain the classic unmerged 
reader for that mode.
-    return !HoodieReaderConfig.REALTIME_SKIP_MERGE.equalsIgnoreCase(mergeType)
-        && shouldUseLsmReader(metaClient, fileSlice);
-  }
-
   /**
    * Create a {@link HoodieFileGroupReader}.
    *
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/TestFormatUtils.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/TestFormatUtils.java
index 515856f3736a..3cf076a5818e 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/TestFormatUtils.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/TestFormatUtils.java
@@ -19,58 +19,18 @@
 
 package org.apache.hudi.table.format;
 
-import org.apache.hudi.common.config.HoodieReaderConfig;
-import org.apache.hudi.common.model.FileSlice;
-import org.apache.hudi.common.model.HoodieLogFile;
-import org.apache.hudi.common.table.HoodieTableConfig;
-import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.util.Option;
-import org.apache.hudi.storage.StoragePath;
 
 import org.apache.flink.configuration.Configuration;
 import org.junit.jupiter.api.Test;
 
 import static 
org.apache.hudi.common.util.TestConfigUtils.TEST_BOOLEAN_CONFIG_PROPERTY;
 import static org.junit.jupiter.api.Assertions.assertEquals;
-import static org.junit.jupiter.api.Assertions.assertFalse;
-import static org.junit.jupiter.api.Assertions.assertTrue;
-import static org.mockito.Mockito.mock;
-import static org.mockito.Mockito.when;
 
 /**
  * Tests {@link FormatUtils}
  */
 public class TestFormatUtils {
-  private static final String INLINE_LOG_PATH = 
"file:///tmp/.file-id_100.log.1_1-0-1";
-  private static final String NATIVE_LOG_PATH = 
"file:///tmp/file-id_1-0-1_100_1.log.parquet";
-
-  @Test
-  public void testLsmReaderSelection() {
-    HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
-    HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
-    when(metaClient.getTableConfig()).thenReturn(tableConfig);
-
-    when(tableConfig.isLSMTreeStorageLayout()).thenReturn(false);
-    assertFalse(FormatUtils.shouldUseLsmReader(metaClient, 
fileSlice(NATIVE_LOG_PATH)));
-
-    when(tableConfig.isLSMTreeStorageLayout()).thenReturn(true);
-    assertTrue(FormatUtils.shouldUseLsmReader(metaClient, fileSlice()));
-    assertTrue(FormatUtils.shouldUseLsmReader(metaClient, 
fileSlice(NATIVE_LOG_PATH)));
-    assertFalse(FormatUtils.shouldUseLsmReader(metaClient, 
fileSlice(INLINE_LOG_PATH, NATIVE_LOG_PATH)));
-    assertTrue(FormatUtils.shouldUseLsmReader(
-        metaClient, fileSlice(NATIVE_LOG_PATH), 
HoodieReaderConfig.REALTIME_PAYLOAD_COMBINE));
-    assertFalse(FormatUtils.shouldUseLsmReader(
-        metaClient, fileSlice(NATIVE_LOG_PATH), 
HoodieReaderConfig.REALTIME_SKIP_MERGE));
-  }
-
-  private static FileSlice fileSlice(String... logPaths) {
-    FileSlice fileSlice = new FileSlice("partition", "100", "file-id");
-    for (String logPath : logPaths) {
-      fileSlice.addLogFile(new HoodieLogFile(new StoragePath(logPath)));
-    }
-    return fileSlice;
-  }
-
   @Test
   public void testGetRawValueWithAltKeys() {
     Configuration flinkConf = new Configuration();
diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieMergeOnReadRDDV2.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieMergeOnReadRDDV2.scala
index bb321616a1eb..3ae45815c679 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieMergeOnReadRDDV2.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieMergeOnReadRDDV2.scala
@@ -30,7 +30,8 @@ import org.apache.hudi.common.schema.HoodieSchema
 import org.apache.hudi.common.table.HoodieTableMetaClient
 import org.apache.hudi.common.table.log.InstantRange
 import org.apache.hudi.common.table.log.InstantRange.RangeType
-import org.apache.hudi.common.table.read.HoodieFileGroupReader
+import org.apache.hudi.common.table.read.{HoodieFileGroupReader, 
HoodieRecordReader}
+import org.apache.hudi.common.table.read.lsm.{HoodieLsmFileGroupReader, 
LsmReaderUtils}
 import org.apache.hudi.common.util.{Option => HOption}
 import org.apache.hudi.common.util.collection.ClosableIterator
 import 
org.apache.hudi.hadoop.utils.HoodieRealtimeRecordReaderUtils.getMaxCompactionMemoryInBytes
@@ -191,18 +192,34 @@ class HoodieMergeOnReadRDDV2(@transient sc: SparkContext,
         } else {
           val readerContext = new 
SparkFileFormatInternalRowReaderContext(fileGroupBaseFileReader.value, 
optionalFilters,
             Seq.empty, storageConf, metaClient.getTableConfig)
-          val fileGroupReader = HoodieFileGroupReader.builder()
-            .withReaderContext(readerContext)
-            .withHoodieTableMetaClient(metaClient)
-            .withLatestCommitTime(tableState.latestCommitTimestamp.orNull)
-            .withLogFiles(logFiles.stream())
-            .withBaseFileOption(baseFileOption)
-            .withPartitionPath(partitionPath)
-            .withProps(properties)
-            .withDataSchema(tableSchema.schema)
-            .withRequestedSchema(requiredSchema.schema)
-            
.withInternalSchemaOpt(HOption.ofNullable(tableSchema.internalSchema.orNull))
-            .build()
+          val fileGroupReader: HoodieRecordReader[InternalRow] =
+            if (LsmReaderUtils.shouldUseLsmReader(metaClient.getTableConfig, 
mergeType)) {
+              HoodieLsmFileGroupReader.builder[InternalRow]()
+                .withReaderContext(readerContext)
+                .withHoodieTableMetaClient(metaClient)
+                .withLatestCommitTime(tableState.latestCommitTimestamp.orNull)
+                .withLogFiles(logFiles.stream())
+                .withBaseFileOption(baseFileOption)
+                .withPartitionPath(partitionPath)
+                .withProps(properties)
+                .withDataSchema(tableSchema.schema)
+                .withRequestedSchema(requiredSchema.schema)
+                
.withInternalSchemaOpt(HOption.ofNullable(tableSchema.internalSchema.orNull))
+                .build()
+            } else {
+              HoodieFileGroupReader.builder[InternalRow]()
+                .withReaderContext(readerContext)
+                .withHoodieTableMetaClient(metaClient)
+                .withLatestCommitTime(tableState.latestCommitTimestamp.orNull)
+                .withLogFiles(logFiles.stream())
+                .withBaseFileOption(baseFileOption)
+                .withPartitionPath(partitionPath)
+                .withProps(properties)
+                .withDataSchema(tableSchema.schema)
+                .withRequestedSchema(requiredSchema.schema)
+                
.withInternalSchemaOpt(HOption.ofNullable(tableSchema.internalSchema.orNull))
+                .build()
+            }
           convertCloseableIterator(fileGroupReader.getClosableIterator)
         }
     }
diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieSparkSqlWriter.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieSparkSqlWriter.scala
index 37f8e569a0e7..c72fcc3ceb8c 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieSparkSqlWriter.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieSparkSqlWriter.scala
@@ -299,6 +299,7 @@ class HoodieSparkSqlWriterInternal {
           .setTableType(tableType)
           .setTableVersion(tableVersion)
           .setTableFormat(tableFormat)
+          
.setTableStorageLayout(hoodieConfig.getStringOrDefault(HoodieTableConfig.TABLE_STORAGE_LAYOUT))
           .setDatabaseName(databaseName)
           .setTableName(tblName)
           .setBaseFileFormat(baseFileFormat)
@@ -768,6 +769,7 @@ class HoodieSparkSqlWriterInternal {
           .setRecordKeyFields(recordKeyFields)
           .setTableVersion(tableVersion)
           .setTableFormat(tableFormat)
+          
.setTableStorageLayout(hoodieConfig.getStringOrDefault(HoodieTableConfig.TABLE_STORAGE_LAYOUT))
           .setArchiveLogFolder(archiveLogFolder)
           .setPayloadClassName(payloadClass)
           
.setRecordMergeMode(RecordMergeMode.getValue(hoodieConfig.getString(HoodieWriteConfig.RECORD_MERGE_MODE)))
diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/execution/datasources/parquet/HoodieFileGroupReaderBasedFileFormat.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/execution/datasources/parquet/HoodieFileGroupReaderBasedFileFormat.scala
index 175239d9d7c2..c776172d3896 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/execution/datasources/parquet/HoodieFileGroupReaderBasedFileFormat.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/execution/datasources/parquet/HoodieFileGroupReaderBasedFileFormat.scala
@@ -21,7 +21,7 @@ import org.apache.hudi.{HoodieFileIndex, 
HoodiePartitionCDCFileGroupMapping, Hoo
 import org.apache.hudi.cdc.{CDCFileGroupIterator, HoodieCDCFileGroupSplit, 
HoodieCDCFileIndex}
 import org.apache.hudi.client.common.HoodieSparkEngineContext
 import org.apache.hudi.client.utils.SparkInternalSchemaConverter
-import org.apache.hudi.common.config.{HoodieMemoryConfig, TypedProperties}
+import org.apache.hudi.common.config.{HoodieMemoryConfig, HoodieReaderConfig, 
TypedProperties}
 import org.apache.hudi.common.fs.FSUtils
 import org.apache.hudi.common.model.HoodieFileFormat
 import org.apache.hudi.common.schema.HoodieSchema
@@ -29,8 +29,9 @@ import org.apache.hudi.common.schema.HoodieSchemaRepair
 import org.apache.hudi.common.schema.HoodieSchemaUtils
 import org.apache.hudi.common.schema.internal.InternalSchema
 import org.apache.hudi.common.table.{HoodieTableConfig, HoodieTableMetaClient, 
ParquetTableSchemaResolver}
-import org.apache.hudi.common.table.read.HoodieFileGroupReader
-import org.apache.hudi.common.util.{Option => HOption}
+import org.apache.hudi.common.table.read.{HoodieFileGroupReader, 
HoodieRecordReader}
+import org.apache.hudi.common.table.read.lsm.{HoodieLsmFileGroupReader, 
LsmReaderUtils}
+import org.apache.hudi.common.util.{ConfigUtils, Option => HOption}
 import org.apache.hudi.common.util.collection.ClosableIterator
 import org.apache.hudi.data.CloseableIteratorListener
 import org.apache.hudi.exception.HoodieNotSupportedException
@@ -314,21 +315,41 @@ class HoodieFileGroupReaderBasedFileFormat(tablePath: 
String,
               } else {
                 0
               }
-              val reader = HoodieFileGroupReader.builder()
-                .withReaderContext(readerContext)
-                .withHoodieTableMetaClient(metaClient)
-                .withLatestCommitTime(queryTimestamp)
-                .withBaseFileOption(fileSlice.getBaseFile)
-                .withLogFiles(fileSlice.getLogFiles)
-                .withPartitionPath(fileSlice.getPartitionPath)
-                .withDataSchema(dataSchema)
-                .withRequestedSchema(requestedSchema)
-                .withInternalSchemaOpt(internalSchemaOpt)
-                .withProps(props)
-                .withStart(file.start)
-                .withLength(baseFileLength)
-                .withShouldUseRecordPosition(shouldUseRecordPosition)
-                .build()
+              val reader: HoodieRecordReader[InternalRow] =
+                if (LsmReaderUtils.shouldUseLsmReader(
+                  metaClient.getTableConfig,
+                  ConfigUtils.getStringWithAltKeys(props, 
HoodieReaderConfig.MERGE_TYPE, true))) {
+                  HoodieLsmFileGroupReader.builder[InternalRow]()
+                    .withReaderContext(readerContext)
+                    .withHoodieTableMetaClient(metaClient)
+                    .withLatestCommitTime(queryTimestamp)
+                    .withBaseFileOption(fileSlice.getBaseFile)
+                    .withLogFiles(fileSlice.getLogFiles)
+                    .withPartitionPath(fileSlice.getPartitionPath)
+                    .withDataSchema(dataSchema)
+                    .withRequestedSchema(requestedSchema)
+                    .withInternalSchemaOpt(internalSchemaOpt)
+                    .withProps(props)
+                    .withStart(file.start)
+                    .withLength(baseFileLength)
+                    .build()
+                } else {
+                  HoodieFileGroupReader.builder[InternalRow]()
+                    .withReaderContext(readerContext)
+                    .withHoodieTableMetaClient(metaClient)
+                    .withLatestCommitTime(queryTimestamp)
+                    .withBaseFileOption(fileSlice.getBaseFile)
+                    .withLogFiles(fileSlice.getLogFiles)
+                    .withPartitionPath(fileSlice.getPartitionPath)
+                    .withDataSchema(dataSchema)
+                    .withRequestedSchema(requestedSchema)
+                    .withInternalSchemaOpt(internalSchemaOpt)
+                    .withProps(props)
+                    .withStart(file.start)
+                    .withLength(baseFileLength)
+                    .withShouldUseRecordPosition(shouldUseRecordPosition)
+                    .build()
+                }
               // Append partition values to rows and project to output schema
               appendPartitionAndProject(
                 reader.getClosableIterator,
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestMORDataSource.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestMORDataSource.scala
index 81d049d43243..0316423c10e8 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestMORDataSource.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestMORDataSource.scala
@@ -38,7 +38,7 @@ import 
org.apache.hudi.metadata.HoodieTableMetadataUtil.{metadataPartitionExists
 import org.apache.hudi.storage.{StoragePath, StoragePathInfo}
 import org.apache.hudi.table.action.compact.CompactionTriggerStrategy
 import org.apache.hudi.table.upgrade.TestUpgradeDowngrade.getFixtureName
-import org.apache.hudi.testutils.{DataSourceTestUtils, 
HoodieSparkClientTestBase}
+import org.apache.hudi.testutils.{DataSourceTestUtils, HoodieClientTestUtils, 
HoodieSparkClientTestBase}
 import org.apache.hudi.util.JFunction
 
 import org.apache.commons.io.FileUtils
@@ -104,6 +104,330 @@ class TestMORDataSource extends HoodieSparkClientTestBase 
with SparkDatasetMixin
         JFunction.toJavaConsumer((receiver: SparkSessionExtensions) => new 
HoodieSparkSessionExtension().apply(receiver)))
     )
 
+  @Test
+  def testLsmUpsertUsesUtf8Ordering(): Unit = {
+    val fullWidthAKey = "A-key"
+    val emojiFaceKey = "😀-key"
+    val _spark = spark
+    import _spark.implicits._
+
+    Seq(HoodieTableType.COPY_ON_WRITE, HoodieTableType.MERGE_ON_READ).foreach 
{ tableType =>
+      val tablePath = s"${basePath}_${tableType.name.toLowerCase}_lsm"
+      val options = Map[String, String](
+        DataSourceWriteOptions.TABLE_TYPE.key -> tableType.name,
+        DataSourceWriteOptions.OPERATION.key -> UPSERT_OPERATION_OPT_VAL,
+        DataSourceWriteOptions.RECORDKEY_FIELD.key -> "id",
+        DataSourceWriteOptions.PARTITIONPATH_FIELD.key -> "",
+        DataSourceWriteOptions.KEYGENERATOR_CLASS_NAME.key -> 
"org.apache.hudi.keygen.NonpartitionedKeyGenerator",
+        HoodieTableConfig.ORDERING_FIELDS.key -> "ts",
+        HoodieTableConfig.TABLE_STORAGE_LAYOUT.key -> 
HoodieTableConfig.TableStorageLayout.LSM_TREE.configValue,
+        HoodieWriteConfig.TBL_NAME.key -> 
s"hoodie_lsm_${tableType.name.toLowerCase}",
+        HoodieMetadataConfig.ENABLE.key -> "false",
+        HoodieCompactionConfig.INLINE_COMPACT.key -> "false",
+        HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet",
+        "hoodie.insert.shuffle.parallelism" -> "1",
+        "hoodie.upsert.shuffle.parallelism" -> "1")
+      val (writeOpts, readOpts) = getWriterReaderOpts(HoodieRecordType.AVRO, 
options)
+
+      Seq(
+        (fullWidthAKey, "full-width-a-v1", 1L),
+        (emojiFaceKey, "emoji-v1", 1L))
+        .toDF("id", "value", "ts")
+        .repartition(1)
+        .write.format("hudi")
+        .options(writeOpts)
+        .mode(SaveMode.Append)
+        .save(tablePath)
+
+      val metaClient = HoodieTableMetaClient.builder()
+        .setBasePath(tablePath)
+        .setConf(storageConf.newInstance())
+        .build()
+      assertTrue(metaClient.getTableConfig.isLSMTreeStorageLayout)
+
+      // All LSM sorted runs use UTF-8 byte order, where U+FF21 (A) sorts 
before U+1F600 (😀).
+      val physicalBaseFileKeys = spark.read.parquet(tablePath)
+        .select(HoodieRecord.RECORD_KEY_METADATA_FIELD)
+        .collect()
+        .map(_.getString(0))
+        .toList
+      assertEquals(Seq(fullWidthAKey, emojiFaceKey), physicalBaseFileKeys)
+
+      Seq(
+        (fullWidthAKey, "full-width-a-v2", 2L),
+        (emojiFaceKey, "emoji-v2", 2L))
+        .toDF("id", "value", "ts")
+        .repartition(1)
+        .write.format("hudi")
+        .options(writeOpts)
+        .mode(SaveMode.Append)
+        .save(tablePath)
+
+      val actual = spark.read.format("hudi")
+        .options(readOpts)
+        .load(tablePath)
+        .select("id", "value")
+        .collect()
+        .map(row => row.getString(0) -> row.getString(1))
+        .toMap
+      assertEquals(Map(
+        fullWidthAKey -> "full-width-a-v2",
+        emojiFaceKey -> "emoji-v2"), actual)
+    }
+  }
+
+  @Test
+  def testLsmBaseFileOnlyReadPreservesDuplicateKeys(): Unit = {
+    val duplicateKey = "duplicate-key"
+    val tablePath = s"${basePath}_mor_lsm_base_file_only_duplicates"
+    val _spark = spark
+    import _spark.implicits._
+
+    val options = Map[String, String](
+      DataSourceWriteOptions.TABLE_TYPE.key -> 
HoodieTableType.MERGE_ON_READ.name,
+      DataSourceWriteOptions.OPERATION.key -> INSERT_OPERATION_OPT_VAL,
+      DataSourceWriteOptions.INSERT_DUP_POLICY.key -> 
DataSourceWriteOptions.NONE_INSERT_DUP_POLICY,
+      DataSourceWriteOptions.RECORDKEY_FIELD.key -> "id",
+      DataSourceWriteOptions.PARTITIONPATH_FIELD.key -> "",
+      DataSourceWriteOptions.KEYGENERATOR_CLASS_NAME.key -> 
"org.apache.hudi.keygen.NonpartitionedKeyGenerator",
+      HoodieTableConfig.ORDERING_FIELDS.key -> "ts",
+      HoodieWriteConfig.COMBINE_BEFORE_INSERT.key -> "false",
+      HoodieWriteConfig.TBL_NAME.key -> 
"hoodie_mor_lsm_base_file_only_duplicates",
+      HoodieMetadataConfig.ENABLE.key -> "false",
+      HoodieCompactionConfig.INLINE_COMPACT.key -> "false",
+      "hoodie.insert.shuffle.parallelism" -> "1")
+    val (writeOpts, readOpts) = getWriterReaderOpts(HoodieRecordType.AVRO, 
options)
+
+    Seq(
+      (duplicateKey, "first", 1L),
+      (duplicateKey, "second", 2L))
+      .toDF("id", "value", "ts")
+      .repartition(1)
+      .write.format("hudi")
+      .options(writeOpts)
+      .mode(SaveMode.Append)
+      .save(tablePath)
+
+    val physicalRows = spark.read.parquet(tablePath)
+      .select("id", "value")
+      .collect()
+    assertEquals(2, physicalRows.length)
+    assertEquals(Set("first", "second"), 
physicalRows.map(_.getString(1)).toSet)
+
+    val dataFiles = storage.listDirectEntries(new 
StoragePath(tablePath)).asScala
+      .filterNot(_.isDirectory)
+    assertEquals(1, 
dataFiles.count(_.getPath.getName.endsWith(HoodieFileFormat.PARQUET.getFileExtension)))
+    assertFalse(dataFiles.exists(pathInfo => 
org.apache.hudi.common.fs.FSUtils.isLogFile(pathInfo.getPath)))
+
+    // Build the duplicate-bearing base file with the supported default-layout 
insert path, then
+    // switch only the test fixture to LSM so this test does not imply LSM 
insert support.
+    val metaClient = HoodieTableMetaClient.builder()
+      .setBasePath(tablePath)
+      .setConf(storageConf.newInstance())
+      .build()
+    assertFalse(metaClient.getTableConfig.isLSMTreeStorageLayout)
+    val tableProps = metaClient.getTableConfig.getProps
+    tableProps.setProperty(
+      HoodieTableConfig.TABLE_STORAGE_LAYOUT.key,
+      HoodieTableConfig.TableStorageLayout.LSM_TREE.configValue)
+    HoodieTableConfig.update(metaClient.getStorage, metaClient.getMetaPath, 
tableProps)
+    metaClient.reloadTableConfig()
+    assertTrue(metaClient.getTableConfig.isLSMTreeStorageLayout)
+
+    // With no log records to merge, the LSM reader delegates directly to the 
base-file iterator and
+    // preserve duplicate record keys just like the classic Spark file-group 
reader.
+    val snapshotRows = spark.read.format("hudi")
+      .options(readOpts)
+      .load(tablePath)
+      .select("id", "value")
+      .collect()
+    assertEquals(2, snapshotRows.length)
+    assertTrue(snapshotRows.forall(_.getString(0) == duplicateKey))
+    assertEquals(Set("first", "second"), 
snapshotRows.map(_.getString(1)).toSet)
+  }
+
+  @Test
+  def testLsmIncrementalAndSkipMergeReads(): Unit = {
+    val fullWidthAKey = "A-key"
+    val emojiFaceKey = "😀-key"
+    val tablePath = s"${basePath}_mor_lsm_read_paths"
+    val _spark = spark
+    import _spark.implicits._
+
+    val options = Map[String, String](
+      DataSourceWriteOptions.TABLE_TYPE.key -> 
HoodieTableType.MERGE_ON_READ.name,
+      DataSourceWriteOptions.OPERATION.key -> UPSERT_OPERATION_OPT_VAL,
+      DataSourceWriteOptions.RECORDKEY_FIELD.key -> "id",
+      DataSourceWriteOptions.PARTITIONPATH_FIELD.key -> "",
+      DataSourceWriteOptions.KEYGENERATOR_CLASS_NAME.key -> 
"org.apache.hudi.keygen.NonpartitionedKeyGenerator",
+      HoodieTableConfig.ORDERING_FIELDS.key -> "ts",
+      HoodieTableConfig.TABLE_STORAGE_LAYOUT.key -> 
HoodieTableConfig.TableStorageLayout.LSM_TREE.configValue,
+      HoodieWriteConfig.TBL_NAME.key -> "hoodie_mor_lsm_read_paths",
+      HoodieMetadataConfig.ENABLE.key -> "false",
+      HoodieCompactionConfig.INLINE_COMPACT.key -> "false",
+      HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet",
+      "hoodie.insert.shuffle.parallelism" -> "1",
+      "hoodie.upsert.shuffle.parallelism" -> "1")
+    val (writeOpts, readOpts) = getWriterReaderOpts(HoodieRecordType.AVRO, 
options)
+
+    Seq(
+      (fullWidthAKey, "full-width-a-v1", 1L),
+      (emojiFaceKey, "emoji-v1", 1L))
+      .toDF("id", "value", "ts")
+      .repartition(1)
+      .write.format("hudi")
+      .options(writeOpts)
+      .mode(SaveMode.Append)
+      .save(tablePath)
+    val firstCompletionTime = 
DataSourceTestUtils.latestDeltaCommitCompletionTime(storage, tablePath)
+
+    // Updating only the first physical key makes the base and log sorted runs 
have different key
+    // sets, exposing an ordering mismatch while the readers merge them.
+    Seq((fullWidthAKey, "full-width-a-v2", 2L))
+      .toDF("id", "value", "ts")
+      .repartition(1)
+      .write.format("hudi")
+      .options(writeOpts)
+      .mode(SaveMode.Append)
+      .save(tablePath)
+
+    val expectedSnapshot = Map(
+      fullWidthAKey -> "full-width-a-v2",
+      emojiFaceKey -> "emoji-v1")
+    val snapshotRows = spark.read.format("hudi")
+      .options(readOpts)
+      .load(tablePath)
+      .select("id", "value")
+      .collect()
+    assertEquals(2, snapshotRows.length)
+    val actualSnapshot = snapshotRows
+      .map(row => row.getString(0) -> row.getString(1))
+      .toMap
+    assertEquals(expectedSnapshot, actualSnapshot)
+
+    val incrementalOptions = readOpts ++ Map(
+      DataSourceReadOptions.QUERY_TYPE.key -> 
DataSourceReadOptions.QUERY_TYPE_INCREMENTAL_OPT_VAL,
+      DataSourceReadOptions.START_COMMIT.key -> firstCompletionTime)
+    val expectedIncremental = Map(fullWidthAKey -> "full-width-a-v2")
+    val incrementalRows = spark.read.format("hudi")
+      .options(incrementalOptions)
+      .load(tablePath)
+      .select("id", "value")
+      .collect()
+    assertEquals(1, incrementalRows.length)
+    val actualIncremental = incrementalRows
+      .map(row => row.getString(0) -> row.getString(1))
+      .toMap
+    assertEquals(expectedIncremental, actualIncremental)
+
+    // Skip-merge intentionally stays on the classic reader and exposes the 
base and log versions.
+    val skipMergeRows = spark.read.format("hudi")
+      .options(readOpts)
+      .option(DataSourceReadOptions.REALTIME_MERGE.key, 
DataSourceReadOptions.REALTIME_SKIP_MERGE_OPT_VAL)
+      .load(tablePath)
+      .select("id", "value")
+      .collect()
+    assertEquals(3, skipMergeRows.length)
+    val skipMergeVersions = skipMergeRows
+      .map(row => row.getString(0) -> row.getString(1))
+      .groupBy(_._1)
+      .mapValues(_.map(_._2).toSet)
+    assertEquals(Set("full-width-a-v1", "full-width-a-v2"), 
skipMergeVersions(fullWidthAKey))
+    assertEquals(Set("emoji-v1"), skipMergeVersions(emojiFaceKey))
+  }
+
+  @Test
+  def testLsmCompactionUsesUtf8Ordering(): Unit = {
+    val fullWidthAKey = "A-key"
+    val emojiFaceKey = "😀-key"
+    val tableName = "hoodie_mor_lsm_compaction"
+    val tablePath = s"${basePath}_mor_lsm_compaction"
+    val _spark = spark
+    import _spark.implicits._
+
+    val options = Map[String, String](
+      DataSourceWriteOptions.TABLE_TYPE.key -> 
HoodieTableType.MERGE_ON_READ.name,
+      DataSourceWriteOptions.OPERATION.key -> UPSERT_OPERATION_OPT_VAL,
+      DataSourceWriteOptions.RECORDKEY_FIELD.key -> "id",
+      DataSourceWriteOptions.PARTITIONPATH_FIELD.key -> "",
+      DataSourceWriteOptions.KEYGENERATOR_CLASS_NAME.key -> 
"org.apache.hudi.keygen.NonpartitionedKeyGenerator",
+      HoodieTableConfig.ORDERING_FIELDS.key -> "ts",
+      HoodieTableConfig.TABLE_STORAGE_LAYOUT.key -> 
HoodieTableConfig.TableStorageLayout.LSM_TREE.configValue,
+      HoodieWriteConfig.TBL_NAME.key -> tableName,
+      HoodieMetadataConfig.ENABLE.key -> "false",
+      HoodieCompactionConfig.INLINE_COMPACT.key -> "false",
+      HoodieCompactionConfig.INLINE_COMPACT_NUM_DELTA_COMMITS.key -> "1",
+      HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet",
+      "hoodie.insert.shuffle.parallelism" -> "1",
+      "hoodie.upsert.shuffle.parallelism" -> "1")
+    val (writeOpts, readOpts) = getWriterReaderOpts(HoodieRecordType.AVRO, 
options)
+
+    Seq(
+      (fullWidthAKey, "full-width-a-v1", 1L),
+      (emojiFaceKey, "emoji-v1", 1L))
+      .toDF("id", "value", "ts")
+      .repartition(1)
+      .write.format("hudi")
+      .options(writeOpts)
+      .mode(SaveMode.Append)
+      .save(tablePath)
+
+    // Update only the full-width A key so the base and log sorted runs have 
different key sets. This
+    // exposes an ordering mismatch during the LSM merge instead of merging 
two identical runs.
+    Seq((fullWidthAKey, "full-width-a-v2", 2L))
+      .toDF("id", "value", "ts")
+      .repartition(1)
+      .write.format("hudi")
+      .options(writeOpts)
+      .mode(SaveMode.Append)
+      .save(tablePath)
+
+    val metaClient = HoodieTableMetaClient.builder()
+      .setBasePath(tablePath)
+      .setConf(storageConf.newInstance())
+      .build()
+    assertTrue(metaClient.getTableConfig.isLSMTreeStorageLayout)
+
+    val client = DataSourceUtils.createHoodieClient(
+      spark.sparkContext, "", tablePath, tableName, writeOpts.asJava)
+      .asInstanceOf[SparkRDDWriteClient[HoodieRecordPayload[Nothing]]]
+    val compactionInstant = try {
+      val instant = client.scheduleCompaction(Option.empty()).get()
+      val statuses = client.compact(instant, true).getWriteStatuses.collect()
+      assertFalse(statuses.isEmpty)
+      assertTrue(statuses.asScala.forall(status => !status.hasErrors))
+      instant
+    } finally {
+      client.close()
+    }
+
+    
assertTrue(metaClient.reloadActiveTimeline().filterCompletedInstants.containsInstant(compactionInstant))
+
+    val latestBaseFiles = HoodieClientTestUtils.getLatestBaseFiles(
+      tablePath, metaClient.getStorage, s"$tablePath/*")
+    assertEquals(1, latestBaseFiles.size())
+    assertEquals(compactionInstant, latestBaseFiles.get(0).getCommitTime)
+
+    // Compacted LSM base files retain the table-level UTF-8 record-key 
ordering contract.
+    val physicalBaseFileKeys = 
spark.read.parquet(latestBaseFiles.get(0).getPath)
+      .select(HoodieRecord.RECORD_KEY_METADATA_FIELD)
+      .collect()
+      .map(_.getString(0))
+      .toList
+    assertEquals(Seq(fullWidthAKey, emojiFaceKey), physicalBaseFileKeys)
+
+    val actual = spark.read.format("hudi")
+      .options(readOpts)
+      .load(tablePath)
+      .select("id", "value")
+      .collect()
+      .map(row => row.getString(0) -> row.getString(1))
+      .toMap
+    assertEquals(Map(
+      fullWidthAKey -> "full-width-a-v2",
+      emojiFaceKey -> "emoji-v1"), actual)
+  }
+
   @ParameterizedTest
   @CsvSource(Array(
     // Inferred as COMMIT_TIME_ORDERING

Reply via email to