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