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