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 e6b8ef6c92f8 fix(core): use incoming record keys in native log writers
(#19874)
e6b8ef6c92f8 is described below
commit e6b8ef6c92f84b4412429199c3392847c4abfe58
Author: Danny Chan <[email protected]>
AuthorDate: Thu Sep 10 13:30:24 2026 +0800
fix(core): use incoming record keys in native log writers (#19874)
* fix(core): use incoming record keys in native log writers
---
.../hudi/io/HoodieNativeLogAppendHandle.java | 10 +---
.../hudi/io/cdc/HoodieNativeLogFormatWriter.java | 9 +--
.../hudi/io/TestHoodieNativeLogAppendHandle.java | 17 +++---
.../io/cdc/TestHoodieNativeLogFormatWriter.java | 66 ++++++++++++++++++++--
.../hudi/io/TestFlinkNativeLogAppendHandle.java | 4 +-
5 files changed, 77 insertions(+), 29 deletions(-)
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 f0b810e4943a..658a24d610e8 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
@@ -123,13 +123,10 @@ public class HoodieNativeLogAppendHandle<T, I, K, O>
extends HoodieAppendHandle<
}
HoodieRecord populatedRecord = hoodieRecord.prependMetaFields(
schema, writeSchemaWithMetaFields,
populateMetadataFields(hoodieRecord), recordProperties);
- String keyField = config.populateMetaFields()
- ? HoodieRecord.RECORD_KEY_METADATA_FIELD
- : hoodieTable.getMetaClient().getTableConfig().getRecordKeyFieldProp();
if (!canWriteDataFile()) {
flushAppend();
}
- writer.appendRecord(populatedRecord, writeSchemaWithMetaFields, keyField);
+ writer.appendRecord(populatedRecord, writeSchemaWithMetaFields);
if (isUpdateRecord || isLogCompaction) {
updatedRecordsWritten++;
} else {
@@ -143,13 +140,10 @@ public class HoodieNativeLogAppendHandle<T, I, K, O>
extends HoodieAppendHandle<
hoodieRecord.unseal();
hoodieRecord.clearNewLocation();
hoodieRecord.seal();
- String keyField =
schema.getField(HoodieRecord.RECORD_KEY_METADATA_FIELD).isPresent()
- ? HoodieRecord.RECORD_KEY_METADATA_FIELD
- : hoodieTable.getMetaClient().getTableConfig().getRecordKeyFieldProp();
if (!canWriteDeleteFile()) {
flushAppend();
}
- writer.appendDeleteRecord(hoodieRecord, schema, keyField);
+ writer.appendDeleteRecord(hoodieRecord, schema);
recordsDeleted++;
}
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/cdc/HoodieNativeLogFormatWriter.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/cdc/HoodieNativeLogFormatWriter.java
index c45dc2fa439b..97aa854dae65 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/cdc/HoodieNativeLogFormatWriter.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/cdc/HoodieNativeLogFormatWriter.java
@@ -56,6 +56,7 @@ import static
org.apache.hudi.common.model.LogExtensions.DELETE_LOG_EXTENSION;
/**
* Writes MOR log blocks as native files, for example {@code .log.parquet} and
{@code .deletes.parquet}.
+ * Uses the key already carried by incoming {@link HoodieRecord}s,
independently of whether metadata fields are populated.
*/
public class HoodieNativeLogFormatWriter extends HoodieLogFormat.Writer {
@@ -158,16 +159,16 @@ public class HoodieNativeLogFormatWriter extends
HoodieLogFormat.Writer {
return deleteFileWriter == null || deleteFileWriter.canWrite();
}
- public void appendRecord(HoodieRecord record, HoodieSchema recordSchema,
String keyFieldName) throws IOException {
+ public void appendRecord(HoodieRecord record, HoodieSchema recordSchema)
throws IOException {
ensureDataFileWriter(recordSchema);
- dataFileWriter.write(record.getRecordKey(recordSchema, keyFieldName),
+ dataFileWriter.write(record.getRecordKey(),
record, recordSchema, recordProperties);
dataRecordPositions.add(record.getCurrentPosition());
}
- public void appendDeleteRecord(HoodieRecord record, HoodieSchema
recordSchema, String keyFieldName) throws IOException {
+ public void appendDeleteRecord(HoodieRecord record, HoodieSchema
recordSchema) throws IOException {
ensureDeleteFileWriter();
- String recordKey = record.getRecordKey(recordSchema, keyFieldName);
+ String recordKey = record.getRecordKey();
Comparable orderingValue = getDeleteOrderingValue(record, recordSchema);
Object deleteEngineRecord = recordContext.constructEngineRecord(
deleteLogSchema, createDeleteLogFieldValues(recordKey, orderingValue));
diff --git
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieNativeLogAppendHandle.java
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieNativeLogAppendHandle.java
index 7787a2b542a5..e0ff62eb0ea4 100644
---
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieNativeLogAppendHandle.java
+++
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieNativeLogAppendHandle.java
@@ -91,12 +91,10 @@ public class TestHoodieNativeLogAppendHandle {
HoodieNativeLogFormatWriter writer = writers.constructed().get(0);
handle.writeData(inputRecord, true);
- verify(writer).appendRecord(eq(populatedRecord), any(HoodieSchema.class),
- eq(HoodieRecord.RECORD_KEY_METADATA_FIELD));
+ verify(writer).appendRecord(eq(populatedRecord),
any(HoodieSchema.class));
handle.writeDeleteRecord(inputRecord);
verify(inputRecord).clearNewLocation();
- verify(writer).appendDeleteRecord(eq(inputRecord),
any(HoodieSchema.class),
- eq(HoodieRecord.RECORD_KEY_METADATA_FIELD));
+ verify(writer).appendDeleteRecord(eq(inputRecord),
any(HoodieSchema.class));
handle.flushWriter();
verify(writer).flushAppend(any());
@@ -133,7 +131,7 @@ public class TestHoodieNativeLogAppendHandle {
handle.writeData(inputRecord, false);
InOrder rolloverOrder = inOrder(writer);
rolloverOrder.verify(writer).flushAppend(any());
- rolloverOrder.verify(writer).appendRecord(eq(populatedRecord),
any(HoodieSchema.class), any());
+ rolloverOrder.verify(writer).appendRecord(eq(populatedRecord),
any(HoodieSchema.class));
handle.flushWriter();
verify(writer, times(2)).flushAppend(any());
@@ -162,7 +160,7 @@ public class TestHoodieNativeLogAppendHandle {
}
@Test
- public void
testUsesConfiguredKeyWithoutMetadataFieldsAndSkipsIgnoredRecords() throws
Exception {
+ public void testWritesWithoutMetadataFieldsAndSkipsIgnoredRecords() throws
Exception {
HoodieWriteConfig config = HoodieWriteConfig.newBuilder()
.withPath("/tmp")
.withSchema(SCHEMA)
@@ -171,7 +169,6 @@ public class TestHoodieNativeLogAppendHandle {
.withWriteRecordPositionsEnabled(false)
.build();
HoodieTable table = table(config);
-
when(table.getMetaClient().getTableConfig().getRecordKeyFieldProp()).thenReturn("id");
try (MockedConstruction<HoodieNativeLogFormatWriter> writers =
mockConstruction(
HoodieNativeLogFormatWriter.class, (writer, context) -> {
@@ -186,16 +183,16 @@ public class TestHoodieNativeLogAppendHandle {
HoodieRecord ignoredRecord = mock(HoodieRecord.class);
when(ignoredRecord.shouldIgnore(any(HoodieSchema.class),
any())).thenReturn(true);
handle.writeData(ignoredRecord, false);
- verify(writer, never()).appendRecord(eq(ignoredRecord), any(), any());
+ verify(writer, never()).appendRecord(eq(ignoredRecord), any());
HoodieRecord inputRecord = mock(HoodieRecord.class);
HoodieRecord populatedRecord = mock(HoodieRecord.class);
when(inputRecord.prependMetaFields(any(HoodieSchema.class),
any(HoodieSchema.class), any(), any()))
.thenReturn(populatedRecord);
handle.writeData(inputRecord, false);
- verify(writer).appendRecord(eq(populatedRecord),
any(HoodieSchema.class), eq("id"));
+ verify(writer).appendRecord(eq(populatedRecord),
any(HoodieSchema.class));
handle.writeDeleteWithoutMetadata(inputRecord);
- verify(writer).appendDeleteRecord(eq(inputRecord),
any(HoodieSchema.class), eq("id"));
+ verify(writer).appendDeleteRecord(eq(inputRecord),
any(HoodieSchema.class));
}
}
diff --git
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/cdc/TestHoodieNativeLogFormatWriter.java
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/cdc/TestHoodieNativeLogFormatWriter.java
index 52f26b8e5ef8..0bd8c9b88127 100644
---
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/cdc/TestHoodieNativeLogFormatWriter.java
+++
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/cdc/TestHoodieNativeLogFormatWriter.java
@@ -18,13 +18,17 @@
package org.apache.hudi.io.cdc;
+import org.apache.hudi.common.avro.AvroRecordContext;
import org.apache.hudi.common.config.TypedProperties;
import org.apache.hudi.common.engine.RecordContext;
import org.apache.hudi.common.engine.TaskContextSupplier;
+import org.apache.hudi.common.model.HoodieAvroIndexedRecord;
import org.apache.hudi.common.model.HoodieFileFormat;
+import org.apache.hudi.common.model.HoodieKey;
import org.apache.hudi.common.model.HoodieRecord;
import org.apache.hudi.common.model.HoodieRecordMerger;
import org.apache.hudi.common.schema.HoodieSchema;
+import org.apache.hudi.common.schema.HoodieSchemaUtils;
import org.apache.hudi.common.table.HoodieTableVersion;
import org.apache.hudi.common.table.log.LogFileCreationCallback;
import org.apache.hudi.common.table.log.LogReaderUtils;
@@ -39,7 +43,12 @@ import org.apache.hudi.storage.HoodieStorage;
import org.apache.hudi.storage.StoragePath;
import org.apache.hudi.storage.StoragePathInfo;
+import org.apache.avro.generic.GenericData;
+import org.apache.avro.generic.GenericRecord;
+import org.apache.avro.generic.IndexedRecord;
import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
import org.mockito.ArgumentCaptor;
import org.mockito.MockedStatic;
@@ -62,6 +71,53 @@ import static org.mockito.Mockito.when;
public class TestHoodieNativeLogFormatWriter {
+ @ParameterizedTest
+ @ValueSource(booleans = {true, false})
+ public void testWritesCompositeRecordKey(boolean populateMetaFields) throws
Exception {
+ String schemaString = "{\"type\":\"record\",\"name\":\"test\",\"fields\":["
+ +
"{\"name\":\"id\",\"type\":\"string\"},{\"name\":\"name\",\"type\":\"string\"}]}";
+ HoodieWriteConfig config = HoodieWriteConfig.newBuilder()
+ .withPath("/tmp")
+ .withSchema(schemaString)
+ .withPopulateMetaFields(populateMetaFields)
+ .build();
+ HoodieSchema dataSchema = HoodieSchema.parse(schemaString);
+ HoodieSchema recordSchema = populateMetaFields ?
HoodieSchemaUtils.addMetadataFields(dataSchema) : dataSchema;
+ String recordKey = "id:1,name:alice";
+ GenericRecord data = new GenericData.Record(recordSchema.toAvroSchema());
+ data.put("id", "1");
+ data.put("name", "alice");
+ if (populateMetaFields) {
+ data.put(HoodieRecord.RECORD_KEY_METADATA_FIELD, recordKey);
+ }
+ HoodieRecord record = new HoodieAvroIndexedRecord(new HoodieKey(recordKey,
"partition"), data);
+ HoodieStorage storage = mock(HoodieStorage.class);
+ HoodieFileWriter dataFileWriter = mock(HoodieFileWriter.class);
+ HoodieFileWriter deleteFileWriter = mock(HoodieFileWriter.class);
+ TaskContextSupplier taskContextSupplier = mock(TaskContextSupplier.class);
+
+ try (MockedStatic<HoodieFileWriterFactory> writerFactory =
mockStatic(HoodieFileWriterFactory.class)) {
+ writerFactory.when(() -> HoodieFileWriterFactory.getFileWriter(
+ eq("100"), any(StoragePath.class), eq(storage), eq(config),
any(HoodieSchema.class),
+ eq(taskContextSupplier), eq(HoodieRecord.HoodieRecordType.AVRO)))
+ .thenReturn(dataFileWriter, deleteFileWriter);
+
+ try (HoodieNativeLogFormatWriter writer = new
HoodieNativeLogFormatWriter(
+ 4096, storage, new StoragePath("/tmp/partition"), "file-1", "100",
1, "1-0-1", 1024L,
+ new LogFileCreationCallback() {
+ }, HoodieTableVersion.current(), config, HoodieFileFormat.PARQUET,
recordSchema,
+ taskContextSupplier, new AvroRecordContext(), new ArrayList<>(),
Option.empty())) {
+ writer.appendRecord(record, recordSchema);
+ writer.appendDeleteRecord(record, recordSchema);
+ }
+ }
+
+ verify(dataFileWriter).write(eq(recordKey), eq(record), eq(recordSchema),
eq(config.getProps()));
+ ArgumentCaptor<IndexedRecord> deleteRecordCaptor =
ArgumentCaptor.forClass(IndexedRecord.class);
+ verify(deleteFileWriter).writeRow(eq(recordKey),
deleteRecordCaptor.capture());
+ assertEquals(recordKey, deleteRecordCaptor.getValue().get(0).toString());
+ }
+
@Test
public void testAddsRecordPositionsToDataLogFooter() throws Exception {
Map<HeaderMetadataType, String> parsedHeader =
writeDataLogFooterWithPositions(2L, 7L);
@@ -118,7 +174,7 @@ public class TestHoodieNativeLogFormatWriter {
public void testUsesDefaultOrderingValueForCommitTimeDeleteLog() throws
Exception {
HoodieSchema schema = mock(HoodieSchema.class);
HoodieRecord record = mock(HoodieRecord.class);
- when(record.getRecordKey(schema,
HoodieRecord.RECORD_KEY_METADATA_FIELD)).thenReturn("key-1");
+ when(record.getRecordKey()).thenReturn("key-1");
when(record.getCurrentPosition()).thenReturn(7L);
doReturn(HoodieRecord.DEFAULT_ORDERING_VALUE +
1).when(record).getOrderingValue(eq(schema), any(), any());
@@ -202,7 +258,7 @@ public class TestHoodieNativeLogFormatWriter {
Option.empty());
writer.appendRecord(recordWithPosition("key-1", 1L, schema),
- schema, HoodieRecord.RECORD_KEY_METADATA_FIELD);
+ schema);
writer.flushAppend(new HashMap<>());
return writer.getLastDataFileFormatMetadata();
}
@@ -253,7 +309,7 @@ public class TestHoodieNativeLogFormatWriter {
for (int i = 0; i < positions.length; i++) {
writer.appendRecord(recordWithPosition("key-" + i, positions[i],
schema),
- schema, HoodieRecord.RECORD_KEY_METADATA_FIELD);
+ schema);
}
Map<HeaderMetadataType, String> header = new HashMap<>();
@@ -323,7 +379,7 @@ public class TestHoodieNativeLogFormatWriter {
orderingFieldNames,
baseFileInstantTimeOfPositions);
- writer.appendDeleteRecord(record, schema,
HoodieRecord.RECORD_KEY_METADATA_FIELD);
+ writer.appendDeleteRecord(record, schema);
Map<HeaderMetadataType, String> header = new HashMap<>();
header.put(HeaderMetadataType.SCHEMA, schemaString);
@@ -339,7 +395,7 @@ public class TestHoodieNativeLogFormatWriter {
private static HoodieRecord recordWithPosition(String key, long position,
HoodieSchema schema) throws Exception {
HoodieRecord record = mock(HoodieRecord.class);
- when(record.getRecordKey(schema,
HoodieRecord.RECORD_KEY_METADATA_FIELD)).thenReturn(key);
+ when(record.getRecordKey()).thenReturn(key);
when(record.getCurrentPosition()).thenReturn(position);
when(record.getOrderingValue(eq(schema), any(),
any())).thenReturn(OrderingValues.getDefault());
return record;
diff --git
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/TestFlinkNativeLogAppendHandle.java
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/TestFlinkNativeLogAppendHandle.java
index 60cc172eafcd..876149987162 100644
---
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/TestFlinkNativeLogAppendHandle.java
+++
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/TestFlinkNativeLogAppendHandle.java
@@ -84,8 +84,8 @@ public class TestFlinkNativeLogAppendHandle {
verify(writer, never()).canWriteDataFile();
verify(writer, never()).canWriteDeleteFile();
verify(writer, never()).flushAppend(any());
- verify(writer, times(2)).appendRecord(eq(populatedRecord),
any(HoodieSchema.class), any());
- verify(writer, times(2)).appendDeleteRecord(eq(inputRecord),
any(HoodieSchema.class), any());
+ verify(writer, times(2)).appendRecord(eq(populatedRecord),
any(HoodieSchema.class));
+ verify(writer, times(2)).appendDeleteRecord(eq(inputRecord),
any(HoodieSchema.class));
handle.flushWriter();
verify(writer).flushAppend(any());