From c3c38cf691e15ecd7632ced09ac77800f1791eb1 Mon Sep 17 00:00:00 2001 From: Lin Liu Date: Mon, 14 Jul 2025 12:21:36 -0700 Subject: [PATCH 01/10] Introduce input based file group record buffer. --- .../table/read/HoodieFileGroupReader.java | 64 ++++++++++++- .../read/InputBasedFileGroupRecordBuffer.java | 94 +++++++++++++++++++ 2 files changed, 154 insertions(+), 4 deletions(-) create mode 100644 hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java index 73094129b14fc..c74025828dee1 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java @@ -55,6 +55,7 @@ import java.io.IOException; import java.util.ArrayList; import java.util.Collections; +import java.util.Iterator; import java.util.List; import java.util.function.UnaryOperator; import java.util.stream.Collectors; @@ -112,6 +113,54 @@ public HoodieFileGroupReader(HoodieReaderContext readerContext, HoodieStorage InputSplit.fromFileSlice(fileSlice, start, length), Option.empty(), false); } + private HoodieFileGroupReader(HoodieReaderContext readerContext, HoodieStorage storage, String tablePath, + String latestCommitTime, Schema dataSchema, Schema requestedSchema, + Option internalSchemaOpt, Iterator inputRecords, + HoodieTableMetaClient hoodieTableMetaClient, TypedProperties props, + boolean shouldUseRecordPosition, boolean allowInflightInstants, + boolean emitDelete, boolean sortOutput, + InputSplit inputSplit, Option updateCallback, boolean enableOptimizedLogBlockScan) { + this.readerContext = readerContext; + this.fileGroupUpdateCallback = updateCallback; + this.metaClient = hoodieTableMetaClient; + this.storage = storage; + this.enableOptimizedLogBlockScan = enableOptimizedLogBlockScan; + this.inputSplit = inputSplit; + readerContext.setHasLogFiles(!this.inputSplit.logFiles.isEmpty()); + readerContext.setPartitionPath(inputSplit.partitionPath); + if (readerContext.getHasLogFiles() && inputSplit.start != 0) { + throw new IllegalArgumentException("Filegroup reader is doing log file merge but not reading from the start of the base file"); + } + this.props = props; + HoodieTableConfig tableConfig = hoodieTableMetaClient.getTableConfig(); + this.partitionPathFields = tableConfig.getPartitionFields(); + readerContext.initRecordMerger(props); + readerContext.setTablePath(tablePath); + readerContext.setLatestCommitTime(latestCommitTime); + boolean isSkipMerge = ConfigUtils.getStringWithAltKeys(props, HoodieReaderConfig.MERGE_TYPE, true).equalsIgnoreCase(HoodieReaderConfig.REALTIME_SKIP_MERGE); + readerContext.setShouldMergeUseRecordPosition(shouldUseRecordPosition && !isSkipMerge && readerContext.getHasLogFiles()); + readerContext.setHasBootstrapBaseFile(inputSplit.baseFileOption.flatMap(HoodieBaseFile::getBootstrapBaseFile).isPresent()); + readerContext.setSchemaHandler(readerContext.supportsParquetRowIndex() + ? new ParquetRowIndexBasedSchemaHandler<>(readerContext, dataSchema, requestedSchema, internalSchemaOpt, tableConfig, props) + : new FileGroupReaderSchemaHandler<>(readerContext, dataSchema, requestedSchema, internalSchemaOpt, tableConfig, props)); + this.outputConverter = readerContext.getSchemaHandler().getOutputConverter(); + this.orderingFieldName = readerContext.getMergeMode() == RecordMergeMode.COMMIT_TIME_ORDERING + ? Option.empty() + : Option.ofNullable(ConfigUtils.getOrderingField(props)) + .or(() -> { + String preCombineField = hoodieTableMetaClient.getTableConfig().getPreCombineField(); + if (StringUtils.isNullOrEmpty(preCombineField)) { + return Option.empty(); + } + return Option.of(preCombineField); + }); + this.readStats = new HoodieReadStats(); + this.recordBuffer = getRecordBuffer(readerContext, hoodieTableMetaClient, + readerContext.getMergeMode(), tableConfig.getPartialUpdateMode(), props, + isSkipMerge, shouldUseRecordPosition, readStats, emitDelete, sortOutput, Option.of(inputRecords)); + this.allowInflightInstants = allowInflightInstants; + } + private HoodieFileGroupReader(HoodieReaderContext readerContext, HoodieStorage storage, String tablePath, String latestCommitTime, Schema dataSchema, Schema requestedSchema, Option internalSchemaOpt, HoodieTableMetaClient hoodieTableMetaClient, TypedProperties props, @@ -154,7 +203,7 @@ private HoodieFileGroupReader(HoodieReaderContext readerContext, HoodieStorag this.readStats = new HoodieReadStats(); this.recordBuffer = getRecordBuffer(readerContext, hoodieTableMetaClient, readerContext.getMergeMode(), tableConfig.getPartialUpdateMode(), props, - isSkipMerge, shouldUseRecordPosition, readStats, emitDelete, sortOutput); + isSkipMerge, shouldUseRecordPosition, readStats, emitDelete, sortOutput, Option.empty()); this.allowInflightInstants = allowInflightInstants; } @@ -170,11 +219,17 @@ private FileGroupRecordBuffer getRecordBuffer(HoodieReaderContext readerCo boolean shouldUseRecordPosition, HoodieReadStats readStats, boolean emitDelete, - boolean sortOutput) { + boolean sortOutput, + Option> inputRecordOpt) { + UpdateProcessor updateProcessor = UpdateProcessor.create(readStats, readerContext, emitDelete, fileGroupUpdateCallback); if (inputSplit.logFiles.isEmpty()) { + if (inputRecordOpt.isPresent()) { + return new InputBasedFileGroupRecordBuffer<>( + readerContext, hoodieTableMetaClient, recordMergeMode, partialUpdateMode, + props, orderingFieldName, inputRecordOpt.get(), updateProcessor); + } return null; } - UpdateProcessor updateProcessor = UpdateProcessor.create(readStats, readerContext, emitDelete, fileGroupUpdateCallback); if (isSkipMerge) { return new UnmergedFileGroupRecordBuffer<>( readerContext, hoodieTableMetaClient, recordMergeMode, partialUpdateMode, props, readStats); @@ -183,7 +238,8 @@ private FileGroupRecordBuffer getRecordBuffer(HoodieReaderContext readerCo readerContext, hoodieTableMetaClient, recordMergeMode, partialUpdateMode, props, orderingFieldName, updateProcessor); } else if (shouldUseRecordPosition && inputSplit.baseFileOption.isPresent()) { return new PositionBasedFileGroupRecordBuffer<>( - readerContext, hoodieTableMetaClient, recordMergeMode, partialUpdateMode, inputSplit.baseFileOption.get().getCommitTime(), props, orderingFieldName, updateProcessor); + readerContext, hoodieTableMetaClient, recordMergeMode, partialUpdateMode, inputSplit.baseFileOption.get().getCommitTime(), + props, orderingFieldName, updateProcessor); } else { return new KeyBasedFileGroupRecordBuffer<>( readerContext, hoodieTableMetaClient, recordMergeMode, partialUpdateMode, props, orderingFieldName, updateProcessor); diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java new file mode 100644 index 0000000000000..664a5e9c780e7 --- /dev/null +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java @@ -0,0 +1,94 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.hudi.common.table.read; + +import org.apache.hudi.common.config.RecordMergeMode; +import org.apache.hudi.common.config.TypedProperties; +import org.apache.hudi.common.engine.HoodieReaderContext; +import org.apache.hudi.common.model.DeleteRecord; +import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.common.table.PartialUpdateMode; +import org.apache.hudi.common.table.log.KeySpec; +import org.apache.hudi.common.table.log.block.HoodieDataBlock; +import org.apache.hudi.common.table.log.block.HoodieDeleteBlock; +import org.apache.hudi.common.util.Option; +import org.apache.hudi.common.util.collection.ClosableIterator; +import org.apache.hudi.exception.HoodieNotSupportedException; + +import java.io.IOException; +import java.io.Serializable; +import java.util.Iterator; + +public class InputBasedFileGroupRecordBuffer extends KeyBasedFileGroupRecordBuffer { + private final Iterator inputRecordIterator; + + public InputBasedFileGroupRecordBuffer(HoodieReaderContext readerContext, + HoodieTableMetaClient hoodieTableMetaClient, + RecordMergeMode recordMergeMode, + PartialUpdateMode partialUpdateMode, + TypedProperties props, + Option orderingFieldName, + Iterator inputRecordIterator, + UpdateProcessor updateProcessor) { + super(readerContext, hoodieTableMetaClient, recordMergeMode, partialUpdateMode, props, orderingFieldName, updateProcessor); + this.inputRecordIterator = inputRecordIterator; + + // Read input records into the buffer. + populateRecordBuffer(); + } + + private void populateRecordBuffer() { + if (null == inputRecordIterator) { + throw new IllegalArgumentException("InputRecordIterator can not be null"); + } + + while (inputRecordIterator.hasNext()) { + T engineRecord = inputRecordIterator.next(); + String recordKey = readerContext.getRecordKey(nextRecord, readerSchema); + boolean isDelete = + isBuiltInDeleteRecord(engineRecord) + || isCustomDeleteRecord(engineRecord) + || isDeleteHoodieOperation(engineRecord); + BufferedRecord bufferedRecord = BufferedRecord.forRecordWithContext( + engineRecord, readerSchema, readerContext, orderingFieldName, isDelete); + records.put(recordKey, bufferedRecord); + } + } + + @Override + public void processDataBlock(HoodieDataBlock dataBlock, Option keySpecOpt) throws IOException { + throw new HoodieNotSupportedException("Method 'processDataBlock' is not supported"); + } + + @Override + public void processNextDataRecord(BufferedRecord record, Serializable recordKey) throws IOException { + throw new HoodieNotSupportedException("Method 'processNextDataRecord' is not supported"); + } + + @Override + public void processDeleteBlock(HoodieDeleteBlock deleteBlock) throws IOException { + throw new HoodieNotSupportedException("Method 'processDeleteBlock' is not supported"); + } + + @Override + public void processNextDeletedRecord(DeleteRecord deleteRecord, Serializable recordKey) { + throw new HoodieNotSupportedException("Method 'processNextDeletedRecord' is not supported"); + } +} From 2861d0bc0e175c3767546fc45b2422d79039d1bd Mon Sep 17 00:00:00 2001 From: Lin Liu Date: Sat, 19 Jul 2025 07:50:13 -0700 Subject: [PATCH 02/10] Support record iterator based file group reader merge handle --- .../io/FileGroupReaderBasedMergeHandle.java | 54 ++++++++++++++++-- .../table/read/HoodieFileGroupReader.java | 55 +++---------------- .../read/InputBasedFileGroupRecordBuffer.java | 4 ++ 3 files changed, 61 insertions(+), 52 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 a8cc7e6864e6d..85da4d50c4952 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 @@ -27,6 +27,7 @@ import org.apache.hudi.common.fs.FSUtils; import org.apache.hudi.common.model.CompactionOperation; import org.apache.hudi.common.model.HoodieLogFile; +import org.apache.hudi.common.model.HoodieBaseFile; import org.apache.hudi.common.model.HoodiePartitionMetadata; import org.apache.hudi.common.model.HoodieRecord; import org.apache.hudi.common.model.HoodieWriteStat; @@ -36,11 +37,13 @@ import org.apache.hudi.common.table.read.HoodieReadStats; import org.apache.hudi.common.util.Option; import org.apache.hudi.common.util.collection.ClosableIterator; +import org.apache.hudi.common.util.collection.MappingIterator; import org.apache.hudi.config.HoodieWriteConfig; import org.apache.hudi.exception.HoodieUpsertException; import org.apache.hudi.internal.schema.InternalSchema; import org.apache.hudi.internal.schema.utils.SerDeHelper; import org.apache.hudi.io.storage.HoodieFileWriterFactory; +import org.apache.hudi.keygen.BaseKeyGenerator; import org.apache.hudi.storage.StoragePath; import org.apache.hudi.table.HoodieTable; import org.apache.hudi.table.action.compact.strategy.CompactionStrategy; @@ -56,6 +59,7 @@ import java.io.IOException; import java.util.Collections; import java.util.HashSet; +import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.function.UnaryOperator; @@ -81,6 +85,7 @@ public class FileGroupReaderBasedMergeHandle extends HoodieWriteMerg private HoodieReadStats readStats; private final HoodieRecord.HoodieRecordType recordType; private final Option cdcLogger; + private final Iterator> recordIterator; public FileGroupReaderBasedMergeHandle(HoodieWriteConfig config, String instantTime, HoodieTable hoodieTable, CompactionOperation operation, TaskContextSupplier taskContextSupplier, @@ -107,6 +112,41 @@ public FileGroupReaderBasedMergeHandle(HoodieWriteConfig config, String instantT this.cdcLogger = Option.empty(); } init(operation, this.partitionPath); + this.recordIterator = null; + } + + /** + * FG reader based generic merge handle, which is not just for compaction. + */ + public FileGroupReaderBasedMergeHandle(HoodieWriteConfig config, String instantTime, HoodieTable hoodieTable, + Iterator> recordItr, String partitionPath, String fileId, + TaskContextSupplier taskContextSupplier, HoodieBaseFile baseFile, + Option keyGeneratorOpt, HoodieReaderContext readerContext, + String maxInstantTime, HoodieRecord.HoodieRecordType enginRecordType) { + super(config, instantTime, hoodieTable, recordItr, partitionPath, fileId, taskContextSupplier, baseFile, keyGeneratorOpt); + this.maxInstantTime = maxInstantTime; + this.keyToNewRecords = Collections.emptyMap(); + this.readerContext = readerContext; + this.recordIterator = recordItr; + this.operation = null; + if (hoodieTable.getMetaClient().getTableConfig().isCDCEnabled()) { + this.cdcLogger = Option.of(new HoodieCDCLogger( + 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(); + } + // 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; + init(operation, this.partitionPath); } private void init(CompactionOperation operation, String partitionPath) { @@ -177,12 +217,18 @@ public void doMerge() { Stream logFiles = operation.getDeltaFileNames().stream().map(logFileName -> new HoodieLogFile(new StoragePath(FSUtils.constructAbsolutePath( config.getBasePath(), operation.getPartitionPath()), logFileName))); + Iterator engineRecordIterator = recordIterator == null + ? null : new MappingIterator<>(recordIterator, HoodieRecord::getData); // Initializes file group reader - try (HoodieFileGroupReader fileGroupReader = HoodieFileGroupReader.newBuilder().withReaderContext(readerContext).withHoodieTableMetaClient(hoodieTable.getMetaClient()) - .withLatestCommitTime(maxInstantTime).withPartitionPath(partitionPath).withBaseFileOption(Option.ofNullable(baseFileToMerge)).withLogFiles(logFiles) - .withDataSchema(writeSchemaWithMetaFields).withRequestedSchema(writeSchemaWithMetaFields).withInternalSchema(internalSchemaOption).withProps(props) + try (HoodieFileGroupReader fileGroupReader = HoodieFileGroupReader.newBuilder() + .withReaderContext(readerContext).withHoodieTableMetaClient(hoodieTable.getMetaClient()) + .withLatestCommitTime(maxInstantTime).withPartitionPath(partitionPath) + .withBaseFileOption(Option.ofNullable(baseFileToMerge)).withLogFiles(logFiles) + .withDataSchema(writeSchemaWithMetaFields).withRequestedSchema(writeSchemaWithMetaFields) + .withInternalSchema(internalSchemaOption).withProps(props) .withShouldUseRecordPosition(usePosition).withSortOutput(hoodieTable.requireSortedRecords()) - .withFileGroupUpdateCallback(cdcLogger.map(logger -> new CDCCallback(logger, readerContext))).build()) { + .withFileGroupUpdateCallback(cdcLogger.map(logger -> new CDCCallback(logger, readerContext))) + .withRecordIterator(engineRecordIterator).build()) { // Reads the records from the file slice try (ClosableIterator> recordIterator = fileGroupReader.getClosableHoodieRecordIterator()) { while (recordIterator.hasNext()) { diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java index c74025828dee1..40b58d26664b7 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java @@ -113,54 +113,6 @@ public HoodieFileGroupReader(HoodieReaderContext readerContext, HoodieStorage InputSplit.fromFileSlice(fileSlice, start, length), Option.empty(), false); } - private HoodieFileGroupReader(HoodieReaderContext readerContext, HoodieStorage storage, String tablePath, - String latestCommitTime, Schema dataSchema, Schema requestedSchema, - Option internalSchemaOpt, Iterator inputRecords, - HoodieTableMetaClient hoodieTableMetaClient, TypedProperties props, - boolean shouldUseRecordPosition, boolean allowInflightInstants, - boolean emitDelete, boolean sortOutput, - InputSplit inputSplit, Option updateCallback, boolean enableOptimizedLogBlockScan) { - this.readerContext = readerContext; - this.fileGroupUpdateCallback = updateCallback; - this.metaClient = hoodieTableMetaClient; - this.storage = storage; - this.enableOptimizedLogBlockScan = enableOptimizedLogBlockScan; - this.inputSplit = inputSplit; - readerContext.setHasLogFiles(!this.inputSplit.logFiles.isEmpty()); - readerContext.setPartitionPath(inputSplit.partitionPath); - if (readerContext.getHasLogFiles() && inputSplit.start != 0) { - throw new IllegalArgumentException("Filegroup reader is doing log file merge but not reading from the start of the base file"); - } - this.props = props; - HoodieTableConfig tableConfig = hoodieTableMetaClient.getTableConfig(); - this.partitionPathFields = tableConfig.getPartitionFields(); - readerContext.initRecordMerger(props); - readerContext.setTablePath(tablePath); - readerContext.setLatestCommitTime(latestCommitTime); - boolean isSkipMerge = ConfigUtils.getStringWithAltKeys(props, HoodieReaderConfig.MERGE_TYPE, true).equalsIgnoreCase(HoodieReaderConfig.REALTIME_SKIP_MERGE); - readerContext.setShouldMergeUseRecordPosition(shouldUseRecordPosition && !isSkipMerge && readerContext.getHasLogFiles()); - readerContext.setHasBootstrapBaseFile(inputSplit.baseFileOption.flatMap(HoodieBaseFile::getBootstrapBaseFile).isPresent()); - readerContext.setSchemaHandler(readerContext.supportsParquetRowIndex() - ? new ParquetRowIndexBasedSchemaHandler<>(readerContext, dataSchema, requestedSchema, internalSchemaOpt, tableConfig, props) - : new FileGroupReaderSchemaHandler<>(readerContext, dataSchema, requestedSchema, internalSchemaOpt, tableConfig, props)); - this.outputConverter = readerContext.getSchemaHandler().getOutputConverter(); - this.orderingFieldName = readerContext.getMergeMode() == RecordMergeMode.COMMIT_TIME_ORDERING - ? Option.empty() - : Option.ofNullable(ConfigUtils.getOrderingField(props)) - .or(() -> { - String preCombineField = hoodieTableMetaClient.getTableConfig().getPreCombineField(); - if (StringUtils.isNullOrEmpty(preCombineField)) { - return Option.empty(); - } - return Option.of(preCombineField); - }); - this.readStats = new HoodieReadStats(); - this.recordBuffer = getRecordBuffer(readerContext, hoodieTableMetaClient, - readerContext.getMergeMode(), tableConfig.getPartialUpdateMode(), props, - isSkipMerge, shouldUseRecordPosition, readStats, emitDelete, sortOutput, Option.of(inputRecords)); - this.allowInflightInstants = allowInflightInstants; - } - private HoodieFileGroupReader(HoodieReaderContext readerContext, HoodieStorage storage, String tablePath, String latestCommitTime, Schema dataSchema, Schema requestedSchema, Option internalSchemaOpt, HoodieTableMetaClient hoodieTableMetaClient, TypedProperties props, @@ -511,6 +463,7 @@ public static class Builder { private boolean sortOutput = false; private boolean enableOptimizedLogBlockScan = false; private Option fileGroupUpdateCallback = Option.empty(); + private Option> recordIteratorOpt = Option.empty(); public Builder withReaderContext(HoodieReaderContext readerContext) { this.readerContext = readerContext; @@ -616,6 +569,12 @@ public Builder withSortOutput(boolean sortOutput) { return this; } + public Builder withRecordIterator(Iterator iterator) { + this.recordIteratorOpt = Option.ofNullable(iterator); + return this; + } + + public HoodieFileGroupReader build() { ValidationUtils.checkArgument(readerContext != null, "Reader context is required"); ValidationUtils.checkArgument(hoodieTableMetaClient != null, "Hoodie table meta client is required"); diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java index 664a5e9c780e7..61307bf7ee326 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java @@ -47,7 +47,11 @@ public InputBasedFileGroupRecordBuffer(HoodieReaderContext readerContext, Option orderingFieldName, Iterator inputRecordIterator, UpdateProcessor updateProcessor) { +<<<<<<< HEAD super(readerContext, hoodieTableMetaClient, recordMergeMode, partialUpdateMode, props, orderingFieldName, updateProcessor); +======= + super(readerContext, hoodieTableMetaClient, recordMergeMode,partialUpdateMode, props, readStats, orderingFieldName, updateProcessor); +>>>>>>> 6d74b6820d3 (Support record iterator based file group reader merge handle) this.inputRecordIterator = inputRecordIterator; // Read input records into the buffer. From 0657b9fbe6096c9b70bd9a1df2b85146f03dad74 Mon Sep 17 00:00:00 2001 From: Lin Liu Date: Sat, 19 Jul 2025 08:54:58 -0700 Subject: [PATCH 03/10] Plug into spark --- .../apache/hudi/io/FileGroupReaderBasedMergeHandle.java | 8 ++++---- .../hudi/common/table/read/HoodieFileGroupReader.java | 1 - .../table/read/InputBasedFileGroupRecordBuffer.java | 1 - 3 files changed, 4 insertions(+), 6 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 85da4d50c4952..cb364e5644b42 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 @@ -120,10 +120,10 @@ public FileGroupReaderBasedMergeHandle(HoodieWriteConfig config, String instantT */ public FileGroupReaderBasedMergeHandle(HoodieWriteConfig config, String instantTime, HoodieTable hoodieTable, Iterator> recordItr, String partitionPath, String fileId, - TaskContextSupplier taskContextSupplier, HoodieBaseFile baseFile, - Option keyGeneratorOpt, HoodieReaderContext readerContext, - String maxInstantTime, HoodieRecord.HoodieRecordType enginRecordType) { - super(config, instantTime, hoodieTable, recordItr, partitionPath, fileId, taskContextSupplier, baseFile, keyGeneratorOpt); + TaskContextSupplier taskContextSupplier, Option keyGeneratorOpt, + HoodieReaderContext readerContext, String maxInstantTime, + HoodieRecord.HoodieRecordType enginRecordType) { + super(config, instantTime, hoodieTable, recordItr, partitionPath, fileId, taskContextSupplier, keyGeneratorOpt); this.maxInstantTime = maxInstantTime; this.keyToNewRecords = Collections.emptyMap(); this.readerContext = readerContext; diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java index 40b58d26664b7..6eaa0a49512df 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java @@ -574,7 +574,6 @@ public Builder withRecordIterator(Iterator iterator) { return this; } - public HoodieFileGroupReader build() { ValidationUtils.checkArgument(readerContext != null, "Reader context is required"); ValidationUtils.checkArgument(hoodieTableMetaClient != null, "Hoodie table meta client is required"); diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java index 61307bf7ee326..01589d85e3773 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java @@ -29,7 +29,6 @@ import org.apache.hudi.common.table.log.block.HoodieDataBlock; import org.apache.hudi.common.table.log.block.HoodieDeleteBlock; import org.apache.hudi.common.util.Option; -import org.apache.hudi.common.util.collection.ClosableIterator; import org.apache.hudi.exception.HoodieNotSupportedException; import java.io.IOException; From 172ef4765d4083c3cfdcc1bdf3ef4bf0597d4f2a Mon Sep 17 00:00:00 2001 From: Lin Liu Date: Sun, 20 Jul 2025 08:17:34 -0700 Subject: [PATCH 04/10] Add support for secondary index --- .../io/FileGroupReaderBasedMergeHandle.java | 115 +++++++++++++++++- .../hudi/io/HoodieWriteMergeHandle.java | 2 +- .../table/read/HoodieFileGroupReader.java | 14 +-- .../read/InputBasedFileGroupRecordBuffer.java | 4 - .../common/table/read/UpdateProcessor.java | 44 +++++-- .../metadata/HoodieTableMetadataUtil.java | 2 +- .../TestKeyBasedFileGroupRecordBuffer.java | 2 +- ...stSortedKeyBasedFileGroupRecordBuffer.java | 2 +- .../hudi/cdc/CDCFileGroupIterator.scala | 2 +- 9 files changed, 159 insertions(+), 28 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 cb364e5644b42..824541db900c6 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 @@ -27,11 +27,13 @@ import org.apache.hudi.common.fs.FSUtils; import org.apache.hudi.common.model.CompactionOperation; import org.apache.hudi.common.model.HoodieLogFile; -import org.apache.hudi.common.model.HoodieBaseFile; +import org.apache.hudi.common.model.HoodieIndexDefinition; +import org.apache.hudi.common.model.HoodieKey; import org.apache.hudi.common.model.HoodiePartitionMetadata; import org.apache.hudi.common.model.HoodieRecord; import org.apache.hudi.common.model.HoodieWriteStat; import org.apache.hudi.common.table.cdc.HoodieCDCUtils; +import org.apache.hudi.common.table.read.BufferedRecord; import org.apache.hudi.common.table.read.BaseFileUpdateCallback; import org.apache.hudi.common.table.read.HoodieFileGroupReader; import org.apache.hudi.common.table.read.HoodieReadStats; @@ -57,11 +59,13 @@ import javax.annotation.concurrent.NotThreadSafe; import java.io.IOException; +import java.util.ArrayList; import java.util.Collections; import java.util.HashSet; import java.util.Iterator; import java.util.List; import java.util.Map; +import java.util.function.Supplier; import java.util.function.UnaryOperator; import java.util.stream.Stream; @@ -227,7 +231,7 @@ public void doMerge() { .withDataSchema(writeSchemaWithMetaFields).withRequestedSchema(writeSchemaWithMetaFields) .withInternalSchema(internalSchemaOption).withProps(props) .withShouldUseRecordPosition(usePosition).withSortOutput(hoodieTable.requireSortedRecords()) - .withFileGroupUpdateCallback(cdcLogger.map(logger -> new CDCCallback(logger, readerContext))) + .withFileGroupUpdateCallback(createCallbacks()) .withRecordIterator(engineRecordIterator).build()) { // Reads the records from the file slice try (ClosableIterator> recordIterator = fileGroupReader.getClosableHoodieRecordIterator()) { @@ -266,6 +270,28 @@ public void doMerge() { } } + private List> createCallbacks() { + List> callbacks = new ArrayList<>(); + // Handle CDC workflow. + if (cdcLogger.isPresent()) { + callbacks.add(new CDCCallback<>(cdcLogger.get(), readerContext)); + } + // Stream secondary index stats. + if (isSecondaryIndexStatsStreamingWritesEnabled || writeStatus.isTrackingSuccessfulWrites()) { + callbacks.add(new SecondaryIndexCallback<>( + partitionPath, + writeSchemaWithMetaFields, + readerContext, + this::getNewSchema, + writeStatus, + secondaryIndexDefns, + keyGeneratorOpt, + config + )); + } + return callbacks; + } + @Override protected void writeIncomingRecords() { // no operation. @@ -335,4 +361,89 @@ private GenericRecord convertOutput(T record) { return convertedRecord == null ? null : readerContext.convertToAvroRecord(convertedRecord, requestedSchema.get()); } } + + private static class SecondaryIndexCallback implements BaseFileUpdateCallback { + private final String partitionPath; + private final Schema writeSchemaWithMetaFields; + private final HoodieReaderContext readerContext; + private final Supplier newSchemaSupplier; + private final WriteStatus writeStatus; + private final List secondaryIndexDefns; + private final Option keyGeneratorOpt; + private final HoodieWriteConfig config; + + public SecondaryIndexCallback(String partitionPath, + Schema writeSchemaWithMetaFields, + HoodieReaderContext readerContext, + Supplier newSchemaSupplier, + WriteStatus writeStatus, + List secondaryIndexDefns, + Option keyGeneratorOpt, + HoodieWriteConfig config) { + this.partitionPath = partitionPath; + this.writeSchemaWithMetaFields = writeSchemaWithMetaFields; + this.readerContext = readerContext; + this.newSchemaSupplier = newSchemaSupplier; + this.secondaryIndexDefns = secondaryIndexDefns; + this.keyGeneratorOpt = keyGeneratorOpt; + this.writeStatus = writeStatus; + this.config = config; + } + + @Override + public void onUpdate(String recordKey, T previousRecord, T mergedRecord) { + HoodieKey hoodieKey = new HoodieKey(recordKey, partitionPath); + BufferedRecord bufferedPrevousRecord = BufferedRecord.forRecordWithContext( + previousRecord, writeSchemaWithMetaFields, readerContext, Option.empty(), false); + BufferedRecord bufferedMergedRecord = BufferedRecord.forRecordWithContext( + mergedRecord, writeSchemaWithMetaFields, readerContext, Option.empty(), false); + SecondaryIndexStreamingTracker.trackSecondaryIndexStats( + hoodieKey, + Option.of(readerContext.constructHoodieRecord(bufferedMergedRecord)), + readerContext.constructHoodieRecord(bufferedPrevousRecord), + false, + writeStatus, + writeSchemaWithMetaFields, + newSchemaSupplier, + secondaryIndexDefns, + keyGeneratorOpt, + config); + } + + @Override + public void onInsert(String recordKey, T newRecord) { + HoodieKey hoodieKey = new HoodieKey(recordKey, partitionPath); + BufferedRecord bufferedNewRecord = BufferedRecord.forRecordWithContext( + newRecord, writeSchemaWithMetaFields, readerContext, Option.empty(), false); + SecondaryIndexStreamingTracker.trackSecondaryIndexStats( + hoodieKey, + Option.of(readerContext.constructHoodieRecord(bufferedNewRecord)), + null, + false, + writeStatus, + writeSchemaWithMetaFields, + newSchemaSupplier, + secondaryIndexDefns, + keyGeneratorOpt, + config); + } + + @Override + public void onDelete(String recordKey, T previousRecord) { + HoodieKey hoodieKey = new HoodieKey(recordKey, partitionPath); + BufferedRecord bufferedPrevousRecord = BufferedRecord.forRecordWithContext( + previousRecord, writeSchemaWithMetaFields, readerContext, Option.empty(), false); + SecondaryIndexStreamingTracker.trackSecondaryIndexStats( + hoodieKey, + null, + readerContext.constructHoodieRecord(bufferedPrevousRecord), + true, + writeStatus, + writeSchemaWithMetaFields, + newSchemaSupplier, + secondaryIndexDefns, + keyGeneratorOpt, + config); + } + } } diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieWriteMergeHandle.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieWriteMergeHandle.java index 8d389d7ba88e6..96c8afbc11752 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieWriteMergeHandle.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieWriteMergeHandle.java @@ -429,7 +429,7 @@ protected void writeIncomingRecords() throws IOException { } } - private Schema getNewSchema() { + protected Schema getNewSchema() { return preserveMetadata ? writeSchemaWithMetaFields : writeSchema; } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java index 6eaa0a49512df..b2ebaa58c3b22 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java @@ -92,7 +92,7 @@ public final class HoodieFileGroupReader implements Closeable { // considers the log records which are inflight. private final boolean allowInflightInstants; // Callback to run custom logic on updates to the base files for the file group - private final Option fileGroupUpdateCallback; + private final List> fileGroupUpdateCallbacks; private final boolean enableOptimizedLogBlockScan; // The list of instant times read from the log blocks, this value is used by the log-compaction to allow optimized log-block scans private List validBlockInstants = Collections.emptyList(); @@ -110,16 +110,16 @@ public HoodieFileGroupReader(HoodieReaderContext readerContext, HoodieStorage long start, long length, boolean shouldUseRecordPosition) { this(readerContext, storage, tablePath, latestCommitTime, dataSchema, requestedSchema, internalSchemaOpt, hoodieTableMetaClient, props, shouldUseRecordPosition, false, false, false, - InputSplit.fromFileSlice(fileSlice, start, length), Option.empty(), false); + InputSplit.fromFileSlice(fileSlice, start, length), Collections.emptyList(), false); } private HoodieFileGroupReader(HoodieReaderContext readerContext, HoodieStorage storage, String tablePath, String latestCommitTime, Schema dataSchema, Schema requestedSchema, Option internalSchemaOpt, HoodieTableMetaClient hoodieTableMetaClient, TypedProperties props, boolean shouldUseRecordPosition, boolean allowInflightInstants, boolean emitDelete, boolean sortOutput, - InputSplit inputSplit, Option updateCallback, boolean enableOptimizedLogBlockScan) { + InputSplit inputSplit, List> updateCallback, boolean enableOptimizedLogBlockScan) { this.readerContext = readerContext; - this.fileGroupUpdateCallback = updateCallback; + this.fileGroupUpdateCallbacks = updateCallback; this.metaClient = hoodieTableMetaClient; this.storage = storage; this.enableOptimizedLogBlockScan = enableOptimizedLogBlockScan; @@ -462,7 +462,7 @@ public static class Builder { private boolean emitDelete; private boolean sortOutput = false; private boolean enableOptimizedLogBlockScan = false; - private Option fileGroupUpdateCallback = Option.empty(); + private List> fileGroupUpdateCallbacks = Collections.emptyList(); private Option> recordIteratorOpt = Option.empty(); public Builder withReaderContext(HoodieReaderContext readerContext) { @@ -548,8 +548,8 @@ public Builder withEmitDelete(boolean emitDelete) { return this; } - public Builder withFileGroupUpdateCallback(Option fileGroupUpdateCallback) { - this.fileGroupUpdateCallback = fileGroupUpdateCallback; + public Builder withFileGroupUpdateCallback(List> fileGroupUpdateCallbacks) { + this.fileGroupUpdateCallbacks = fileGroupUpdateCallbacks; return this; } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java index 01589d85e3773..4b4523450d2af 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java @@ -46,11 +46,7 @@ public InputBasedFileGroupRecordBuffer(HoodieReaderContext readerContext, Option orderingFieldName, Iterator inputRecordIterator, UpdateProcessor updateProcessor) { -<<<<<<< HEAD super(readerContext, hoodieTableMetaClient, recordMergeMode, partialUpdateMode, props, orderingFieldName, updateProcessor); -======= - super(readerContext, hoodieTableMetaClient, recordMergeMode,partialUpdateMode, props, readStats, orderingFieldName, updateProcessor); ->>>>>>> 6d74b6820d3 (Support record iterator based file group reader merge handle) this.inputRecordIterator = inputRecordIterator; // Read input records into the buffer. diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java index a0f3e0eecea05..df5ea6dea93bf 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java @@ -19,8 +19,11 @@ package org.apache.hudi.common.table.read; import org.apache.hudi.common.engine.HoodieReaderContext; -import org.apache.hudi.common.util.Option; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.List; /** * Interface used within the {@link HoodieFileGroupReader} for processing updates to records in Merge-on-Read tables. * Note that the updates are always relative to the base file's current state. @@ -38,10 +41,10 @@ public interface UpdateProcessor { T processUpdate(String recordKey, T previousRecord, T mergedRecord, boolean isDelete); static UpdateProcessor create(HoodieReadStats readStats, HoodieReaderContext readerContext, - boolean emitDeletes, Option updateCallback) { + boolean emitDeletes, List> updateCallbacks) { UpdateProcessor handler = new StandardUpdateProcessor<>(readStats, readerContext, emitDeletes); - if (updateCallback.isPresent()) { - return new CallbackProcessor<>(updateCallback.get(), handler); + if (!updateCallbacks.isEmpty()) { + return new CallbackProcessor<>(updateCallbacks, handler, readerContext); } return handler; } @@ -86,11 +89,14 @@ public T processUpdate(String recordKey, T previousRecord, T mergedRecord, boole * @param the engine specific record type */ class CallbackProcessor implements UpdateProcessor { - private final BaseFileUpdateCallback callback; + private static final Logger LOG = LoggerFactory.getLogger(CallbackProcessor.class); + private final List> callbacks; private final UpdateProcessor delegate; - public CallbackProcessor(BaseFileUpdateCallback callback, UpdateProcessor delegate) { - this.callback = callback; + public CallbackProcessor(List> callbacks, + UpdateProcessor delegate, + HoodieReaderContext readerContext) { + this.callbacks = callbacks; this.delegate = delegate; } @@ -99,11 +105,29 @@ public T processUpdate(String recordKey, T previousRecord, T currentRecord, bool T result = delegate.processUpdate(recordKey, previousRecord, currentRecord, isDelete); if (isDelete) { - callback.onDelete(recordKey, previousRecord); + for (BaseFileUpdateCallback callback : callbacks) { + try { + callback.onDelete(recordKey, previousRecord); + } catch (Exception e) { + LOG.error("Callback failed", e); + } + } } else if (previousRecord != null && previousRecord != currentRecord) { - callback.onUpdate(recordKey, previousRecord, currentRecord); + for (BaseFileUpdateCallback callback : callbacks) { + try { + callback.onUpdate(recordKey, previousRecord, currentRecord); + } catch (Exception e) { + LOG.error("Callback failed", e); + } + } } else { - callback.onInsert(recordKey, currentRecord); + for (BaseFileUpdateCallback callback : callbacks) { + try { + callback.onInsert(recordKey, currentRecord); + } catch (Exception e) { + LOG.error("Callback failed", e); + } + } } return result; } diff --git a/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java b/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java index e8429383179a7..1b156b3cdd4a2 100644 --- a/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java +++ b/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java @@ -1034,7 +1034,7 @@ private static ClosableIterator> getLogRecords(List recordBuffer = new KeyBasedFileGroupRecordBuffer<>(readerContext, datasetMetaClient, readerContext.getMergeMode(), PartialUpdateMode.NONE, properties, Option.ofNullable(tableConfig.getPreCombineField()), - UpdateProcessor.create(readStats, readerContext, true, Option.empty())); + UpdateProcessor.create(readStats, readerContext, true, Collections.emptyList())); // CRITICAL: Ensure allowInflightInstants is set to true try (HoodieMergedLogRecordReader mergedLogRecordReader = HoodieMergedLogRecordReader.newBuilder() diff --git a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestKeyBasedFileGroupRecordBuffer.java b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestKeyBasedFileGroupRecordBuffer.java index dd2a281b295d1..8d988fd7ad9eb 100644 --- a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestKeyBasedFileGroupRecordBuffer.java +++ b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestKeyBasedFileGroupRecordBuffer.java @@ -270,7 +270,7 @@ private static KeyBasedFileGroupRecordBuffer buildKeyBasedFileGro HoodieTableMetaClient mockMetaClient = mock(HoodieTableMetaClient.class, RETURNS_DEEP_STUBS); when(mockMetaClient.getTableConfig()).thenReturn(tableConfig); TypedProperties props = new TypedProperties(); - UpdateProcessor updateProcessor = UpdateProcessor.create(readStats, readerContext, false, Option.empty()); + UpdateProcessor updateProcessor = UpdateProcessor.create(readStats, readerContext, false, Collections.emptyList()); return new KeyBasedFileGroupRecordBuffer<>( readerContext, mockMetaClient, recordMergeMode, PartialUpdateMode.NONE, props, orderingFieldName, updateProcessor); } diff --git a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestSortedKeyBasedFileGroupRecordBuffer.java b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestSortedKeyBasedFileGroupRecordBuffer.java index 439bcf8f920e1..b8268f2d2e518 100644 --- a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestSortedKeyBasedFileGroupRecordBuffer.java +++ b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestSortedKeyBasedFileGroupRecordBuffer.java @@ -127,7 +127,7 @@ private SortedKeyBasedFileGroupRecordBuffer buildSortedKeyBasedFileG RecordMergeMode recordMergeMode = RecordMergeMode.COMMIT_TIME_ORDERING; PartialUpdateMode partialUpdateMode = PartialUpdateMode.NONE; TypedProperties props = new TypedProperties(); - UpdateProcessor updateProcessor = UpdateProcessor.create(readStats, mockReaderContext, false, Option.empty()); + UpdateProcessor updateProcessor = UpdateProcessor.create(readStats, mockReaderContext, false, Collections.emptyList()); return new SortedKeyBasedFileGroupRecordBuffer<>( mockReaderContext, mockMetaClient, recordMergeMode, partialUpdateMode, props, Option.empty(), updateProcessor); } diff --git a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/cdc/CDCFileGroupIterator.scala b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/cdc/CDCFileGroupIterator.scala index 855a93c4dc8e0..34f422f7e9a96 100644 --- a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/cdc/CDCFileGroupIterator.scala +++ b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/cdc/CDCFileGroupIterator.scala @@ -503,7 +503,7 @@ class CDCFileGroupIterator(split: HoodieCDCFileGroupSplit, val recordBuffer = new KeyBasedFileGroupRecordBuffer[InternalRow](readerContext, metaClient, readerContext.getMergeMode, metaClient.getTableConfig.getPartialUpdateMode, readerProperties, Option.ofNullable(metaClient.getTableConfig.getPreCombineField), - UpdateProcessor.create(stats, readerContext, true, Option.empty())) + UpdateProcessor.create(stats, readerContext, true, Collections.emptyList())) HoodieMergedLogRecordReader.newBuilder[InternalRow] .withStorage(metaClient.getStorage) From a8d862f10b998704399c0cb4f8e20742ce7f0862 Mon Sep 17 00:00:00 2001 From: Lin Liu Date: Tue, 22 Jul 2025 07:36:49 -0700 Subject: [PATCH 05/10] Refactored --- .../io/FileGroupReaderBasedMergeHandle.java | 100 ++++++++---------- .../hudi/io/HoodieMergeHandleFactory.java | 33 ++++++ .../commit/BaseJavaCommitActionExecutor.java | 8 ++ .../commit/BaseSparkCommitActionExecutor.java | 14 ++- .../table/read/HoodieFileGroupReader.java | 35 +++--- .../common/table/read/UpdateProcessor.java | 1 + 6 files changed, 119 insertions(+), 72 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 824541db900c6..8ae011bd44052 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 @@ -75,91 +75,72 @@ /** * A merge handle implementation based on the {@link HoodieFileGroupReader}. *

- * This merge handle is used for compaction, which passes a file slice from the - * compaction operation of a single file group to a file group reader, get an iterator of - * the records, and writes the records to a new base file. + * This merge handle is used for + * 1. compaction, which passes a file slice from the compaction operation of a single file group + * to a file group reader, get an iterator of the records, and writes the records to a new base file. + * 2. cow write, which takes an iterator of hoodie records, and a base file, and merge them, + * and write to a new base file. */ @NotThreadSafe public class FileGroupReaderBasedMergeHandle extends HoodieWriteMergeHandle { private static final Logger LOG = LoggerFactory.getLogger(FileGroupReaderBasedMergeHandle.class); - private final HoodieReaderContext readerContext; - private final CompactionOperation operation; private final String maxInstantTime; - private HoodieReadStats readStats; private final HoodieRecord.HoodieRecordType recordType; private final Option cdcLogger; - private final Iterator> recordIterator; + private final Option operation; + private final Option>> recordIterator; + private HoodieReadStats readStats; + /** + * For compactor. + */ public FileGroupReaderBasedMergeHandle(HoodieWriteConfig config, String instantTime, HoodieTable hoodieTable, CompactionOperation operation, TaskContextSupplier taskContextSupplier, + Option keyGeneratorOpt, HoodieReaderContext readerContext, String maxInstantTime, HoodieRecord.HoodieRecordType enginRecordType) { - super(config, instantTime, operation.getPartitionPath(), operation.getFileId(), hoodieTable, taskContextSupplier); - this.maxInstantTime = maxInstantTime; - this.keyToNewRecords = Collections.emptyMap(); - this.readerContext = readerContext; - this.operation = operation; - // 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( - 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(); - } - init(operation, this.partitionPath); - this.recordIterator = null; + this(config, instantTime, hoodieTable, null, operation.getPartitionPath(), operation.getFileId(), + taskContextSupplier, keyGeneratorOpt, readerContext, maxInstantTime, enginRecordType, Option.of(operation)); } /** - * FG reader based generic merge handle, which is not just for compaction. + * For generic FG reader based merge handle. */ public FileGroupReaderBasedMergeHandle(HoodieWriteConfig config, String instantTime, HoodieTable hoodieTable, Iterator> recordItr, String partitionPath, String fileId, TaskContextSupplier taskContextSupplier, Option keyGeneratorOpt, HoodieReaderContext readerContext, String maxInstantTime, - HoodieRecord.HoodieRecordType enginRecordType) { + HoodieRecord.HoodieRecordType enginRecordType, Option operation) { super(config, instantTime, hoodieTable, recordItr, partitionPath, fileId, taskContextSupplier, keyGeneratorOpt); + this.recordIterator = Option.ofNullable(recordItr); + this.operation = operation; + + // Common attributes. this.maxInstantTime = maxInstantTime; this.keyToNewRecords = Collections.emptyMap(); this.readerContext = readerContext; - this.recordIterator = recordItr; - this.operation = null; - if (hoodieTable.getMetaClient().getTableConfig().isCDCEnabled()) { - this.cdcLogger = Option.of(new HoodieCDCLogger( - 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(); - } - // 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())) + this.recordType = hoodieTable.isMetadataTable() || HFILE.getFileExtension().equals(hoodieTable.getBaseFileExtension()) ? HoodieRecord.HoodieRecordType.AVRO : enginRecordType; + this.cdcLogger = hoodieTable.getMetaClient().getTableConfig().isCDCEnabled() + ? Option.of(new HoodieCDCLogger(instantTime, config, hoodieTable.getMetaClient().getTableConfig(), + partitionPath, storage, getWriterSchema(), + createLogWriter(instantTime, HoodieCDCUtils.CDC_LOGFILE_SUFFIX, Option.empty()), + IOUtils.getMaxMemoryPerPartitionMerge(taskContextSupplier, config))) + : Option.empty(); init(operation, this.partitionPath); } - private void init(CompactionOperation operation, String partitionPath) { + private void init(Option operation, String partitionPath) { LOG.info("partitionPath:{}, fileId to be merged:{}", partitionPath, fileId); - this.baseFileToMerge = operation.getBaseFile(config.getBasePath(), operation.getPartitionPath()).orElse(null); + this.writtenRecordKeys = new HashSet<>(); writeStatus.setStat(new HoodieWriteStat()); - writeStatus.getStat().setTotalLogSizeCompacted( - operation.getMetrics().get(CompactionStrategy.TOTAL_LOG_FILE_SIZE).longValue()); + if (operation.isPresent()) { + writeStatus.getStat().setTotalLogSizeCompacted( + operation.get().getMetrics().get(CompactionStrategy.TOTAL_LOG_FILE_SIZE).longValue()); + } + try { Option latestValidFilePath = Option.empty(); if (baseFileToMerge != null) { @@ -218,11 +199,11 @@ public void doMerge() { TypedProperties props = TypedProperties.copy(config.getProps()); long maxMemoryPerCompaction = IOUtils.getMaxMemoryPerCompaction(taskContextSupplier, config); props.put(HoodieMemoryConfig.MAX_MEMORY_FOR_MERGE.key(), String.valueOf(maxMemoryPerCompaction)); - Stream logFiles = operation.getDeltaFileNames().stream().map(logFileName -> + Stream logFiles = operation.get().getDeltaFileNames().stream().map(logFileName -> new HoodieLogFile(new StoragePath(FSUtils.constructAbsolutePath( - config.getBasePath(), operation.getPartitionPath()), logFileName))); - Iterator engineRecordIterator = recordIterator == null - ? null : new MappingIterator<>(recordIterator, HoodieRecord::getData); + config.getBasePath(), operation.get().getPartitionPath()), logFileName))); + Option> engineRecordIterator = recordIterator.isEmpty() + ? Option.empty() : Option.of(new MappingIterator<>(recordIterator.get(), HoodieRecord::getData)); // Initializes file group reader try (HoodieFileGroupReader fileGroupReader = HoodieFileGroupReader.newBuilder() .withReaderContext(readerContext).withHoodieTableMetaClient(hoodieTable.getMetaClient()) @@ -312,7 +293,10 @@ public List close() { writeStatus.getStat().setTotalLogBlocks(readStats.getTotalLogBlocks()); writeStatus.getStat().setTotalCorruptLogBlock(readStats.getTotalCorruptLogBlock()); writeStatus.getStat().setTotalRollbackBlocks(readStats.getTotalRollbackBlocks()); - writeStatus.getStat().setTotalLogSizeCompacted(operation.getMetrics().get(CompactionStrategy.TOTAL_LOG_FILE_SIZE).longValue()); + if (operation.isPresent()) { + writeStatus.getStat().setTotalLogSizeCompacted( + operation.get().getMetrics().get(CompactionStrategy.TOTAL_LOG_FILE_SIZE).longValue()); + } if (writeStatus.getStat().getRuntimeStats() != null) { writeStatus.getStat().getRuntimeStats().setTotalScanTime(readStats.getTotalLogReadTimeMs()); 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 a0fe1291a9322..cbca7f006ad06 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 @@ -77,6 +77,39 @@ public static HoodieMergeHandle create( writeConfig, instantTime, table, recordItr, partitionPath, fileId, taskContextSupplier, keyGeneratorOpt); } + /** + * Creates a FG reader based handle for generic write path. + */ + + public static HoodieMergeHandle create( + WriteOperationType operationType, + HoodieWriteConfig writeConfig, + String instantTime, + HoodieTable table, + Iterator> recordItr, + String partitionPath, + String fileId, + TaskContextSupplier taskContextSupplier, + Option keyGeneratorOpt, + HoodieReaderContext readerContext, + HoodieRecord.HoodieRecordType recordType) { + + boolean isFallbackEnabled = writeConfig.isMergeHandleFallbackEnabled(); + Pair mergeHandleClasses = getMergeHandleClassesWrite(operationType, writeConfig, table); + String logContext = String.format("for fileId %s and partition path %s at commit %s", fileId, partitionPath, instantTime); + LOG.info("Create HoodieMergeHandle implementation {} {}", mergeHandleClasses.getLeft(), logContext); + + Class[] constructorParamTypes = new Class[] { + HoodieWriteConfig.class, String.class, HoodieTable.class, Iterator.class, + String.class, String.class, TaskContextSupplier.class, Option.class, HoodieReaderContext.class, HoodieRecord.HoodieRecordType.class, + Option.class + }; + + return instantiateMergeHandle( + isFallbackEnabled, mergeHandleClasses.getLeft(), mergeHandleClasses.getRight(), logContext, constructorParamTypes, + writeConfig, instantTime, table, recordItr, partitionPath, fileId, taskContextSupplier, keyGeneratorOpt, readerContext, recordType, Option.empty()); + } + /** * Creates a merge handle for compaction path. */ diff --git a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/commit/BaseJavaCommitActionExecutor.java b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/commit/BaseJavaCommitActionExecutor.java index 092daf2cead25..291523474d620 100644 --- a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/commit/BaseJavaCommitActionExecutor.java +++ b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/commit/BaseJavaCommitActionExecutor.java @@ -21,6 +21,7 @@ import org.apache.hudi.client.WriteStatus; import org.apache.hudi.common.data.HoodieListData; import org.apache.hudi.common.engine.HoodieEngineContext; +import org.apache.hudi.common.engine.HoodieReaderContext; import org.apache.hudi.common.model.HoodieKey; import org.apache.hudi.common.model.HoodieRecord; import org.apache.hudi.common.model.HoodieRecordLocation; @@ -37,6 +38,7 @@ import org.apache.hudi.exception.HoodieUpsertException; import org.apache.hudi.execution.JavaLazyInsertIterable; import org.apache.hudi.io.CreateHandleFactory; +import org.apache.hudi.io.FileGroupReaderBasedMergeHandle; import org.apache.hudi.io.HoodieMergeHandle; import org.apache.hudi.io.HoodieMergeHandleFactory; import org.apache.hudi.io.IOUtils; @@ -256,6 +258,12 @@ public Iterator> handleUpdate(String partitionPath, String fil + "columns are disabled. Please choose the right key generator if you wish to disable meta fields.", e); } } + if (config.getMergeHandleClassName().equals(FileGroupReaderBasedMergeHandle.class.getName())) { + HoodieReaderContext readerContext = table.getContext().getReaderContextFactory(table.getMetaClient()).getContext(); + return HoodieMergeHandleFactory.create( + operationType, config, instantTime, table, recordItr, partitionPath, fileId, taskContextSupplier, keyGeneratorOpt, + readerContext, HoodieRecord.HoodieRecordType.AVRO); + } return HoodieMergeHandleFactory.create(operationType, config, instantTime, table, recordItr, partitionPath, fileId, taskContextSupplier, keyGeneratorOpt); } diff --git a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java index b2119ddf7ce2f..cd680b1f4db4c 100644 --- a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java +++ b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java @@ -19,6 +19,7 @@ package org.apache.hudi.table.action.commit; import org.apache.hudi.client.utils.SparkPartitionUtils; +import org.apache.hudi.common.engine.HoodieReaderContext; import org.apache.hudi.index.HoodieSparkIndexClient; import org.apache.hudi.client.WriteStatus; import org.apache.hudi.client.clustering.update.strategy.SparkAllowUpdateStrategy; @@ -47,6 +48,7 @@ import org.apache.hudi.exception.HoodieUpsertException; import org.apache.hudi.execution.SparkLazyInsertIterable; import org.apache.hudi.index.HoodieIndex; +import org.apache.hudi.io.FileGroupReaderBasedMergeHandle; import org.apache.hudi.io.HoodieMergeHandle; import org.apache.hudi.io.CreateHandleFactory; import org.apache.hudi.io.HoodieMergeHandleFactory; @@ -387,8 +389,16 @@ public Iterator> handleUpdate(String partitionPath, String fil } protected HoodieMergeHandle getUpdateHandle(String partitionPath, String fileId, Iterator> recordItr) { - HoodieMergeHandle mergeHandle = HoodieMergeHandleFactory.create(operationType, config, instantTime, table, recordItr, partitionPath, fileId, - taskContextSupplier, keyGeneratorOpt); + HoodieMergeHandle mergeHandle; + if (config.getMergeHandleClassName().equals(FileGroupReaderBasedMergeHandle.class.getName())) { + HoodieReaderContext readerContext = table.getContext().getReaderContextFactory(table.getMetaClient()).getContext(); + mergeHandle = HoodieMergeHandleFactory.create( + operationType, config, instantTime, table, recordItr, partitionPath, fileId, taskContextSupplier, keyGeneratorOpt, + readerContext, HoodieRecord.HoodieRecordType.SPARK); + } else { + mergeHandle = HoodieMergeHandleFactory.create( + operationType, config, instantTime, table, recordItr, partitionPath, fileId, taskContextSupplier, keyGeneratorOpt); + } if (mergeHandle.getOldFilePath() != null && mergeHandle.baseFileForMerge().getBootstrapBaseFile().isPresent()) { Option partitionFields = table.getMetaClient().getTableConfig().getPartitionFields(); Object[] partitionValues = SparkPartitionUtils.getPartitionFieldVals(partitionFields, mergeHandle.getPartitionPath(), diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java index b2ebaa58c3b22..1b529389d5073 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java @@ -75,7 +75,7 @@ public final class HoodieFileGroupReader implements Closeable { private final HoodieReaderContext readerContext; private final HoodieTableMetaClient metaClient; - private final InputSplit inputSplit; + private final InputSplit inputSplit; private final Option partitionPathFields; private final Option orderingFieldName; private final HoodieStorage storage; @@ -117,7 +117,7 @@ private HoodieFileGroupReader(HoodieReaderContext readerContext, HoodieStorag String latestCommitTime, Schema dataSchema, Schema requestedSchema, Option internalSchemaOpt, HoodieTableMetaClient hoodieTableMetaClient, TypedProperties props, boolean shouldUseRecordPosition, boolean allowInflightInstants, boolean emitDelete, boolean sortOutput, - InputSplit inputSplit, List> updateCallback, boolean enableOptimizedLogBlockScan) { + InputSplit inputSplit, List> updateCallback, boolean enableOptimizedLogBlockScan) { this.readerContext = readerContext; this.fileGroupUpdateCallbacks = updateCallback; this.metaClient = hoodieTableMetaClient; @@ -155,7 +155,7 @@ private HoodieFileGroupReader(HoodieReaderContext readerContext, HoodieStorag this.readStats = new HoodieReadStats(); this.recordBuffer = getRecordBuffer(readerContext, hoodieTableMetaClient, readerContext.getMergeMode(), tableConfig.getPartialUpdateMode(), props, - isSkipMerge, shouldUseRecordPosition, readStats, emitDelete, sortOutput, Option.empty()); + isSkipMerge, shouldUseRecordPosition, readStats, emitDelete, sortOutput, inputSplit.recordIterator); this.allowInflightInstants = allowInflightInstants; } @@ -173,7 +173,7 @@ private FileGroupRecordBuffer getRecordBuffer(HoodieReaderContext readerCo boolean emitDelete, boolean sortOutput, Option> inputRecordOpt) { - UpdateProcessor updateProcessor = UpdateProcessor.create(readStats, readerContext, emitDelete, fileGroupUpdateCallback); + UpdateProcessor updateProcessor = UpdateProcessor.create(readStats, readerContext, emitDelete, fileGroupUpdateCallbacks); if (inputSplit.logFiles.isEmpty()) { if (inputRecordOpt.isPresent()) { return new InputBasedFileGroupRecordBuffer<>( @@ -463,7 +463,7 @@ public static class Builder { private boolean sortOutput = false; private boolean enableOptimizedLogBlockScan = false; private List> fileGroupUpdateCallbacks = Collections.emptyList(); - private Option> recordIteratorOpt = Option.empty(); + private Option> recordIterator = Option.empty(); public Builder withReaderContext(HoodieReaderContext readerContext) { this.readerContext = readerContext; @@ -569,8 +569,8 @@ public Builder withSortOutput(boolean sortOutput) { return this; } - public Builder withRecordIterator(Iterator iterator) { - this.recordIteratorOpt = Option.ofNullable(iterator); + public Builder withRecordIterator(Option> iterator) { + this.recordIterator = iterator; return this; } @@ -590,14 +590,14 @@ public HoodieFileGroupReader build() { ValidationUtils.checkArgument(logFiles != null, "Log files stream is required"); ValidationUtils.checkArgument(partitionPath != null, "Partition path is required"); - InputSplit inputSplit = new InputSplit(baseFileOption, logFiles, partitionPath, start, length); + InputSplit inputSplit = new InputSplit<>(baseFileOption, logFiles, partitionPath, start, length, recordIterator); return new HoodieFileGroupReader<>( readerContext, storage, tablePath, latestCommitTime, dataSchema, requestedSchema, internalSchemaOpt, hoodieTableMetaClient, - props, shouldUseRecordPosition, allowInflightInstants, emitDelete, sortOutput, inputSplit, fileGroupUpdateCallback, enableOptimizedLogBlockScan); + props, shouldUseRecordPosition, allowInflightInstants, emitDelete, sortOutput, inputSplit, fileGroupUpdateCallbacks, enableOptimizedLogBlockScan); } } - private static class InputSplit { + private static class InputSplit { private final Option baseFileOption; private final List logFiles; private final String partitionPath; @@ -605,6 +605,18 @@ private static class InputSplit { private final long start; // Length of bytes to read from the base file private final long length; + // For input record case. + private Option> recordIterator; + + InputSplit(Option baseFileOption, + Stream logFiles, + String partitionPath, + long start, + long length, + Option> recordIterator) { + this(baseFileOption, logFiles, partitionPath, start, length); + this.recordIterator = recordIterator; + } InputSplit(Option baseFileOption, Stream logFiles, String partitionPath, long start, long length) { this.baseFileOption = baseFileOption; @@ -617,8 +629,7 @@ private static class InputSplit { } static InputSplit fromFileSlice(FileSlice fileSlice, long start, long length) { - return new InputSplit(fileSlice.getBaseFile(), fileSlice.getLogFiles(), fileSlice.getPartitionPath(), - start, length); + return new InputSplit(fileSlice.getBaseFile(), fileSlice.getLogFiles(), fileSlice.getPartitionPath(), start, length); } } } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java index df5ea6dea93bf..7f916a78d4504 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java @@ -24,6 +24,7 @@ import org.slf4j.LoggerFactory; import java.util.List; + /** * Interface used within the {@link HoodieFileGroupReader} for processing updates to records in Merge-on-Read tables. * Note that the updates are always relative to the base file's current state. From e10caf4bafb96a6e7089a0c343be634d732b2426 Mon Sep 17 00:00:00 2001 From: Lin Liu Date: Tue, 22 Jul 2025 07:58:16 -0700 Subject: [PATCH 06/10] Address comments --- .../io/FileGroupReaderBasedMergeHandle.java | 12 +++++- .../table/read/BaseFileUpdateCallback.java | 5 +++ .../common/table/read/UpdateProcessor.java | 43 +++++++++---------- 3 files changed, 37 insertions(+), 23 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 8ae011bd44052..8efa0a072bd23 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 @@ -258,7 +258,7 @@ private List> createCallbacks() { callbacks.add(new CDCCallback<>(cdcLogger.get(), readerContext)); } // Stream secondary index stats. - if (isSecondaryIndexStatsStreamingWritesEnabled || writeStatus.isTrackingSuccessfulWrites()) { + if (isSecondaryIndexStatsStreamingWritesEnabled) { callbacks.add(new SecondaryIndexCallback<>( partitionPath, writeSchemaWithMetaFields, @@ -340,6 +340,11 @@ public void onDelete(String recordKey, T previousRecord) { } + @Override + public String getName() { + return "CdcCallBack"; + } + private GenericRecord convertOutput(T record) { T convertedRecord = outputConverter.get().map(converter -> record == null ? null : converter.apply(record)).orElse(record); return convertedRecord == null ? null : readerContext.convertToAvroRecord(convertedRecord, requestedSchema.get()); @@ -429,5 +434,10 @@ public void onDelete(String recordKey, T previousRecord) { keyGeneratorOpt, config); } + + @Override + public String getName() { + return "SecondaryIndex Callback"; + } } } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/BaseFileUpdateCallback.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/BaseFileUpdateCallback.java index 6685cd6cef459..2bbb4f62a8044 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/BaseFileUpdateCallback.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/BaseFileUpdateCallback.java @@ -44,4 +44,9 @@ public interface BaseFileUpdateCallback { * @param previousRecord the record in the base file before deletion */ void onDelete(String recordKey, T previousRecord); + + /** + * Callback method to return the name of the callback. + */ + String getName(); } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java index 7f916a78d4504..774ec6294e21a 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java @@ -87,7 +87,8 @@ public T processUpdate(String recordKey, T previousRecord, T mergedRecord, boole /** * A processor that wraps the standard update processor and invokes a customizable callback for each update. - * @param the engine specific record type + * + * @param the engine-specific record type */ class CallbackProcessor implements UpdateProcessor { private static final Logger LOG = LoggerFactory.getLogger(CallbackProcessor.class); @@ -106,31 +107,29 @@ public T processUpdate(String recordKey, T previousRecord, T currentRecord, bool T result = delegate.processUpdate(recordKey, previousRecord, currentRecord, isDelete); if (isDelete) { - for (BaseFileUpdateCallback callback : callbacks) { - try { - callback.onDelete(recordKey, previousRecord); - } catch (Exception e) { - LOG.error("Callback failed", e); - } - } + invokeCallbacks(callback -> callback.onDelete(recordKey, previousRecord)); } else if (previousRecord != null && previousRecord != currentRecord) { - for (BaseFileUpdateCallback callback : callbacks) { - try { - callback.onUpdate(recordKey, previousRecord, currentRecord); - } catch (Exception e) { - LOG.error("Callback failed", e); - } - } + invokeCallbacks(callback -> callback.onUpdate(recordKey, previousRecord, currentRecord)); } else { - for (BaseFileUpdateCallback callback : callbacks) { - try { - callback.onInsert(recordKey, currentRecord); - } catch (Exception e) { - LOG.error("Callback failed", e); - } - } + invokeCallbacks(callback -> callback.onInsert(recordKey, currentRecord)); } + return result; } + + private void invokeCallbacks(CallbackInvoker invoker) { + for (BaseFileUpdateCallback callback : callbacks) { + try { + invoker.invoke(callback); + } catch (Exception e) { + LOG.error(String.format("Callback %s failed: ", callback.getName()), e); + } + } + } + + @FunctionalInterface + private interface CallbackInvoker { + void invoke(BaseFileUpdateCallback callback) throws Exception; + } } } From a16e427f5d9f1bb7743a24ec081cabc5b193ddd3 Mon Sep 17 00:00:00 2001 From: Lin Liu Date: Wed, 23 Jul 2025 06:48:23 -0700 Subject: [PATCH 07/10] Address comments --- .../io/FileGroupReaderBasedMergeHandle.java | 60 ++++++++++++------- .../hudi/io/HoodieMergeHandleFactory.java | 7 ++- .../commit/BaseSparkCommitActionExecutor.java | 8 ++- 3 files changed, 47 insertions(+), 28 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 8efa0a072bd23..5f27a12dccaf3 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 @@ -100,7 +100,7 @@ public FileGroupReaderBasedMergeHandle(HoodieWriteConfig config, String instantT Option keyGeneratorOpt, HoodieReaderContext readerContext, String maxInstantTime, HoodieRecord.HoodieRecordType enginRecordType) { - this(config, instantTime, hoodieTable, null, operation.getPartitionPath(), operation.getFileId(), + this(config, instantTime, hoodieTable, Option.empty(), operation.getPartitionPath(), operation.getFileId(), taskContextSupplier, keyGeneratorOpt, readerContext, maxInstantTime, enginRecordType, Option.of(operation)); } @@ -108,12 +108,12 @@ public FileGroupReaderBasedMergeHandle(HoodieWriteConfig config, String instantT * For generic FG reader based merge handle. */ public FileGroupReaderBasedMergeHandle(HoodieWriteConfig config, String instantTime, HoodieTable hoodieTable, - Iterator> recordItr, String partitionPath, String fileId, + Option>> recordItr, String partitionPath, String fileId, TaskContextSupplier taskContextSupplier, Option keyGeneratorOpt, HoodieReaderContext readerContext, String maxInstantTime, HoodieRecord.HoodieRecordType enginRecordType, Option operation) { - super(config, instantTime, hoodieTable, recordItr, partitionPath, fileId, taskContextSupplier, keyGeneratorOpt); - this.recordIterator = Option.ofNullable(recordItr); + super(config, instantTime, hoodieTable, recordItr.get(), partitionPath, fileId, taskContextSupplier, keyGeneratorOpt); + this.recordIterator = recordItr; this.operation = operation; // Common attributes. @@ -197,21 +197,35 @@ public void doMerge() { boolean usePosition = config.getBooleanOrDefault(MERGE_USE_RECORD_POSITIONS); Option internalSchemaOption = SerDeHelper.fromJson(config.getInternalSchema()); TypedProperties props = TypedProperties.copy(config.getProps()); - long maxMemoryPerCompaction = IOUtils.getMaxMemoryPerCompaction(taskContextSupplier, config); - props.put(HoodieMemoryConfig.MAX_MEMORY_FOR_MERGE.key(), String.valueOf(maxMemoryPerCompaction)); - Stream logFiles = operation.get().getDeltaFileNames().stream().map(logFileName -> - new HoodieLogFile(new StoragePath(FSUtils.constructAbsolutePath( - config.getBasePath(), operation.get().getPartitionPath()), logFileName))); - Option> engineRecordIterator = recordIterator.isEmpty() - ? Option.empty() : Option.of(new MappingIterator<>(recordIterator.get(), HoodieRecord::getData)); + Stream logFiles = Stream.empty(); + Option> engineRecordIterator = Option.empty(); + // For compaction. + if (operation.isPresent()) { + long maxMemoryPerCompaction = IOUtils.getMaxMemoryPerCompaction(taskContextSupplier, config); + props.put(HoodieMemoryConfig.MAX_MEMORY_FOR_MERGE.key(), String.valueOf(maxMemoryPerCompaction)); + logFiles = operation.get().getDeltaFileNames().stream().map(logFileName -> + new HoodieLogFile(new StoragePath(FSUtils.constructAbsolutePath( + config.getBasePath(), operation.get().getPartitionPath()), logFileName))); + } + // For generic merge. + if (recordIterator.isPresent()) { + engineRecordIterator = Option.of(new MappingIterator<>(recordIterator.get(), HoodieRecord::getData)); + usePosition = false; + } // Initializes file group reader try (HoodieFileGroupReader fileGroupReader = HoodieFileGroupReader.newBuilder() - .withReaderContext(readerContext).withHoodieTableMetaClient(hoodieTable.getMetaClient()) - .withLatestCommitTime(maxInstantTime).withPartitionPath(partitionPath) - .withBaseFileOption(Option.ofNullable(baseFileToMerge)).withLogFiles(logFiles) - .withDataSchema(writeSchemaWithMetaFields).withRequestedSchema(writeSchemaWithMetaFields) - .withInternalSchema(internalSchemaOption).withProps(props) - .withShouldUseRecordPosition(usePosition).withSortOutput(hoodieTable.requireSortedRecords()) + .withReaderContext(readerContext) + .withHoodieTableMetaClient(hoodieTable.getMetaClient()) + .withLatestCommitTime(maxInstantTime) + .withPartitionPath(partitionPath) + .withBaseFileOption(Option.ofNullable(baseFileToMerge)) + .withLogFiles(logFiles) + .withDataSchema(writeSchemaWithMetaFields) + .withRequestedSchema(writeSchemaWithMetaFields) + .withInternalSchema(internalSchemaOption) + .withProps(props) + .withShouldUseRecordPosition(usePosition) + .withSortOutput(hoodieTable.requireSortedRecords()) .withFileGroupUpdateCallback(createCallbacks()) .withRecordIterator(engineRecordIterator).build()) { // Reads the records from the file slice @@ -342,7 +356,7 @@ public void onDelete(String recordKey, T previousRecord) { @Override public String getName() { - return "CdcCallBack"; + return "CDCCallback"; } private GenericRecord convertOutput(T record) { @@ -382,14 +396,14 @@ public SecondaryIndexCallback(String partitionPath, @Override public void onUpdate(String recordKey, T previousRecord, T mergedRecord) { HoodieKey hoodieKey = new HoodieKey(recordKey, partitionPath); - BufferedRecord bufferedPrevousRecord = BufferedRecord.forRecordWithContext( + BufferedRecord bufferedPreviousRecord = BufferedRecord.forRecordWithContext( previousRecord, writeSchemaWithMetaFields, readerContext, Option.empty(), false); BufferedRecord bufferedMergedRecord = BufferedRecord.forRecordWithContext( mergedRecord, writeSchemaWithMetaFields, readerContext, Option.empty(), false); SecondaryIndexStreamingTracker.trackSecondaryIndexStats( hoodieKey, Option.of(readerContext.constructHoodieRecord(bufferedMergedRecord)), - readerContext.constructHoodieRecord(bufferedPrevousRecord), + readerContext.constructHoodieRecord(bufferedPreviousRecord), false, writeStatus, writeSchemaWithMetaFields, @@ -420,12 +434,12 @@ public void onInsert(String recordKey, T newRecord) { @Override public void onDelete(String recordKey, T previousRecord) { HoodieKey hoodieKey = new HoodieKey(recordKey, partitionPath); - BufferedRecord bufferedPrevousRecord = BufferedRecord.forRecordWithContext( + BufferedRecord bufferedPreviousRecord = BufferedRecord.forRecordWithContext( previousRecord, writeSchemaWithMetaFields, readerContext, Option.empty(), false); SecondaryIndexStreamingTracker.trackSecondaryIndexStats( hoodieKey, null, - readerContext.constructHoodieRecord(bufferedPrevousRecord), + readerContext.constructHoodieRecord(bufferedPreviousRecord), true, writeStatus, writeSchemaWithMetaFields, @@ -437,7 +451,7 @@ public void onDelete(String recordKey, T previousRecord) { @Override public String getName() { - return "SecondaryIndex Callback"; + return "SecondaryIndexCallback"; } } } 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 cbca7f006ad06..a3a3c0e2c9751 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 @@ -80,7 +80,6 @@ public static HoodieMergeHandle create( /** * Creates a FG reader based handle for generic write path. */ - public static HoodieMergeHandle create( WriteOperationType operationType, HoodieWriteConfig writeConfig, @@ -158,11 +157,15 @@ public static HoodieMergeHandle create( 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); + if (mergeHandleClass.equals(FileGroupReaderBasedMergeHandle.class.getName())) { + return new FileGroupReaderBasedMergeHandle<>( + config, instantTime, hoodieTable, operation, taskContextSupplier, Option.empty(), readerContext, maxInstantTime, recordType); + } + Class[] constructorParamTypes = new Class[] { HoodieWriteConfig.class, String.class, HoodieTable.class, CompactionOperation.class, TaskContextSupplier.class, HoodieReaderContext.class, String.class, HoodieRecord.HoodieRecordType.class }; - return instantiateMergeHandle( isFallbackEnabled, mergeHandleClass, COMPACT_MERGE_HANDLE_CLASS_NAME.defaultValue(), logContext, constructorParamTypes, config, instantTime, hoodieTable, operation, taskContextSupplier, readerContext, maxInstantTime, recordType); diff --git a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java index cd680b1f4db4c..e1de53c0e2c02 100644 --- a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java +++ b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java @@ -20,6 +20,8 @@ import org.apache.hudi.client.utils.SparkPartitionUtils; 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.index.HoodieSparkIndexClient; import org.apache.hudi.client.WriteStatus; import org.apache.hudi.client.clustering.update.strategy.SparkAllowUpdateStrategy; @@ -392,9 +394,9 @@ protected HoodieMergeHandle getUpdateHandle(String partitionPath, String fileId, HoodieMergeHandle mergeHandle; if (config.getMergeHandleClassName().equals(FileGroupReaderBasedMergeHandle.class.getName())) { HoodieReaderContext readerContext = table.getContext().getReaderContextFactory(table.getMetaClient()).getContext(); - mergeHandle = HoodieMergeHandleFactory.create( - operationType, config, instantTime, table, recordItr, partitionPath, fileId, taskContextSupplier, keyGeneratorOpt, - readerContext, HoodieRecord.HoodieRecordType.SPARK); + mergeHandle = new FileGroupReaderBasedMergeHandle( + config, instantTime, table, Option.of(recordItr), partitionPath, fileId, taskContextSupplier, keyGeneratorOpt, + readerContext, instantTime, table.getConfig().getRecordMerger().getRecordType(), Option.empty()); } else { mergeHandle = HoodieMergeHandleFactory.create( operationType, config, instantTime, table, recordItr, partitionPath, fileId, taskContextSupplier, keyGeneratorOpt); From 183531557aa88efe7e7f5bc450d1e56f86ea5545 Mon Sep 17 00:00:00 2001 From: Lin Liu Date: Wed, 23 Jul 2025 08:59:45 -0700 Subject: [PATCH 08/10] Address comments --- .../hudi/io/HoodieMergeHandleFactory.java | 38 ------------------- .../commit/BaseJavaCommitActionExecutor.java | 6 +-- .../commit/BaseSparkCommitActionExecutor.java | 2 - 3 files changed, 3 insertions(+), 43 deletions(-) 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 a3a3c0e2c9751..94fcd10313925 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 @@ -77,38 +77,6 @@ public static HoodieMergeHandle create( writeConfig, instantTime, table, recordItr, partitionPath, fileId, taskContextSupplier, keyGeneratorOpt); } - /** - * Creates a FG reader based handle for generic write path. - */ - public static HoodieMergeHandle create( - WriteOperationType operationType, - HoodieWriteConfig writeConfig, - String instantTime, - HoodieTable table, - Iterator> recordItr, - String partitionPath, - String fileId, - TaskContextSupplier taskContextSupplier, - Option keyGeneratorOpt, - HoodieReaderContext readerContext, - HoodieRecord.HoodieRecordType recordType) { - - boolean isFallbackEnabled = writeConfig.isMergeHandleFallbackEnabled(); - Pair mergeHandleClasses = getMergeHandleClassesWrite(operationType, writeConfig, table); - String logContext = String.format("for fileId %s and partition path %s at commit %s", fileId, partitionPath, instantTime); - LOG.info("Create HoodieMergeHandle implementation {} {}", mergeHandleClasses.getLeft(), logContext); - - Class[] constructorParamTypes = new Class[] { - HoodieWriteConfig.class, String.class, HoodieTable.class, Iterator.class, - String.class, String.class, TaskContextSupplier.class, Option.class, HoodieReaderContext.class, HoodieRecord.HoodieRecordType.class, - Option.class - }; - - return instantiateMergeHandle( - isFallbackEnabled, mergeHandleClasses.getLeft(), mergeHandleClasses.getRight(), logContext, constructorParamTypes, - writeConfig, instantTime, table, recordItr, partitionPath, fileId, taskContextSupplier, keyGeneratorOpt, readerContext, recordType, Option.empty()); - } - /** * Creates a merge handle for compaction path. */ @@ -156,12 +124,6 @@ public static HoodieMergeHandle create( String mergeHandleClass = config.getCompactionMergeHandleClassName(); 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); - - if (mergeHandleClass.equals(FileGroupReaderBasedMergeHandle.class.getName())) { - return new FileGroupReaderBasedMergeHandle<>( - config, instantTime, hoodieTable, operation, taskContextSupplier, Option.empty(), readerContext, maxInstantTime, recordType); - } - Class[] constructorParamTypes = new Class[] { HoodieWriteConfig.class, String.class, HoodieTable.class, CompactionOperation.class, TaskContextSupplier.class, HoodieReaderContext.class, String.class, HoodieRecord.HoodieRecordType.class diff --git a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/commit/BaseJavaCommitActionExecutor.java b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/commit/BaseJavaCommitActionExecutor.java index 291523474d620..db100993ee23e 100644 --- a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/commit/BaseJavaCommitActionExecutor.java +++ b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/commit/BaseJavaCommitActionExecutor.java @@ -260,9 +260,9 @@ public Iterator> handleUpdate(String partitionPath, String fil } if (config.getMergeHandleClassName().equals(FileGroupReaderBasedMergeHandle.class.getName())) { HoodieReaderContext readerContext = table.getContext().getReaderContextFactory(table.getMetaClient()).getContext(); - return HoodieMergeHandleFactory.create( - operationType, config, instantTime, table, recordItr, partitionPath, fileId, taskContextSupplier, keyGeneratorOpt, - readerContext, HoodieRecord.HoodieRecordType.AVRO); + return new FileGroupReaderBasedMergeHandle<>( + config, instantTime, table, Option.of(recordItr), partitionPath, fileId, + taskContextSupplier, keyGeneratorOpt, readerContext, instantTime, config.getRecordMerger().getRecordType(), Option.empty()); } return HoodieMergeHandleFactory.create(operationType, config, instantTime, table, recordItr, partitionPath, fileId, taskContextSupplier, keyGeneratorOpt); diff --git a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java index e1de53c0e2c02..afcf5ff8915d5 100644 --- a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java +++ b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java @@ -20,8 +20,6 @@ import org.apache.hudi.client.utils.SparkPartitionUtils; 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.index.HoodieSparkIndexClient; import org.apache.hudi.client.WriteStatus; import org.apache.hudi.client.clustering.update.strategy.SparkAllowUpdateStrategy; From 60ea304e4f3cfb29e801970bc56a2ba13ddb7e30 Mon Sep 17 00:00:00 2001 From: Lin Liu Date: Wed, 23 Jul 2025 20:15:02 -0700 Subject: [PATCH 09/10] Add funtional tests --- .../apache/hudi/config/HoodieWriteConfig.java | 3 +- .../io/FileGroupReaderBasedMergeHandle.java | 97 +++---- .../main/java/org/apache/hudi/io/IOUtils.java | 1 + .../commit/BaseSparkCommitActionExecutor.java | 5 +- .../common/table/read/BufferedRecord.java | 14 + .../table/read/HoodieFileGroupReader.java | 20 +- .../read/InputBasedFileGroupRecordBuffer.java | 32 ++- .../common/table/read/UpdateProcessor.java | 2 +- .../TestPayloadDeprecationFlow.scala | 257 +++++++++++++----- 9 files changed, 282 insertions(+), 149 deletions(-) diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java index a1df49505c48f..11e006ee7b2d2 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java @@ -73,7 +73,6 @@ import org.apache.hudi.index.HoodieIndex; import org.apache.hudi.io.FileGroupReaderBasedMergeHandle; import org.apache.hudi.io.HoodieConcatHandle; -import org.apache.hudi.io.HoodieWriteMergeHandle; import org.apache.hudi.keygen.SimpleAvroKeyGenerator; import org.apache.hudi.keygen.constant.KeyGeneratorOptions; import org.apache.hudi.keygen.constant.KeyGeneratorType; @@ -855,7 +854,7 @@ public class HoodieWriteConfig extends HoodieConfig { public static final ConfigProperty MERGE_HANDLE_CLASS_NAME = ConfigProperty .key("hoodie.write.merge.handle.class") - .defaultValue(HoodieWriteMergeHandle.class.getName()) + .defaultValue(FileGroupReaderBasedMergeHandle.class.getName()) .markAdvanced() .sinceVersion("1.1.0") .withDocumentation("The merge handle class that implements interface{@link HoodieMergeHandle} to merge the records " 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 5f27a12dccaf3..6cc3c0ff247ce 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 @@ -39,12 +39,10 @@ import org.apache.hudi.common.table.read.HoodieReadStats; import org.apache.hudi.common.util.Option; import org.apache.hudi.common.util.collection.ClosableIterator; -import org.apache.hudi.common.util.collection.MappingIterator; import org.apache.hudi.config.HoodieWriteConfig; import org.apache.hudi.exception.HoodieUpsertException; import org.apache.hudi.internal.schema.InternalSchema; import org.apache.hudi.internal.schema.utils.SerDeHelper; -import org.apache.hudi.io.storage.HoodieFileWriterFactory; import org.apache.hudi.keygen.BaseKeyGenerator; import org.apache.hudi.storage.StoragePath; import org.apache.hudi.table.HoodieTable; @@ -74,12 +72,8 @@ /** * A merge handle implementation based on the {@link HoodieFileGroupReader}. - *

- * This merge handle is used for - * 1. compaction, which passes a file slice from the compaction operation of a single file group - * to a file group reader, get an iterator of the records, and writes the records to a new base file. - * 2. cow write, which takes an iterator of hoodie records, and a base file, and merge them, - * and write to a new base file. + * This handle uses {@link HoodieFileGroupReader} to merge either a file slice or merge + * a base file and a set of input records. */ @NotThreadSafe public class FileGroupReaderBasedMergeHandle extends HoodieWriteMergeHandle { @@ -112,10 +106,11 @@ public FileGroupReaderBasedMergeHandle(HoodieWriteConfig config, String instantT TaskContextSupplier taskContextSupplier, Option keyGeneratorOpt, HoodieReaderContext readerContext, String maxInstantTime, HoodieRecord.HoodieRecordType enginRecordType, Option operation) { - super(config, instantTime, hoodieTable, recordItr.get(), partitionPath, fileId, taskContextSupplier, keyGeneratorOpt); + super(config, instantTime, hoodieTable, Collections.emptyIterator(), partitionPath, fileId, taskContextSupplier, keyGeneratorOpt); + // For regular merge process, recordIterator exists and operation does not. this.recordIterator = recordItr; + // For compaction process, operation exists and recordIterator does not. this.operation = operation; - // Common attributes. this.maxInstantTime = maxInstantTime; this.keyToNewRecords = Collections.emptyMap(); @@ -141,51 +136,35 @@ private void init(Option operation, String partitionPath) { operation.get().getMetrics().get(CompactionStrategy.TOTAL_LOG_FILE_SIZE).longValue()); } - try { - Option latestValidFilePath = Option.empty(); - if (baseFileToMerge != null) { - latestValidFilePath = Option.of(baseFileToMerge.getFileName()); - writeStatus.getStat().setPrevCommit(baseFileToMerge.getCommitTime()); - // At the moment, we only support SI for overwrite with latest payload. So, we don't need to embed entire file slice here. - // HUDI-8518 will be taken up to fix it for any payload during which we might require entire file slice to be set here. - // Already AppendHandle adds all logs file from current file slice to HoodieDeltaWriteStat. - writeStatus.getStat().setPrevBaseFile(latestValidFilePath.get()); - } else { - writeStatus.getStat().setPrevCommit(HoodieWriteStat.NULL_COMMIT); - } - - HoodiePartitionMetadata partitionMetadata = new HoodiePartitionMetadata(storage, instantTime, - new StoragePath(config.getBasePath()), - FSUtils.constructAbsolutePath(config.getBasePath(), partitionPath), - hoodieTable.getPartitionMetafileFormat()); - partitionMetadata.trySave(); - - String newFileName = FSUtils.makeBaseFileName(instantTime, writeToken, fileId, hoodieTable.getBaseFileExtension()); - makeOldAndNewFilePaths(partitionPath, - latestValidFilePath.isPresent() ? latestValidFilePath.get() : null, newFileName); - - LOG.info("Merging data from file group {}, to a new base file {}", fileId, newFilePath); - // file name is same for all records, in this bunch - writeStatus.setFileId(fileId); - writeStatus.setPartitionPath(partitionPath); - writeStatus.getStat().setPartitionPath(partitionPath); - writeStatus.getStat().setFileId(fileId); - setWriteStatusPath(); - - // Create Marker file, - // uses name of `newFilePath` instead of `newFileName` - // in case the sub-class may roll over the file handle name. - createMarkerFile(partitionPath, newFilePath.getName()); - - // Create the writer for writing the new version file - fileWriter = HoodieFileWriterFactory.getFileWriter(instantTime, newFilePath, hoodieTable.getStorage(), - config, writeSchemaWithMetaFields, taskContextSupplier, recordType); - } catch (IOException io) { - LOG.error("Error in update task at commit {}", instantTime, io); - writeStatus.setGlobalError(io); - throw new HoodieUpsertException("Failed to initialize HoodieUpdateHandle for FileId: " + fileId + " on commit " - + instantTime + " on path " + hoodieTable.getMetaClient().getBasePath(), io); + Option latestValidFilePath = Option.empty(); + if (baseFileToMerge != null) { + latestValidFilePath = Option.of(baseFileToMerge.getFileName()); + writeStatus.getStat().setPrevCommit(baseFileToMerge.getCommitTime()); + // At the moment, we only support SI for overwrite with latest payload. So, we don't need to embed entire file slice here. + // HUDI-8518 will be taken up to fix it for any payload during which we might require entire file slice to be set here. + // Already AppendHandle adds all logs file from current file slice to HoodieDeltaWriteStat. + writeStatus.getStat().setPrevBaseFile(latestValidFilePath.get()); + } else { + writeStatus.getStat().setPrevCommit(HoodieWriteStat.NULL_COMMIT); } + + HoodiePartitionMetadata partitionMetadata = new HoodiePartitionMetadata(storage, instantTime, + new StoragePath(config.getBasePath()), + FSUtils.constructAbsolutePath(config.getBasePath(), partitionPath), + hoodieTable.getPartitionMetafileFormat()); + partitionMetadata.trySave(); + + String newFileName = FSUtils.makeBaseFileName(instantTime, writeToken, fileId, hoodieTable.getBaseFileExtension()); + makeOldAndNewFilePaths(partitionPath, + latestValidFilePath.isPresent() ? latestValidFilePath.get() : null, newFileName); + + LOG.info("Merging data from file group {}, to a new base file {}", fileId, newFilePath); + // file name is same for all records, in this bunch + writeStatus.setFileId(fileId); + writeStatus.setPartitionPath(partitionPath); + writeStatus.getStat().setPartitionPath(partitionPath); + writeStatus.getStat().setFileId(fileId); + setWriteStatusPath(); } /** @@ -198,7 +177,6 @@ public void doMerge() { Option internalSchemaOption = SerDeHelper.fromJson(config.getInternalSchema()); TypedProperties props = TypedProperties.copy(config.getProps()); Stream logFiles = Stream.empty(); - Option> engineRecordIterator = Option.empty(); // For compaction. if (operation.isPresent()) { long maxMemoryPerCompaction = IOUtils.getMaxMemoryPerCompaction(taskContextSupplier, config); @@ -209,10 +187,9 @@ public void doMerge() { } // For generic merge. if (recordIterator.isPresent()) { - engineRecordIterator = Option.of(new MappingIterator<>(recordIterator.get(), HoodieRecord::getData)); usePosition = false; } - // Initializes file group reader + // Initializes file group reader. try (HoodieFileGroupReader fileGroupReader = HoodieFileGroupReader.newBuilder() .withReaderContext(readerContext) .withHoodieTableMetaClient(hoodieTable.getMetaClient()) @@ -220,14 +197,14 @@ public void doMerge() { .withPartitionPath(partitionPath) .withBaseFileOption(Option.ofNullable(baseFileToMerge)) .withLogFiles(logFiles) - .withDataSchema(writeSchemaWithMetaFields) + .withDataSchema(writeSchema) .withRequestedSchema(writeSchemaWithMetaFields) .withInternalSchema(internalSchemaOption) .withProps(props) .withShouldUseRecordPosition(usePosition) .withSortOutput(hoodieTable.requireSortedRecords()) .withFileGroupUpdateCallback(createCallbacks()) - .withRecordIterator(engineRecordIterator).build()) { + .withRecordIterator(recordIterator).build()) { // Reads the records from the file slice try (ClosableIterator> recordIterator = fileGroupReader.getClosableHoodieRecordIterator()) { while (recordIterator.hasNext()) { @@ -241,7 +218,7 @@ public void doMerge() { writeStatus.markFailure(record, failureEx, recordMetadata); continue; } - // Writes the record + // Writes the record. try { writeToFile(record.getKey(), record, writeSchemaWithMetaFields, config.getPayloadConfig().getProps(), preserveMetadata); diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/IOUtils.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/IOUtils.java index 29fe5246b613e..b136e2e773003 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/IOUtils.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/IOUtils.java @@ -119,6 +119,7 @@ public static Iterator> runMerge(HoodieMergeHandle "Error in finding the old file path at commit " + instantTime + " for fileId: " + fileId); } else { mergeHandle.doMerge(); + mergeHandle.close(); } // TODO(vc): This needs to be revisited diff --git a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java index afcf5ff8915d5..0e495928602f2 100644 --- a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java +++ b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/commit/BaseSparkCommitActionExecutor.java @@ -178,7 +178,7 @@ public HoodieWriteMetadata> execute(HoodieData> inputRecordsWithClusteringUpdate = clusteringHandleUpdate(inputRecords); LOG.info("Num spark partitions for inputRecords before triggering workload profile {}", inputRecordsWithClusteringUpdate.getNumPartitions()); @@ -391,7 +391,8 @@ public Iterator> handleUpdate(String partitionPath, String fil protected HoodieMergeHandle getUpdateHandle(String partitionPath, String fileId, Iterator> recordItr) { HoodieMergeHandle mergeHandle; if (config.getMergeHandleClassName().equals(FileGroupReaderBasedMergeHandle.class.getName())) { - HoodieReaderContext readerContext = table.getContext().getReaderContextFactory(table.getMetaClient()).getContext(); + HoodieReaderContext readerContext = + table.getContext().getReaderContextFactory(table.getMetaClient()).getContext(); mergeHandle = new FileGroupReaderBasedMergeHandle( config, instantTime, table, Option.of(recordItr), partitionPath, fileId, taskContextSupplier, keyGeneratorOpt, readerContext, instantTime, table.getConfig().getRecordMerger().getRecordType(), Option.empty()); diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecord.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecord.java index 21515392624aa..ea26967aa7e3a 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecord.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecord.java @@ -74,6 +74,20 @@ public static BufferedRecord forRecordWithContext(T record, Schema schema return new BufferedRecord<>(recordKey, orderingValue, record, schemaId, isDelete); } + /** + * When HoodieRecord is given, recordKey is definitely available. + */ + public static BufferedRecord forRecordWithContext(String recordKey, + T record, + Schema schema, + HoodieReaderContext readerContext, + Option orderingFieldName, + boolean isDelete) { + Integer schemaId = readerContext.encodeAvroSchema(schema); + Comparable orderingValue = readerContext.getOrderingValue(record, schema, orderingFieldName); + return new BufferedRecord<>(recordKey, orderingValue, record, schemaId, isDelete); + } + public static BufferedRecord forDeleteRecord(DeleteRecord deleteRecord, Comparable orderingValue) { return new BufferedRecord<>(deleteRecord.getRecordKey(), orderingValue, null, null, true); } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java index 1b529389d5073..517a645c88291 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java @@ -172,7 +172,7 @@ private FileGroupRecordBuffer getRecordBuffer(HoodieReaderContext readerCo HoodieReadStats readStats, boolean emitDelete, boolean sortOutput, - Option> inputRecordOpt) { + Option>> inputRecordOpt) { UpdateProcessor updateProcessor = UpdateProcessor.create(readStats, readerContext, emitDelete, fileGroupUpdateCallbacks); if (inputSplit.logFiles.isEmpty()) { if (inputRecordOpt.isPresent()) { @@ -204,7 +204,12 @@ private FileGroupRecordBuffer getRecordBuffer(HoodieReaderContext readerCo private void initRecordIterators() throws IOException { ClosableIterator iter = makeBaseFileIterator(); if (inputSplit.logFiles.isEmpty()) { - this.baseFileIterator = new CloseableMappingIterator<>(iter, readerContext::seal); + if (inputSplit.recordIterator.isPresent()) { + this.baseFileIterator = iter; + recordBuffer.setBaseFileIterator(baseFileIterator); + } else { + this.baseFileIterator = new CloseableMappingIterator<>(iter, readerContext::seal); + } } else { this.baseFileIterator = iter; scanLogFiles(); @@ -385,7 +390,8 @@ public ClosableIterator getClosableIterator() throws IOException { */ public ClosableIterator> getClosableHoodieRecordIterator() throws IOException { return new CloseableMappingIterator<>(getClosableIterator(), nextRecord -> { - BufferedRecord bufferedRecord = BufferedRecord.forRecordWithContext(nextRecord, readerContext.getSchemaHandler().getRequestedSchema(), readerContext, orderingFieldName, false); + BufferedRecord bufferedRecord = BufferedRecord.forRecordWithContext( + nextRecord, readerContext.getSchemaHandler().getRequestedSchema(), readerContext, orderingFieldName, false); return readerContext.constructHoodieRecord(bufferedRecord); }); } @@ -463,7 +469,7 @@ public static class Builder { private boolean sortOutput = false; private boolean enableOptimizedLogBlockScan = false; private List> fileGroupUpdateCallbacks = Collections.emptyList(); - private Option> recordIterator = Option.empty(); + private Option>> recordIterator = Option.empty(); public Builder withReaderContext(HoodieReaderContext readerContext) { this.readerContext = readerContext; @@ -569,7 +575,7 @@ public Builder withSortOutput(boolean sortOutput) { return this; } - public Builder withRecordIterator(Option> iterator) { + public Builder withRecordIterator(Option>> iterator) { this.recordIterator = iterator; return this; } @@ -606,14 +612,14 @@ private static class InputSplit { // Length of bytes to read from the base file private final long length; // For input record case. - private Option> recordIterator; + private Option>> recordIterator; InputSplit(Option baseFileOption, Stream logFiles, String partitionPath, long start, long length, - Option> recordIterator) { + Option>> recordIterator) { this(baseFileOption, logFiles, partitionPath, start, length); this.recordIterator = recordIterator; } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java index 4b4523450d2af..f2fc56c2a7c0f 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java @@ -23,20 +23,26 @@ import org.apache.hudi.common.config.TypedProperties; import org.apache.hudi.common.engine.HoodieReaderContext; import org.apache.hudi.common.model.DeleteRecord; +import org.apache.hudi.common.model.HoodieRecord; import org.apache.hudi.common.table.HoodieTableMetaClient; import org.apache.hudi.common.table.PartialUpdateMode; import org.apache.hudi.common.table.log.KeySpec; import org.apache.hudi.common.table.log.block.HoodieDataBlock; import org.apache.hudi.common.table.log.block.HoodieDeleteBlock; import org.apache.hudi.common.util.Option; +import org.apache.hudi.exception.HoodieException; import org.apache.hudi.exception.HoodieNotSupportedException; import java.io.IOException; import java.io.Serializable; import java.util.Iterator; +/** + * This is a special record buffer that caches input records, + * which is used to merge with records from base file. + */ public class InputBasedFileGroupRecordBuffer extends KeyBasedFileGroupRecordBuffer { - private final Iterator inputRecordIterator; + private final Iterator> inputRecordIterator; public InputBasedFileGroupRecordBuffer(HoodieReaderContext readerContext, HoodieTableMetaClient hoodieTableMetaClient, @@ -44,7 +50,7 @@ public InputBasedFileGroupRecordBuffer(HoodieReaderContext readerContext, PartialUpdateMode partialUpdateMode, TypedProperties props, Option orderingFieldName, - Iterator inputRecordIterator, + Iterator> inputRecordIterator, UpdateProcessor updateProcessor) { super(readerContext, hoodieTableMetaClient, recordMergeMode, partialUpdateMode, props, orderingFieldName, updateProcessor); this.inputRecordIterator = inputRecordIterator; @@ -59,15 +65,19 @@ private void populateRecordBuffer() { } while (inputRecordIterator.hasNext()) { - T engineRecord = inputRecordIterator.next(); - String recordKey = readerContext.getRecordKey(nextRecord, readerSchema); - boolean isDelete = - isBuiltInDeleteRecord(engineRecord) - || isCustomDeleteRecord(engineRecord) - || isDeleteHoodieOperation(engineRecord); - BufferedRecord bufferedRecord = BufferedRecord.forRecordWithContext( - engineRecord, readerSchema, readerContext, orderingFieldName, isDelete); - records.put(recordKey, bufferedRecord); + HoodieRecord hoodieRecord = inputRecordIterator.next(); + try { + BufferedRecord bufferedRecord = BufferedRecord.forRecordWithContext( + hoodieRecord.getRecordKey(), + hoodieRecord.getData(), + readerContext.getSchemaHandler().tableSchema, + readerContext, + orderingFieldName, + hoodieRecord.isDelete(readerContext.getSchemaHandler().tableSchema, props)); + records.put(hoodieRecord.getRecordKey(), bufferedRecord); + } catch (IOException e) { + throw new HoodieException("Failed to populate data into the record buffer", e); + } } } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java index 774ec6294e21a..88e37d7f2afd4 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java @@ -122,7 +122,7 @@ private void invokeCallbacks(CallbackInvoker invoker) { try { invoker.invoke(callback); } catch (Exception e) { - LOG.error(String.format("Callback %s failed: ", callback.getName()), e); + LOG.error("Callback {} failed: ", callback.getName(), e); } } } diff --git a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestPayloadDeprecationFlow.scala b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestPayloadDeprecationFlow.scala index 48685ca62714d..12ef0632729fb 100644 --- a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestPayloadDeprecationFlow.scala +++ b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestPayloadDeprecationFlow.scala @@ -19,114 +19,239 @@ package org.apache.hudi.functional -import org.apache.hudi.DataSourceWriteOptions +import org.apache.hudi.{DataSourceWriteOptions, DefaultSparkRecordMerger, OverwriteWithLatestSparkRecordMerger} import org.apache.hudi.DataSourceWriteOptions.{OPERATION, PRECOMBINE_FIELD, RECORDKEY_FIELD, TABLE_TYPE} -import org.apache.hudi.common.model.{AWSDmsAvroPayload, EventTimeAvroPayload, OverwriteNonDefaultsWithLatestAvroPayload, OverwriteWithLatestAvroPayload, PartialUpdateAvroPayload} -import org.apache.hudi.common.table.HoodieTableConfig +import org.apache.hudi.common.config.{HoodieMetadataConfig, HoodieStorageConfig, TypedProperties} +import org.apache.hudi.common.model._ +import org.apache.hudi.common.table.{HoodieTableConfig, HoodieTableMetaClient, HoodieTableVersion} import org.apache.hudi.config.{HoodieCompactionConfig, HoodieWriteConfig} +import org.apache.hudi.io.FileGroupReaderBasedMergeHandle +import org.apache.hudi.keygen.constant.KeyGeneratorOptions +import org.apache.hudi.table.upgrade.{SparkUpgradeDowngradeHelper, UpgradeDowngrade} import org.apache.hudi.testutils.SparkClientFunctionalTestHarness -import org.apache.spark.sql.SaveMode -import org.junit.jupiter.api.Assertions.assertTrue +import org.apache.spark.sql.{Row, SaveMode} +import org.apache.spark.sql.types._ +import org.junit.jupiter.api.Assertions.{assertEquals, assertTrue} +import org.junit.jupiter.api.Disabled import org.junit.jupiter.params.ParameterizedTest import org.junit.jupiter.params.provider.{Arguments, MethodSource} +import scala.jdk.CollectionConverters._ + +// TODO: After binary Avro HoodieRecord is supported, we enable this test. +@Disabled class TestPayloadDeprecationFlow extends SparkClientFunctionalTestHarness { + // Create a custom schema that includes the Op field for AWSDmsAvroPayload testing + private val CUSTOM_SCHEMA_WITH_OP = StructType(Seq( + StructField("timestamp", LongType, false), + StructField("_row_key", StringType, false), + StructField("rider", StringType, false), + StructField("driver", StringType, false), + StructField("fare", StructType(Seq( + StructField("amount", DoubleType, false), + StructField("currency", StringType, false) + )), false), + StructField("_hoodie_is_deleted", BooleanType, false), + StructField("Op", StringType, false) + )) + + /** + * Test if the payload based read have the same behavior for different table versions. + */ @ParameterizedTest - @MethodSource(Array("provideParams")) + @MethodSource(Array("provideParamsForPayloadBehavior")) def testMergerBuiltinPayload(tableType: String, - payloadClazz: String): Unit = { + payloadClazz: String, + tableVersion: String, + mergeMode: String): Unit = { + val mergers = List(classOf[DefaultSparkRecordMerger].getName, classOf[OverwriteWithLatestSparkRecordMerger].getName) + val mergerClasses = mergers.mkString(",") + val compactionEnabled = if (tableType.equals(HoodieTableType.MERGE_ON_READ.name)) "true" else "false" val opts: Map[String, String] = Map( + HoodieCompactionConfig.INLINE_COMPACT.key() -> compactionEnabled, + HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key() -> "parquet", + HoodieWriteConfig.MERGE_HANDLE_CLASS_NAME.key() -> classOf[FileGroupReaderBasedMergeHandle[_, _, _, _]].getName, + HoodieMetadataConfig.ENABLE.key() -> "true", HoodieWriteConfig.WRITE_PAYLOAD_CLASS_NAME.key() -> payloadClazz, - HoodieTableConfig.MERGE_PROPERTIES.key() -> - "hoodie.payload.delete.field=xp,hoodie.payload.delete.marker=d") - val columns = Seq("ts", "key", "rider", "driver", "fare", "Op") - - // 1. Add an insert. - val data = Seq( - (10, "1", "rider-A", "driver-A", 19.10, "i"), - (10, "2", "rider-B", "driver-B", 27.70, "i"), - (10, "3", "rider-C", "driver-C", 33.90, "i"), - (10, "4", "rider-D", "driver-D", 34.15, "i"), - (10, "5", "rider-E", "driver-E", 17.85, "i")) - val inserts = spark.createDataFrame(data).toDF(columns: _*) - inserts.write.format("hudi"). - option(RECORDKEY_FIELD.key(), "key"). - option(PRECOMBINE_FIELD.key(), "ts"). + HoodieWriteConfig.RECORD_MERGE_IMPL_CLASSES.key() -> mergerClasses, + KeyGeneratorOptions.PARTITIONPATH_FIELD_NAME.key() -> "_row_key", + HoodieTableConfig.RECORDKEY_FIELDS.key() -> "_row_key", + HoodieTableConfig.RECORD_MERGE_MODE.key() -> mergeMode) + + var metaClient: HoodieTableMetaClient = getHoodieMetaClient(storageConf(), basePath()) + new UpgradeDowngrade(metaClient, getWriteConfig(opts), context, SparkUpgradeDowngradeHelper.getInstance) + .run(HoodieTableVersion.SIX, null) + + // 1. Add initial inserts using Spark Rows + val initialRecords = List( + createTestRecord("1", "rider-A", "driver-A", 19.10, 10L, false), + createTestRecord("2", "rider-B", "driver-B", 27.70, 10L, false), + createTestRecord("3", "rider-C", "driver-C", 33.90, 10L, false), + createTestRecord("4", "rider-D", "driver-D", 34.15, 10L, false), + createTestRecord("5", "rider-E", "driver-E", 17.85, 10L, false) + ) + val initialDf = createDataFrameFromRows(initialRecords) + initialDf.write.format("hudi"). + option(RECORDKEY_FIELD.key(), "_row_key"). + option(PRECOMBINE_FIELD.key(), "timestamp"). option(TABLE_TYPE.key(), tableType). option(DataSourceWriteOptions.TABLE_NAME.key(), "test_table"). option(HoodieCompactionConfig.INLINE_COMPACT.key(), "false"). + option(HoodieWriteConfig.WRITE_TABLE_VERSION.key(), tableVersion). + option(HoodieTableConfig.INITIAL_VERSION.key(), tableVersion). options(opts). mode(SaveMode.Overwrite). save(basePath) - // 2. Add an update. - val firstUpdateData = Seq( - (11, "1", "rider-X", "driver-X", 19.10, "d"), - (11, "2", "rider-Y", "driver-Y", 27.70, "u")) - val firstUpdate = spark.createDataFrame(firstUpdateData).toDF(columns: _*) - firstUpdate.write.format("hudi"). + // Validate table version. + metaClient = HoodieTableMetaClient.reload(metaClient) + assertEquals( + Integer.valueOf(tableVersion), + metaClient.getTableConfig.getTableVersion.versionCode()) + // 2. Add first update using Spark Rows + val firstUpdateRecords = List( + createTestRecord("1", "rider-X", "driver-X", 19.10, 11L, true), // delete + createTestRecord("2", "rider-Y", "driver-Y", 27.70, 11L, false) // update + ) + val firstUpdateDf = createDataFrameFromRows(firstUpdateRecords) + firstUpdateDf.write.format("hudi"). + option(RECORDKEY_FIELD.key(), "_row_key"). + option(PRECOMBINE_FIELD.key(), "timestamp"). option(OPERATION.key(), "upsert"). option(HoodieCompactionConfig.INLINE_COMPACT.key(), "false"). + option(HoodieWriteConfig.WRITE_TABLE_VERSION.key(), tableVersion). options(opts). mode(SaveMode.Append). save(basePath) - // 3. Add an update. - val secondUpdateData = Seq( - (12, "3", "rider-CC", "driver-CC", 33.90, "i"), - (9, "4", "rider-DD", "driver-DD", 34.15, "i"), - (12, "5", "rider-EE", "driver-EE", 17.85, "i")) - val secondUpdate = spark.createDataFrame(secondUpdateData).toDF(columns: _*) - secondUpdate.write.format("hudi"). + // Validate table version. + metaClient = HoodieTableMetaClient.reload(metaClient); + assertEquals( + Integer.valueOf(tableVersion), + metaClient.getTableConfig.getTableVersion.versionCode()) + + val df1 = spark.read.format("hudi").options(opts).load(basePath) + val finalDf1 = df1.select("timestamp", "_row_key", "rider", "driver", "fare").sort("_row_key") + finalDf1.show(false) + + // 3. Add second update using Spark Rows + val secondUpdateRecords = List( + createTestRecord("3", "rider-CC", "driver-CC", 33.90, 12L, false), + createTestRecord("4", "rider-DD", "driver-DD", 34.15, 9L, false), + createTestRecord("5", "rider-EE", "driver-EE", 17.85, 12L, false) + ) + + val secondUpdateDf = createDataFrameFromRows(secondUpdateRecords) + secondUpdateDf.write.format("hudi"). option(OPERATION.key(), "upsert"). - option(HoodieCompactionConfig.INLINE_COMPACT.key(), "true"). - option(HoodieCompactionConfig.INLINE_COMPACT_NUM_DELTA_COMMITS.key(), "1"). + option(HoodieCompactionConfig.INLINE_COMPACT.key(), "false"). + option(HoodieWriteConfig.WRITE_TABLE_VERSION.key(), tableVersion). options(opts). mode(SaveMode.Append). save(basePath) - // 4. Validate. + // Validate table version. + metaClient = HoodieTableMetaClient.reload(metaClient); + assertEquals( + Integer.valueOf(tableVersion), + metaClient.getTableConfig.getTableVersion.versionCode()) + + // 4. Validate final results using Spark Rows val df = spark.read.format("hudi").options(opts).load(basePath) - val finalDf = df.select("ts", "key", "rider", "driver", "fare", "Op").sort("key") + val finalDf = df.select("timestamp", "_row_key", "rider", "driver", "fare").sort("_row_key") - val expectedData = if (!payloadClazz.equals(classOf[AWSDmsAvroPayload].getName)) { + val expectedRecords = if (!payloadClazz.equals(classOf[AWSDmsAvroPayload].getName)) { if (payloadClazz.equals(classOf[PartialUpdateAvroPayload].getName) || payloadClazz.equals(classOf[EventTimeAvroPayload].getName)) { - Seq( - (11, "1", "rider-X", "driver-X", 19.10, "d"), - (11, "2", "rider-Y", "driver-Y", 27.70, "u"), - (12, "3", "rider-CC", "driver-CC", 33.90, "i"), - (10, "4", "rider-D", "driver-D", 34.15, "i"), - (12, "5", "rider-EE", "driver-EE", 17.85, "i")) + List( + createTestRecord("2", "rider-Y", "driver-Y", 27.70, 11L, false), + createTestRecord("3", "rider-CC", "driver-CC", 33.90, 12L, false), + createTestRecord("4", "rider-D", "driver-D", 34.15, 10L, false), + createTestRecord("5", "rider-EE", "driver-EE", 17.85, 12L, false) + ) } else { - Seq( - (11, "1", "rider-X", "driver-X", 19.10, "d"), - (11, "2", "rider-Y", "driver-Y", 27.70, "u"), - (12, "3", "rider-CC", "driver-CC", 33.90, "i"), - (9, "4", "rider-DD", "driver-DD", 34.15, "i"), - (12, "5", "rider-EE", "driver-EE", 17.85, "i")) + List( + createTestRecord("2", "rider-Y", "driver-Y", 27.70, 11L, false), + createTestRecord("3", "rider-CC", "driver-CC", 33.90, 12L, false), + createTestRecord("4", "rider-DD", "driver-DD", 34.15, 9L, false), + createTestRecord("5", "rider-EE", "driver-EE", 17.85, 12L, false) + ) } } else { - Seq( - (11, "2", "rider-Y", "driver-Y", 27.70, "u"), - (12, "3", "rider-CC", "driver-CC", 33.90, "i"), - (9, "4", "rider-DD", "driver-DD", 34.15, "i"), - (12, "5", "rider-EE", "driver-EE", 17.85, "i")) + List( + createTestRecord("2", "rider-Y", "driver-Y", 27.70, 11L, false), + createTestRecord("3", "rider-CC", "driver-CC", 33.90, 12L, false), + createTestRecord("4", "rider-DD", "driver-DD", 34.15, 9L, false), + createTestRecord("5", "rider-EE", "driver-EE", 17.85, 12L, false) + ) } - val expectedDf = spark.createDataFrame( - spark.sparkContext.parallelize(expectedData)).toDF(columns: _*).sort("key") + val expectedDf = createDataFrameFromRows(expectedRecords) + .select("timestamp", "_row_key", "rider", "driver", "fare").sort("_row_key") + + expectedDf.show(false) + finalDf.show(false) assertTrue( expectedDf.except(finalDf).isEmpty && finalDf.except(expectedDf).isEmpty) } + + /** + * Helper method to create test Spark Row + */ + private def createTestRecord(key: String, + rider: String, + driver: String, + fare: Double, + timestamp: Long, + isDelete: Boolean): Row = { + // Create nested fare row + val fareRow = Row(fare, "USD") + // Set Op field for operation indication + val opValue = if (isDelete) "D" else "I" + // Create the main row + Row(timestamp, key, rider, driver, fareRow, isDelete, opValue) + } + + /** + * Helper method to convert Spark Row list to Spark DataFrame + */ + private def createDataFrameFromRows(records: List[Row]): org.apache.spark.sql.DataFrame = { + spark.createDataFrame(records.asJava, CUSTOM_SCHEMA_WITH_OP) + } + + def getWriteConfig(hudiOpts: Map[String, String]): HoodieWriteConfig = { + val props = TypedProperties.fromMap(hudiOpts.asJava) + HoodieWriteConfig.newBuilder() + .withProps(props) + .withPath(basePath()) + .build() + } } -// TODO: Add COPY_ON_WRITE table type tests when write path is updated accordingly. object TestPayloadDeprecationFlow { - def provideParams(): java.util.List[Arguments] = { + def provideParamsForPayloadBehavior(): java.util.List[Arguments] = { java.util.Arrays.asList( - Arguments.of("MERGE_ON_READ", classOf[OverwriteWithLatestAvroPayload].getName), - Arguments.of("MERGE_ON_READ", classOf[OverwriteNonDefaultsWithLatestAvroPayload].getName), - Arguments.of("MERGE_ON_READ", classOf[PartialUpdateAvroPayload].getName), - Arguments.of("MERGE_ON_READ", classOf[EventTimeAvroPayload].getName), - Arguments.of("MERGE_ON_READ", classOf[AWSDmsAvroPayload].getName) + // For COW merge. + Arguments.of("COPY_ON_WRITE", classOf[OverwriteWithLatestAvroPayload].getName, "6", "COMMIT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", classOf[OverwriteNonDefaultsWithLatestAvroPayload].getName, "6", "COMMIT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", classOf[PartialUpdateAvroPayload].getName, "6", "EVENT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", classOf[EventTimeAvroPayload].getName, "6", "EVENT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", classOf[AWSDmsAvroPayload].getName, "6", "COMMIT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", classOf[OverwriteWithLatestAvroPayload].getName, "8", "COMMIT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", classOf[OverwriteNonDefaultsWithLatestAvroPayload].getName, "8", "COMMIT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", classOf[PartialUpdateAvroPayload].getName, "8", "EVENT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", classOf[EventTimeAvroPayload].getName, "8", "EVENT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", classOf[AWSDmsAvroPayload].getName, "8", "COMMIT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", classOf[OverwriteWithLatestAvroPayload].getName, "9", "COMMIT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", classOf[OverwriteNonDefaultsWithLatestAvroPayload].getName, "9", "COMMIT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", classOf[PartialUpdateAvroPayload].getName, "9", "EVENT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", classOf[EventTimeAvroPayload].getName, "9", "EVENT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", classOf[AWSDmsAvroPayload].getName, "9", "COMMIT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", "", "6", "COMMIT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", "", "8", "COMMIT_TIME_ORDERING"), + Arguments.of("COPY_ON_WRITE", "", "9", "EVENT_TIME_ORDERING"), + // For compaction. + Arguments.of("MERGE_ON_WRITE", classOf[OverwriteWithLatestAvroPayload].getName, "6", "COMMIT_TIME_ORDERING"), + Arguments.of("MERGE_ON_WRITE", classOf[PartialUpdateAvroPayload].getName, "9", "EVENT_TIME_ORDERING"), + Arguments.of("MERGE_ON_WRITE", classOf[EventTimeAvroPayload].getName, "9", "EVENT_TIME_ORDERING"), + Arguments.of("MERGE_ON_WRITE", classOf[AWSDmsAvroPayload].getName, "9", "COMMIT_TIME_ORDERING") ) } } From 79e9f89d1fad549b6f87355c31e3bd4c55835bcc Mon Sep 17 00:00:00 2001 From: danny0405 Date: Fri, 25 Jul 2025 11:46:50 +0800 Subject: [PATCH 10/10] Introduce composite update callback, also fix the ordering value discrepancies in BufferedRecord --- .../io/FileGroupReaderBasedMergeHandle.java | 35 +++++++++++----- .../hudi/io/HoodieMergeHandleFactory.java | 2 + .../table/read/BaseFileUpdateCallback.java | 5 --- .../common/table/read/BufferedRecord.java | 23 ++++------ .../read/BufferedRecordMergerFactory.java | 21 ++++++---- .../table/read/HoodieFileGroupReader.java | 18 ++++---- .../read/InputBasedFileGroupRecordBuffer.java | 20 ++++----- .../common/table/read/UpdateProcessor.java | 42 +++++-------------- .../metadata/HoodieTableMetadataUtil.java | 2 +- .../TestKeyBasedFileGroupRecordBuffer.java | 2 +- ...stSortedKeyBasedFileGroupRecordBuffer.java | 2 +- .../hudi/cdc/CDCFileGroupIterator.scala | 2 +- 12 files changed, 77 insertions(+), 97 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 6cc3c0ff247ce..68c493e4b3d8c 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 @@ -203,7 +203,7 @@ public void doMerge() { .withProps(props) .withShouldUseRecordPosition(usePosition) .withSortOutput(hoodieTable.requireSortedRecords()) - .withFileGroupUpdateCallback(createCallbacks()) + .withFileGroupUpdateCallback(createCallback()) .withRecordIterator(recordIterator).build()) { // Reads the records from the file slice try (ClosableIterator> recordIterator = fileGroupReader.getClosableHoodieRecordIterator()) { @@ -242,7 +242,7 @@ public void doMerge() { } } - private List> createCallbacks() { + private Option> createCallback() { List> callbacks = new ArrayList<>(); // Handle CDC workflow. if (cdcLogger.isPresent()) { @@ -261,7 +261,9 @@ private List> createCallbacks() { config )); } - return callbacks; + return callbacks.isEmpty() + ? Option.empty() + : callbacks.size() == 1 ? Option.of(callbacks.get(0)) : Option.of(new CompositeCallback<>(callbacks)); } @Override @@ -331,11 +333,6 @@ public void onDelete(String recordKey, T previousRecord) { } - @Override - public String getName() { - return "CDCCallback"; - } - private GenericRecord convertOutput(T record) { T convertedRecord = outputConverter.get().map(converter -> record == null ? null : converter.apply(record)).orElse(record); return convertedRecord == null ? null : readerContext.convertToAvroRecord(convertedRecord, requestedSchema.get()); @@ -425,10 +422,28 @@ public void onDelete(String recordKey, T previousRecord) { keyGeneratorOpt, config); } + } + + private static class CompositeCallback implements BaseFileUpdateCallback { + private final List> callbacks; + + public CompositeCallback(List> callbacks) { + this.callbacks = callbacks; + } + + @Override + public void onUpdate(String recordKey, T previousRecord, T mergedRecord) { + this.callbacks.forEach(callback -> callback.onUpdate(recordKey, previousRecord, mergedRecord)); + } + + @Override + public void onInsert(String recordKey, T newRecord) { + this.callbacks.forEach(callback -> callback.onInsert(recordKey, newRecord)); + } @Override - public String getName() { - return "SecondaryIndexCallback"; + public void onDelete(String recordKey, T previousRecord) { + this.callbacks.forEach(callback -> callback.onDelete(recordKey, previousRecord)); } } } 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 94fcd10313925..a0fe1291a9322 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 @@ -124,10 +124,12 @@ public static HoodieMergeHandle create( String mergeHandleClass = config.getCompactionMergeHandleClassName(); 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); + Class[] constructorParamTypes = new Class[] { HoodieWriteConfig.class, String.class, HoodieTable.class, CompactionOperation.class, TaskContextSupplier.class, HoodieReaderContext.class, String.class, HoodieRecord.HoodieRecordType.class }; + return instantiateMergeHandle( isFallbackEnabled, mergeHandleClass, COMPACT_MERGE_HANDLE_CLASS_NAME.defaultValue(), logContext, constructorParamTypes, config, instantTime, hoodieTable, operation, taskContextSupplier, readerContext, maxInstantTime, recordType); diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/BaseFileUpdateCallback.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/BaseFileUpdateCallback.java index 2bbb4f62a8044..6685cd6cef459 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/BaseFileUpdateCallback.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/BaseFileUpdateCallback.java @@ -44,9 +44,4 @@ public interface BaseFileUpdateCallback { * @param previousRecord the record in the base file before deletion */ void onDelete(String recordKey, T previousRecord); - - /** - * Callback method to return the name of the callback. - */ - String getName(); } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecord.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecord.java index ea26967aa7e3a..dfed419b07350 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecord.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecord.java @@ -54,7 +54,11 @@ public BufferedRecord(String recordKey, Comparable orderingValue, T record, Inte this.isDelete = isDelete; } - public static BufferedRecord forRecordWithContext(HoodieRecord record, Schema schema, HoodieReaderContext readerContext, Properties props) { + public static BufferedRecord forRecordWithContext(HoodieRecord record, + Schema schema, + HoodieReaderContext readerContext, + Option orderingFieldName, + Properties props) { HoodieKey hoodieKey = record.getKey(); String recordKey = hoodieKey == null ? readerContext.getRecordKey(record.getData(), schema) : hoodieKey.getRecordKey(); Integer schemaId = readerContext.encodeAvroSchema(schema); @@ -64,7 +68,8 @@ public static BufferedRecord forRecordWithContext(HoodieRecord record, } catch (IOException e) { throw new HoodieException("Failed to get isDelete from record.", e); } - return new BufferedRecord<>(recordKey, record.getOrderingValue(schema, props), record.getData(), schemaId, isDelete); + T row = record.getData(); + return new BufferedRecord<>(recordKey, readerContext.getOrderingValue(row, schema, orderingFieldName), row, schemaId, isDelete); } public static BufferedRecord forRecordWithContext(T record, Schema schema, HoodieReaderContext readerContext, Option orderingFieldName, boolean isDelete) { @@ -74,20 +79,6 @@ public static BufferedRecord forRecordWithContext(T record, Schema schema return new BufferedRecord<>(recordKey, orderingValue, record, schemaId, isDelete); } - /** - * When HoodieRecord is given, recordKey is definitely available. - */ - public static BufferedRecord forRecordWithContext(String recordKey, - T record, - Schema schema, - HoodieReaderContext readerContext, - Option orderingFieldName, - boolean isDelete) { - Integer schemaId = readerContext.encodeAvroSchema(schema); - Comparable orderingValue = readerContext.getOrderingValue(record, schema, orderingFieldName); - return new BufferedRecord<>(recordKey, orderingValue, record, schemaId, isDelete); - } - public static BufferedRecord forDeleteRecord(DeleteRecord deleteRecord, Comparable orderingValue) { return new BufferedRecord<>(deleteRecord.getRecordKey(), orderingValue, null, null, true); } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecordMergerFactory.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecordMergerFactory.java index 5a821fe8338a3..d986bb64ead23 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecordMergerFactory.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/BufferedRecordMergerFactory.java @@ -69,7 +69,7 @@ public static BufferedRecordMerger create(HoodieReaderContext readerCo if (enablePartialMerging) { BufferedRecordMerger deleteRecordMerger = create( readerContext, recordMergeMode, false, recordMerger, orderingFieldName, payloadClass, readerSchema, props, partialUpdateMode); - return new PartialUpdateBufferedRecordMerger<>(readerContext, recordMerger, deleteRecordMerger, readerSchema, props); + return new PartialUpdateBufferedRecordMerger<>(readerContext, recordMerger, deleteRecordMerger, readerSchema, orderingFieldName, props); } switch (recordMergeMode) { @@ -88,7 +88,7 @@ public static BufferedRecordMerger create(HoodieReaderContext readerCo return new CustomPayloadBufferedRecordMerger<>( readerContext, recordMerger, orderingFieldName, payloadClass.get(), readerSchema, props); } else { - return new CustomBufferedRecordMerger<>(readerContext, recordMerger, readerSchema, props); + return new CustomBufferedRecordMerger<>(readerContext, recordMerger, readerSchema, orderingFieldName, props); } } } @@ -262,6 +262,7 @@ private static class PartialUpdateBufferedRecordMerger implements BufferedRec private final Option recordMerger; private final BufferedRecordMerger deleteRecordMerger; private final Schema readerSchema; + private final Option orderingFieldName; private final TypedProperties props; public PartialUpdateBufferedRecordMerger( @@ -269,10 +270,12 @@ public PartialUpdateBufferedRecordMerger( Option recordMerger, BufferedRecordMerger deleteRecordMerger, Schema readerSchema, + Option orderingFieldName, TypedProperties props) { this.readerContext = readerContext; this.recordMerger = recordMerger; this.deleteRecordMerger = deleteRecordMerger; + this.orderingFieldName = orderingFieldName; this.readerSchema = readerSchema; this.props = props; } @@ -300,7 +303,7 @@ public Option> deltaMerge(BufferedRecord newRecord, Buffere // If pre-combine returns existing record, no need to update it if (combinedRecord.getData() != existingRecord.getRecord()) { - return Option.of(BufferedRecord.forRecordWithContext(combinedRecord, combinedRecordAndSchema.getRight(), readerContext, props)); + return Option.of(BufferedRecord.forRecordWithContext(combinedRecord, combinedRecordAndSchema.getRight(), readerContext, orderingFieldName, props)); } return Option.empty(); } @@ -341,8 +344,9 @@ public CustomBufferedRecordMerger( HoodieReaderContext readerContext, Option recordMerger, Schema readerSchema, + Option orderingFieldName, TypedProperties props) { - super(readerContext, recordMerger, readerSchema, props); + super(readerContext, recordMerger, readerSchema, orderingFieldName, props); } @Override @@ -364,7 +368,7 @@ public Option> deltaMergeNonDeleteRecord(BufferedRecord new // If pre-combine returns existing record, no need to update it if (combinedRecord.getData() != existingRecord.getRecord()) { - return Option.of(BufferedRecord.forRecordWithContext(combinedRecord, combinedRecordAndSchema.getRight(), readerContext, props)); + return Option.of(BufferedRecord.forRecordWithContext(combinedRecord, combinedRecordAndSchema.getRight(), readerContext, orderingFieldName, props)); } return Option.empty(); } @@ -391,7 +395,6 @@ public Pair mergeNonDeleteRecord(BufferedRecord olderRecord, Buff * based on {@code CUSTOM} merge mode and a given record payload class. */ private static class CustomPayloadBufferedRecordMerger extends BaseCustomMerger { - private final Option orderingFieldName; private final String payloadClass; public CustomPayloadBufferedRecordMerger( @@ -401,8 +404,7 @@ public CustomPayloadBufferedRecordMerger( String payloadClass, Schema readerSchema, TypedProperties props) { - super(readerContext, recordMerger, readerSchema, props); - this.orderingFieldName = orderingFieldName; + super(readerContext, recordMerger, readerSchema, orderingFieldName, props); this.payloadClass = payloadClass; } @@ -477,16 +479,19 @@ private abstract static class BaseCustomMerger implements BufferedRecordMerge protected final HoodieReaderContext readerContext; protected final Option recordMerger; protected final Schema readerSchema; + protected final Option orderingFieldName; protected final TypedProperties props; public BaseCustomMerger( HoodieReaderContext readerContext, Option recordMerger, Schema readerSchema, + Option orderingFieldName, TypedProperties props) { this.readerContext = readerContext; this.recordMerger = recordMerger; this.readerSchema = readerSchema; + this.orderingFieldName = orderingFieldName; this.props = props; } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java index 517a645c88291..fc1f027f0ab9c 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java @@ -92,7 +92,7 @@ public final class HoodieFileGroupReader implements Closeable { // considers the log records which are inflight. private final boolean allowInflightInstants; // Callback to run custom logic on updates to the base files for the file group - private final List> fileGroupUpdateCallbacks; + private final Option> fileGroupUpdateCallback; private final boolean enableOptimizedLogBlockScan; // The list of instant times read from the log blocks, this value is used by the log-compaction to allow optimized log-block scans private List validBlockInstants = Collections.emptyList(); @@ -110,16 +110,16 @@ public HoodieFileGroupReader(HoodieReaderContext readerContext, HoodieStorage long start, long length, boolean shouldUseRecordPosition) { this(readerContext, storage, tablePath, latestCommitTime, dataSchema, requestedSchema, internalSchemaOpt, hoodieTableMetaClient, props, shouldUseRecordPosition, false, false, false, - InputSplit.fromFileSlice(fileSlice, start, length), Collections.emptyList(), false); + InputSplit.fromFileSlice(fileSlice, start, length), Option.empty(), false); } private HoodieFileGroupReader(HoodieReaderContext readerContext, HoodieStorage storage, String tablePath, String latestCommitTime, Schema dataSchema, Schema requestedSchema, Option internalSchemaOpt, HoodieTableMetaClient hoodieTableMetaClient, TypedProperties props, boolean shouldUseRecordPosition, boolean allowInflightInstants, boolean emitDelete, boolean sortOutput, - InputSplit inputSplit, List> updateCallback, boolean enableOptimizedLogBlockScan) { + InputSplit inputSplit, Option> updateCallback, boolean enableOptimizedLogBlockScan) { this.readerContext = readerContext; - this.fileGroupUpdateCallbacks = updateCallback; + this.fileGroupUpdateCallback = updateCallback; this.metaClient = hoodieTableMetaClient; this.storage = storage; this.enableOptimizedLogBlockScan = enableOptimizedLogBlockScan; @@ -173,7 +173,7 @@ private FileGroupRecordBuffer getRecordBuffer(HoodieReaderContext readerCo boolean emitDelete, boolean sortOutput, Option>> inputRecordOpt) { - UpdateProcessor updateProcessor = UpdateProcessor.create(readStats, readerContext, emitDelete, fileGroupUpdateCallbacks); + UpdateProcessor updateProcessor = UpdateProcessor.create(readStats, readerContext, emitDelete, fileGroupUpdateCallback); if (inputSplit.logFiles.isEmpty()) { if (inputRecordOpt.isPresent()) { return new InputBasedFileGroupRecordBuffer<>( @@ -468,7 +468,7 @@ public static class Builder { private boolean emitDelete; private boolean sortOutput = false; private boolean enableOptimizedLogBlockScan = false; - private List> fileGroupUpdateCallbacks = Collections.emptyList(); + private Option> fileGroupUpdateCallback = Option.empty(); private Option>> recordIterator = Option.empty(); public Builder withReaderContext(HoodieReaderContext readerContext) { @@ -554,8 +554,8 @@ public Builder withEmitDelete(boolean emitDelete) { return this; } - public Builder withFileGroupUpdateCallback(List> fileGroupUpdateCallbacks) { - this.fileGroupUpdateCallbacks = fileGroupUpdateCallbacks; + public Builder withFileGroupUpdateCallback(Option> fileGroupUpdateCallbacks) { + this.fileGroupUpdateCallback = fileGroupUpdateCallbacks; return this; } @@ -599,7 +599,7 @@ public HoodieFileGroupReader build() { InputSplit inputSplit = new InputSplit<>(baseFileOption, logFiles, partitionPath, start, length, recordIterator); return new HoodieFileGroupReader<>( readerContext, storage, tablePath, latestCommitTime, dataSchema, requestedSchema, internalSchemaOpt, hoodieTableMetaClient, - props, shouldUseRecordPosition, allowInflightInstants, emitDelete, sortOutput, inputSplit, fileGroupUpdateCallbacks, enableOptimizedLogBlockScan); + props, shouldUseRecordPosition, allowInflightInstants, emitDelete, sortOutput, inputSplit, fileGroupUpdateCallback, enableOptimizedLogBlockScan); } } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java index f2fc56c2a7c0f..34e0ff5b7724d 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/InputBasedFileGroupRecordBuffer.java @@ -30,7 +30,6 @@ import org.apache.hudi.common.table.log.block.HoodieDataBlock; import org.apache.hudi.common.table.log.block.HoodieDeleteBlock; import org.apache.hudi.common.util.Option; -import org.apache.hudi.exception.HoodieException; import org.apache.hudi.exception.HoodieNotSupportedException; import java.io.IOException; @@ -66,18 +65,13 @@ private void populateRecordBuffer() { while (inputRecordIterator.hasNext()) { HoodieRecord hoodieRecord = inputRecordIterator.next(); - try { - BufferedRecord bufferedRecord = BufferedRecord.forRecordWithContext( - hoodieRecord.getRecordKey(), - hoodieRecord.getData(), - readerContext.getSchemaHandler().tableSchema, - readerContext, - orderingFieldName, - hoodieRecord.isDelete(readerContext.getSchemaHandler().tableSchema, props)); - records.put(hoodieRecord.getRecordKey(), bufferedRecord); - } catch (IOException e) { - throw new HoodieException("Failed to populate data into the record buffer", e); - } + BufferedRecord bufferedRecord = BufferedRecord.forRecordWithContext( + hoodieRecord, + readerContext.getSchemaHandler().tableSchema, + readerContext, + orderingFieldName, + props); + records.put(hoodieRecord.getRecordKey(), bufferedRecord); } } diff --git a/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java b/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java index 88e37d7f2afd4..0173dc31b06e1 100644 --- a/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java +++ b/hudi-common/src/main/java/org/apache/hudi/common/table/read/UpdateProcessor.java @@ -19,11 +19,7 @@ package org.apache.hudi.common.table.read; import org.apache.hudi.common.engine.HoodieReaderContext; - -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import java.util.List; +import org.apache.hudi.common.util.Option; /** * Interface used within the {@link HoodieFileGroupReader} for processing updates to records in Merge-on-Read tables. @@ -42,10 +38,10 @@ public interface UpdateProcessor { T processUpdate(String recordKey, T previousRecord, T mergedRecord, boolean isDelete); static UpdateProcessor create(HoodieReadStats readStats, HoodieReaderContext readerContext, - boolean emitDeletes, List> updateCallbacks) { + boolean emitDeletes, Option> updateCallback) { UpdateProcessor handler = new StandardUpdateProcessor<>(readStats, readerContext, emitDeletes); - if (!updateCallbacks.isEmpty()) { - return new CallbackProcessor<>(updateCallbacks, handler, readerContext); + if (updateCallback.isPresent()) { + return new CallbackProcessor<>(updateCallback.get(), handler); } return handler; } @@ -91,14 +87,11 @@ public T processUpdate(String recordKey, T previousRecord, T mergedRecord, boole * @param the engine-specific record type */ class CallbackProcessor implements UpdateProcessor { - private static final Logger LOG = LoggerFactory.getLogger(CallbackProcessor.class); - private final List> callbacks; + private final BaseFileUpdateCallback callback; private final UpdateProcessor delegate; - public CallbackProcessor(List> callbacks, - UpdateProcessor delegate, - HoodieReaderContext readerContext) { - this.callbacks = callbacks; + public CallbackProcessor(BaseFileUpdateCallback callback, UpdateProcessor delegate) { + this.callback = callback; this.delegate = delegate; } @@ -107,29 +100,14 @@ public T processUpdate(String recordKey, T previousRecord, T currentRecord, bool T result = delegate.processUpdate(recordKey, previousRecord, currentRecord, isDelete); if (isDelete) { - invokeCallbacks(callback -> callback.onDelete(recordKey, previousRecord)); + callback.onDelete(recordKey, previousRecord); } else if (previousRecord != null && previousRecord != currentRecord) { - invokeCallbacks(callback -> callback.onUpdate(recordKey, previousRecord, currentRecord)); + callback.onUpdate(recordKey, previousRecord, currentRecord); } else { - invokeCallbacks(callback -> callback.onInsert(recordKey, currentRecord)); + callback.onInsert(recordKey, currentRecord); } return result; } - - private void invokeCallbacks(CallbackInvoker invoker) { - for (BaseFileUpdateCallback callback : callbacks) { - try { - invoker.invoke(callback); - } catch (Exception e) { - LOG.error("Callback {} failed: ", callback.getName(), e); - } - } - } - - @FunctionalInterface - private interface CallbackInvoker { - void invoke(BaseFileUpdateCallback callback) throws Exception; - } } } diff --git a/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java b/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java index 1b156b3cdd4a2..e8429383179a7 100644 --- a/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java +++ b/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java @@ -1034,7 +1034,7 @@ private static ClosableIterator> getLogRecords(List recordBuffer = new KeyBasedFileGroupRecordBuffer<>(readerContext, datasetMetaClient, readerContext.getMergeMode(), PartialUpdateMode.NONE, properties, Option.ofNullable(tableConfig.getPreCombineField()), - UpdateProcessor.create(readStats, readerContext, true, Collections.emptyList())); + UpdateProcessor.create(readStats, readerContext, true, Option.empty())); // CRITICAL: Ensure allowInflightInstants is set to true try (HoodieMergedLogRecordReader mergedLogRecordReader = HoodieMergedLogRecordReader.newBuilder() diff --git a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestKeyBasedFileGroupRecordBuffer.java b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestKeyBasedFileGroupRecordBuffer.java index 8d988fd7ad9eb..dd2a281b295d1 100644 --- a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestKeyBasedFileGroupRecordBuffer.java +++ b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestKeyBasedFileGroupRecordBuffer.java @@ -270,7 +270,7 @@ private static KeyBasedFileGroupRecordBuffer buildKeyBasedFileGro HoodieTableMetaClient mockMetaClient = mock(HoodieTableMetaClient.class, RETURNS_DEEP_STUBS); when(mockMetaClient.getTableConfig()).thenReturn(tableConfig); TypedProperties props = new TypedProperties(); - UpdateProcessor updateProcessor = UpdateProcessor.create(readStats, readerContext, false, Collections.emptyList()); + UpdateProcessor updateProcessor = UpdateProcessor.create(readStats, readerContext, false, Option.empty()); return new KeyBasedFileGroupRecordBuffer<>( readerContext, mockMetaClient, recordMergeMode, PartialUpdateMode.NONE, props, orderingFieldName, updateProcessor); } diff --git a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestSortedKeyBasedFileGroupRecordBuffer.java b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestSortedKeyBasedFileGroupRecordBuffer.java index b8268f2d2e518..439bcf8f920e1 100644 --- a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestSortedKeyBasedFileGroupRecordBuffer.java +++ b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestSortedKeyBasedFileGroupRecordBuffer.java @@ -127,7 +127,7 @@ private SortedKeyBasedFileGroupRecordBuffer buildSortedKeyBasedFileG RecordMergeMode recordMergeMode = RecordMergeMode.COMMIT_TIME_ORDERING; PartialUpdateMode partialUpdateMode = PartialUpdateMode.NONE; TypedProperties props = new TypedProperties(); - UpdateProcessor updateProcessor = UpdateProcessor.create(readStats, mockReaderContext, false, Collections.emptyList()); + UpdateProcessor updateProcessor = UpdateProcessor.create(readStats, mockReaderContext, false, Option.empty()); return new SortedKeyBasedFileGroupRecordBuffer<>( mockReaderContext, mockMetaClient, recordMergeMode, partialUpdateMode, props, Option.empty(), updateProcessor); } diff --git a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/cdc/CDCFileGroupIterator.scala b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/cdc/CDCFileGroupIterator.scala index 34f422f7e9a96..855a93c4dc8e0 100644 --- a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/cdc/CDCFileGroupIterator.scala +++ b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/cdc/CDCFileGroupIterator.scala @@ -503,7 +503,7 @@ class CDCFileGroupIterator(split: HoodieCDCFileGroupSplit, val recordBuffer = new KeyBasedFileGroupRecordBuffer[InternalRow](readerContext, metaClient, readerContext.getMergeMode, metaClient.getTableConfig.getPartialUpdateMode, readerProperties, Option.ofNullable(metaClient.getTableConfig.getPreCombineField), - UpdateProcessor.create(stats, readerContext, true, Collections.emptyList())) + UpdateProcessor.create(stats, readerContext, true, Option.empty())) HoodieMergedLogRecordReader.newBuilder[InternalRow] .withStorage(metaClient.getStorage)