This is an automated email from the ASF dual-hosted git repository. voonhous pushed a commit to branch release-1.2.1 in repository https://gitbox.apache.org/repos/asf/hudi.git
commit e3e003b6d1dd008c7afd2cfe04c8f34adad6ef74 Author: Shuo Cheng <[email protected]> AuthorDate: Fri Jun 5 12:20:03 2026 +0800 fix(flink): Use the execution mode without rocksdb cache by default … (#18894) (cherry picked from commit 6f81c51fae1d964989edfd634c69720b4569a5e6) --- .../metadata/HoodieBackedTableMetadataWriter.java | 2 +- .../apache/hudi/configuration/OptionsResolver.java | 41 +++++++++++ .../sink/partitioner/BucketAssignFunction.java | 8 +-- .../partitioner/GlobalRecordIndexPartitioner.java | 7 ++ .../partitioner/MinibatchBucketAssignFunction.java | 8 +-- .../sink/partitioner/RecordIndexPartitioner.java | 27 +------ .../index/GlobalRecordLevelIndexBackend.java | 45 ++++++++---- .../java/org/apache/hudi/sink/utils/Pipelines.java | 6 +- .../org/apache/hudi/table/HoodieTableFactory.java | 4 -- .../org/apache/hudi/util/FlinkWriteClients.java | 4 +- .../hudi/configuration/TestOptionsResolver.java | 34 +++++++++ .../TestMinibatchBucketAssignFunction.java | 82 +++++++++++++++++++++- .../partitioner/TestRecordIndexPartitioner.java | 1 + .../index/TestGlobalRecordLevelIndexBackend.java | 4 +- .../apache/hudi/table/ITTestHoodieDataSource.java | 20 +++++- .../apache/hudi/table/TestHoodieTableFactory.java | 17 +++++ 16 files changed, 250 insertions(+), 60 deletions(-) diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java index 85701fcb715f..546ce10bdfbe 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java @@ -172,7 +172,7 @@ public abstract class HoodieBackedTableMetadataWriter<I, O> implements HoodieTab // Average size of a record saved within the record index. // Record index has a fixed size schema. This has been calculated based on experiments with default settings // for block size (1MB), compression (GZ) and disabling the hudi metadata fields. - private static final int RECORD_INDEX_AVERAGE_RECORD_SIZE = 48; + public static final int RECORD_INDEX_AVERAGE_RECORD_SIZE = 48; private transient BaseHoodieWriteClient<?, I, ?, O> writeClient; protected HoodieWriteConfig metadataWriteConfig; diff --git a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/configuration/OptionsResolver.java b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/configuration/OptionsResolver.java index cc32351e62cf..42cdad7032c0 100644 --- a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/configuration/OptionsResolver.java +++ b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/configuration/OptionsResolver.java @@ -42,6 +42,8 @@ import org.apache.hudi.index.HoodieIndex; import org.apache.hudi.index.bucket.partition.PartitionBucketIndexUtils; import org.apache.hudi.keygen.KeyGenUtils; import org.apache.hudi.keygen.constant.KeyGeneratorOptions; +import org.apache.hudi.metadata.HoodieTableMetadataUtil; +import org.apache.hudi.metadata.MetadataPartitionType; import org.apache.hudi.sink.buffer.BufferMemoryType; import org.apache.hudi.sink.overwrite.PartitionOverwriteMode; import org.apache.hudi.table.format.FilePathUtils; @@ -61,12 +63,22 @@ import java.util.Locale; import java.util.Map; import static org.apache.hudi.common.config.HoodieCommonConfig.INCREMENTAL_READ_HANDLE_HOLLOW_COMMIT; +import static org.apache.hudi.common.config.HoodieMetadataConfig.GLOBAL_RECORD_LEVEL_INDEX_MAX_FILE_GROUP_COUNT_PROP; +import static org.apache.hudi.common.config.HoodieMetadataConfig.GLOBAL_RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_PROP; +import static org.apache.hudi.common.config.HoodieMetadataConfig.RECORD_INDEX_GROWTH_FACTOR_PROP; +import static org.apache.hudi.common.config.HoodieMetadataConfig.RECORD_INDEX_MAX_FILE_GROUP_SIZE_BYTES_PROP; +import static org.apache.hudi.common.config.HoodieMetadataConfig.RECORD_LEVEL_INDEX_MAX_FILE_GROUP_COUNT_PROP; +import static org.apache.hudi.common.config.HoodieMetadataConfig.RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_PROP; +import static org.apache.hudi.metadata.HoodieBackedTableMetadataWriter.RECORD_INDEX_AVERAGE_RECORD_SIZE; /** * Tool helping to resolve the flink options {@link FlinkOptions}. */ public class OptionsResolver { + // Value to override the default minimum file group count for global record level index. + public static String GLOBAL_RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_DEFAULT = "8"; + /** * Returns whether the current runtime mode is adaptive batch execution. */ @@ -227,6 +239,35 @@ public class OptionsResolver { return indexType == HoodieIndex.IndexType.GLOBAL_RECORD_LEVEL_INDEX; } + /** + * Estimates the file group count to use for RLI partition of a new table. + */ + public static int estimateFileGroupCountForRLI(Configuration conf) { + int minFileGroupCount; + int maxFileGroupCount; + if (OptionsResolver.isRecordLevelIndex(conf)) { + minFileGroupCount = Integer.parseInt(conf.getString(RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_PROP.key(), + RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_PROP.defaultValue() + "")); + maxFileGroupCount = Integer.parseInt(conf.getString(RECORD_LEVEL_INDEX_MAX_FILE_GROUP_COUNT_PROP.key(), + RECORD_LEVEL_INDEX_MAX_FILE_GROUP_COUNT_PROP.defaultValue() + "")); + } else { + minFileGroupCount = Integer.parseInt(conf.getString(GLOBAL_RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_PROP.key(), + GLOBAL_RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_DEFAULT)); + maxFileGroupCount = Integer.parseInt(conf.getString(GLOBAL_RECORD_LEVEL_INDEX_MAX_FILE_GROUP_COUNT_PROP.key(), + GLOBAL_RECORD_LEVEL_INDEX_MAX_FILE_GROUP_COUNT_PROP.defaultValue() + "")); + } + return HoodieTableMetadataUtil.estimateFileGroupCount( + MetadataPartitionType.RECORD_INDEX, + () -> 0L, + RECORD_INDEX_AVERAGE_RECORD_SIZE, + minFileGroupCount, + maxFileGroupCount, + Float.parseFloat(conf.getString(RECORD_INDEX_GROWTH_FACTOR_PROP.key(), + RECORD_INDEX_GROWTH_FACTOR_PROP.defaultValue() + "")), + Long.parseLong(conf.getString(RECORD_INDEX_MAX_FILE_GROUP_SIZE_BYTES_PROP.key(), + RECORD_INDEX_MAX_FILE_GROUP_SIZE_BYTES_PROP.defaultValue() + ""))); + } + /** * Returns whether it is a MERGE_ON_READ table, and updates by bucket index. */ diff --git a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/BucketAssignFunction.java b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/BucketAssignFunction.java index 7faff523e689..d2deba8ef030 100644 --- a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/BucketAssignFunction.java +++ b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/BucketAssignFunction.java @@ -172,9 +172,7 @@ public class BucketAssignFunction protected void processChangingRecord( HoodieFlinkInternalRow record, String recordKey, - Collector<HoodieFlinkInternalRow> out, - HoodieRecordGlobalLocation prefetchedOldLoc, - boolean prefetched) throws Exception { + Collector<HoodieFlinkInternalRow> out) throws Exception { // 1. put the record into the BucketAssigner; // 2. look up the state for location, if the record has a location, just send it out; // 3. if it is an INSERT, decide the location using the BucketAssigner then send it out. @@ -183,7 +181,7 @@ public class BucketAssignFunction // Only changing records need looking up the index for the location, // append only records are always recognized as INSERT. // Structured as Tuple(partition, fileId, instantTime). - HoodieRecordGlobalLocation oldLoc = prefetched ? prefetchedOldLoc : indexBackend.get(recordKey); + HoodieRecordGlobalLocation oldLoc = indexBackend.get(recordKey); if (oldLoc != null) { // Set up the instant time as "U" to mark the bucket as an update bucket. String partitionFromState = oldLoc.getPartitionPath(); @@ -238,7 +236,7 @@ public class BucketAssignFunction */ private Processor initRecordProcessor() { if (isChangingRecords) { - return (value, out) -> processChangingRecord(value, value.getRecordKey(), out, null, false); + return (value, out) -> processChangingRecord(value, value.getRecordKey(), out); } else { return this::processInsertRecord; } diff --git a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/GlobalRecordIndexPartitioner.java b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/GlobalRecordIndexPartitioner.java index 8b6f1ae22ad9..258c3f4a8fd3 100644 --- a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/GlobalRecordIndexPartitioner.java +++ b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/GlobalRecordIndexPartitioner.java @@ -22,6 +22,7 @@ import org.apache.hudi.client.common.HoodieFlinkEngineContext; import org.apache.hudi.common.model.HoodieKey; import org.apache.hudi.common.table.HoodieTableMetaClient; import org.apache.hudi.configuration.FlinkOptions; +import org.apache.hudi.configuration.OptionsResolver; import org.apache.hudi.exception.HoodieException; import org.apache.hudi.metadata.HoodieTableMetadata; import org.apache.hudi.metadata.HoodieTableMetadataUtil; @@ -79,6 +80,12 @@ public class GlobalRecordIndexPartitioner implements Partitioner<HoodieKey> { */ private int getNumFileGroupsForRecordIndexPartition() { HoodieTableMetaClient metaClient = StreamerUtil.createMetaClient(conf); + // For flink adaptive batch execution, writer coordinator is not started yet, so metadata table + // is not initialized for a new table. + if (!metaClient.getTableConfig().isMetadataPartitionAvailable(MetadataPartitionType.RECORD_INDEX)) { + // estimate the minimum file group count used to initialize global record level index + return OptionsResolver.estimateFileGroupCountForRLI(conf); + } try (HoodieTableMetadata metadataTable = metaClient.getTableFormat().getMetadataFactory().create( HoodieFlinkEngineContext.DEFAULT, metaClient.getStorage(), diff --git a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/MinibatchBucketAssignFunction.java b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/MinibatchBucketAssignFunction.java index cdf30c932eb7..3a44267c42a9 100644 --- a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/MinibatchBucketAssignFunction.java +++ b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/MinibatchBucketAssignFunction.java @@ -20,7 +20,6 @@ package org.apache.hudi.sink.partitioner; import org.apache.hudi.adapter.ProcessFunctionAdapter; import org.apache.hudi.client.model.HoodieFlinkInternalRow; -import org.apache.hudi.common.model.HoodieRecordGlobalLocation; import org.apache.hudi.common.model.WriteOperationType; import org.apache.hudi.common.util.VisibleForTesting; import org.apache.hudi.configuration.FlinkOptions; @@ -40,7 +39,6 @@ import org.apache.flink.util.Collector; import java.io.Serializable; import java.util.ArrayList; import java.util.List; -import java.util.Map; import java.util.stream.Collectors; /** @@ -152,11 +150,11 @@ public class MinibatchBucketAssignFunction public void process(List<HoodieFlinkInternalRow> records, Collector<HoodieFlinkInternalRow> out) throws Exception { List<String> recordKeys = records.stream().map(HoodieFlinkInternalRow::getRecordKey).collect(Collectors.toList()); MinibatchIndexBackend minibatchIndexBackend = (MinibatchIndexBackend) delegateFunction.getIndexBackend(); - // get record locations by minibatch - Map<String, HoodieRecordGlobalLocation> recordLocations = minibatchIndexBackend.get(recordKeys); + // warm up the in-memory cache for record level index + minibatchIndexBackend.get(recordKeys); for (HoodieFlinkInternalRow record: records) { String recordKey = record.getRecordKey(); - delegateFunction.processChangingRecord(record, recordKey, out, recordLocations.get(recordKey), true); + delegateFunction.processChangingRecord(record, recordKey, out); } } }; diff --git a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/RecordIndexPartitioner.java b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/RecordIndexPartitioner.java index a4597e00f099..e1d2a5f203ae 100644 --- a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/RecordIndexPartitioner.java +++ b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/RecordIndexPartitioner.java @@ -19,12 +19,12 @@ package org.apache.hudi.sink.partitioner; import org.apache.hudi.client.common.HoodieFlinkEngineContext; -import org.apache.hudi.common.config.HoodieMetadataConfig; import org.apache.hudi.common.model.HoodieKey; import org.apache.hudi.common.table.HoodieTableMetaClient; import org.apache.hudi.common.util.Functions; import org.apache.hudi.common.util.hash.BucketIndexUtil; import org.apache.hudi.configuration.FlinkOptions; +import org.apache.hudi.configuration.OptionsResolver; import org.apache.hudi.exception.HoodieException; import org.apache.hudi.metadata.HoodieTableMetadata; import org.apache.hudi.metadata.HoodieTableMetadataUtil; @@ -34,7 +34,6 @@ import org.apache.hudi.util.StreamerUtil; import org.apache.flink.api.common.functions.Partitioner; import org.apache.flink.configuration.Configuration; -import java.util.Collections; import java.util.HashMap; import java.util.Map; @@ -77,7 +76,7 @@ public class RecordIndexPartitioner implements Partitioner<HoodieKey> { if (partitionedRLIFileGroupCounts == null) { partitionedRLIFileGroupCounts = getPartitionedRLIFileGroupCounts(); } - int fileGroupCount = getFileGroupCountForPartitionedRLI(recordKey.getPartitionPath()); + int fileGroupCount = partitionedRLIFileGroupCounts.computeIfAbsent(recordKey.getPartitionPath(), s -> OptionsResolver.estimateFileGroupCountForRLI(conf)); int fgIndex = HoodieTableMetadataUtil.mapRecordKeyToFileGroupIndex(recordKey.getRecordKey(), fileGroupCount); if (partitionIndexFunc == null) { partitionIndexFunc = BucketIndexUtil.getPartitionIndexFunc(numPartitions); @@ -91,7 +90,7 @@ public class RecordIndexPartitioner implements Partitioner<HoodieKey> { private Map<String, Integer> getPartitionedRLIFileGroupCounts() { HoodieTableMetaClient metaClient = StreamerUtil.createMetaClient(conf); if (!metaClient.getTableConfig().isMetadataPartitionAvailable(MetadataPartitionType.RECORD_INDEX)) { - return Collections.emptyMap(); + return new HashMap<>(); } try (HoodieTableMetadata metadataTable = metaClient.getTableFormat().getMetadataFactory().create( HoodieFlinkEngineContext.DEFAULT, @@ -106,24 +105,4 @@ public class RecordIndexPartitioner implements Partitioner<HoodieKey> { throw new HoodieException("Failed to get file group counts for partitioned record index.", e); } } - - /** - * Get the partitioned record index file group count for the given data partition. - */ - private int getFileGroupCountForPartitionedRLI(String partitionPath) { - int fileGroupCount = partitionedRLIFileGroupCounts.getOrDefault(partitionPath, 0); - // HoodieBackedTableMetadataWriter initializes record-index file groups for a newly seen - // data partition with RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_PROP, so the writer-side - // partitioner should use the same count before that partition appears in the MDT view. - return fileGroupCount > 0 ? fileGroupCount : getMinFileGroupCountForPartitionedRLI(); - } - - /** - * Get the minimum file group count used to initialize newly seen partitioned record index partitions. - */ - private int getMinFileGroupCountForPartitionedRLI() { - return Integer.parseInt(conf.getString( - HoodieMetadataConfig.RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_PROP.key(), - HoodieMetadataConfig.RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_PROP.defaultValue().toString())); - } } diff --git a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/index/GlobalRecordLevelIndexBackend.java b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/index/GlobalRecordLevelIndexBackend.java index a95e18e863a1..93a6303009f8 100644 --- a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/index/GlobalRecordLevelIndexBackend.java +++ b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/index/GlobalRecordLevelIndexBackend.java @@ -20,6 +20,7 @@ package org.apache.hudi.sink.partitioner.index; import org.apache.hudi.client.common.HoodieFlinkEngineContext; import org.apache.hudi.common.data.HoodieListData; +import org.apache.hudi.common.data.HoodieListPairData; import org.apache.hudi.common.data.HoodiePairData; import org.apache.hudi.common.model.HoodieRecordGlobalLocation; import org.apache.hudi.common.table.HoodieTableMetaClient; @@ -27,7 +28,7 @@ import org.apache.hudi.common.util.VisibleForTesting; import org.apache.hudi.common.util.collection.Pair; import org.apache.hudi.configuration.FlinkOptions; import org.apache.hudi.exception.HoodieException; -import org.apache.hudi.metadata.HoodieTableMetadata; +import org.apache.hudi.metadata.HoodieBackedTableMetadata; import org.apache.hudi.sink.event.Correspondent; import org.apache.hudi.util.StreamerUtil; @@ -37,6 +38,7 @@ import org.apache.flink.configuration.Configuration; import java.io.IOException; import java.util.ArrayList; +import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -55,7 +57,7 @@ public class GlobalRecordLevelIndexBackend implements MinibatchIndexBackend { private final RecordIndexCache recordIndexCache; private final Configuration conf; private final HoodieTableMetaClient metaClient; - private HoodieTableMetadata metadataTable; + private HoodieBackedTableMetadata tableMetadata; /** * Creates a global RLI backend with a checkpoint-aware cache. @@ -72,7 +74,9 @@ public class GlobalRecordLevelIndexBackend implements MinibatchIndexBackend { @Override public HoodieRecordGlobalLocation get(String recordKey) throws IOException { - throw new UnsupportedOperationException(this.getClass().getSimpleName() + " doesn't support lookup with a single key."); + // note: always fetch record location from the cache, since this backend is only used for minibatch mode, + // and the cache has been warmed up by calling `get(List<String> recordKeys)` previously. + return recordIndexCache.get(recordKey); } @Override @@ -93,8 +97,7 @@ public class GlobalRecordLevelIndexBackend implements MinibatchIndexBackend { } } if (!missedKeys.isEmpty()) { - HoodiePairData<String, HoodieRecordGlobalLocation> recordIndexData = - metadataTable.readRecordIndexLocationsWithKeys(HoodieListData.eager(missedKeys)); + HoodiePairData<String, HoodieRecordGlobalLocation> recordIndexData = lookupLocationsForMissedKeys(missedKeys); recordIndexData.forEach(keyAndLocation -> { recordIndexCache.update(keyAndLocation.getKey(), keyAndLocation.getValue()); keysAndLocations.put(keyAndLocation.getKey(), keyAndLocation.getValue()); @@ -103,6 +106,15 @@ public class GlobalRecordLevelIndexBackend implements MinibatchIndexBackend { return keysAndLocations; } + private HoodiePairData<String, HoodieRecordGlobalLocation> lookupLocationsForMissedKeys(List<String> missedKeys) { + // For flink adaptive batch execution, writer coordinator is not started yet, so metadata table + // is not initialized for a new table. + if (!tableMetadata.enabled()) { + return HoodieListPairData.eager(Collections.emptyList()); + } + return tableMetadata.readRecordIndexLocationsWithKeys(HoodieListData.eager(missedKeys)); + } + @Override public void update(List<Pair<String, HoodieRecordGlobalLocation>> recordKeysAndLocations) throws IOException { recordKeysAndLocations.forEach(keyAndLocation -> recordIndexCache.update(keyAndLocation.getKey(), keyAndLocation.getValue())); @@ -141,21 +153,30 @@ public class GlobalRecordLevelIndexBackend implements MinibatchIndexBackend { } private void reloadMetadataTable() { - this.metadataTable = metaClient.getTableFormat().getMetadataFactory().create( - HoodieFlinkEngineContext.DEFAULT, - metaClient.getStorage(), - StreamerUtil.metadataConfig(conf), - conf.get(FlinkOptions.PATH)); + if (this.tableMetadata != null) { + this.tableMetadata.close(); + } + this.tableMetadata = + new HoodieBackedTableMetadata( + HoodieFlinkEngineContext.DEFAULT, + metaClient.getStorage(), + StreamerUtil.metadataConfig(conf), + conf.get(FlinkOptions.PATH)); + if (!tableMetadata.enabled()) { + if (metaClient.getTableConfig().isMetadataTableAvailable()) { + throw new RuntimeException("Can not initialize the table metadata"); + } + } } @Override public void close() throws IOException { this.recordIndexCache.close(); - if (this.metadataTable == null) { + if (this.tableMetadata == null) { return; } try { - this.metadataTable.close(); + this.tableMetadata.close(); } catch (Exception e) { throw new HoodieException("Exception happened during close metadata table.", e); } diff --git a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java index b17be182b9de..036f724b00d5 100644 --- a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java +++ b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java @@ -291,13 +291,13 @@ public class Pipelines { boolean bounded) { DataStream<HoodieFlinkInternalRow> dataStream1 = rowDataToHoodieRecord(conf, rowType, dataStream); - if (conf.get(FlinkOptions.INDEX_BOOTSTRAP_ENABLED) || bounded) { - boolean isRliBootstrap = OptionsResolver.isGlobalRecordLevelIndex(conf); + boolean isGlobalRLI = OptionsResolver.isGlobalRecordLevelIndex(conf); + if (conf.get(FlinkOptions.INDEX_BOOTSTRAP_ENABLED) || (bounded && !isGlobalRLI)) { dataStream1 = dataStream1 .transform( "index_bootstrap", new HoodieFlinkInternalRowTypeInfo(rowType), - isRliBootstrap ? new RLIBootstrapOperator(conf) : new BootstrapOperator(conf)) + isGlobalRLI ? new RLIBootstrapOperator(conf) : new BootstrapOperator(conf)) .setParallelism(conf.getOptional(FlinkOptions.INDEX_BOOTSTRAP_TASKS).orElse(dataStream1.getParallelism())) .uid(opUID("index_bootstrap", conf)); ((OneInputTransformation<?, ?>) dataStream1.getTransformation()).setChainingStrategy(ChainingStrategy.ALWAYS); diff --git a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/HoodieTableFactory.java b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/HoodieTableFactory.java index ae8dbb1e963b..43d18e849b3d 100644 --- a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/HoodieTableFactory.java +++ b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/HoodieTableFactory.java @@ -475,10 +475,6 @@ public class HoodieTableFactory implements DynamicTableSourceFactory, DynamicTab conf.setString(HoodieMetadataConfig.GLOBAL_RECORD_LEVEL_INDEX_ENABLE_PROP.key(), "true"); conf.set(FlinkOptions.INDEX_GLOBAL_ENABLED, true); conf.setString(HoodieMetadataConfig.STREAMING_WRITE_ENABLED.key(), "true"); - // set index bootstrap as true if not specified by user explicitly. - if (!conf.contains(FlinkOptions.INDEX_BOOTSTRAP_ENABLED)) { - conf.set(FlinkOptions.INDEX_BOOTSTRAP_ENABLED, true); - } // generally size of index data is much smaller than data record, so set the buffer size of // the index writer as 1/4 of that for data writer if it's not set by user explicitly. if (!conf.contains(FlinkOptions.INDEX_RLI_WRITE_BUFFER_SIZE)) { diff --git a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/FlinkWriteClients.java b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/FlinkWriteClients.java index a389d721908b..0293aad632df 100644 --- a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/FlinkWriteClients.java +++ b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/FlinkWriteClients.java @@ -50,6 +50,7 @@ import org.apache.flink.configuration.Configuration; import java.io.IOException; import java.util.Locale; +import static org.apache.hudi.configuration.OptionsResolver.GLOBAL_RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_DEFAULT; import static org.apache.hudi.util.StreamerUtil.flinkConf2TypedProperties; import static org.apache.hudi.util.StreamerUtil.getLockConfig; import static org.apache.hudi.util.StreamerUtil.getPayloadConfig; @@ -235,7 +236,8 @@ public class FlinkWriteClients { .withEngineType(EngineType.FLINK) // this affects the default value inference .enable(conf.get(FlinkOptions.METADATA_ENABLED)) .withRecordIndexFileGroupCount( - Integer.parseInt(conf.getString(HoodieMetadataConfig.GLOBAL_RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_PROP.key(), "8")), + Integer.parseInt(conf.getString(HoodieMetadataConfig.GLOBAL_RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_PROP.key(), + GLOBAL_RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_DEFAULT)), Integer.parseInt(conf.getString(HoodieMetadataConfig.GLOBAL_RECORD_LEVEL_INDEX_MAX_FILE_GROUP_COUNT_PROP.key(), HoodieMetadataConfig.GLOBAL_RECORD_LEVEL_INDEX_MAX_FILE_GROUP_COUNT_PROP.defaultValue() + ""))) .withMaxNumDeltaCommitsBeforeCompaction(conf.get(FlinkOptions.METADATA_COMPACTION_DELTA_COMMITS)) diff --git a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/configuration/TestOptionsResolver.java b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/configuration/TestOptionsResolver.java index 0ed8e61fd350..a70e7dc9f7df 100644 --- a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/configuration/TestOptionsResolver.java +++ b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/configuration/TestOptionsResolver.java @@ -18,6 +18,7 @@ package org.apache.hudi.configuration; +import org.apache.hudi.common.config.HoodieMetadataConfig; import org.apache.hudi.common.model.HoodieFailedWritesCleaningPolicy; import org.apache.hudi.common.model.HoodieTableType; import org.apache.hudi.common.model.WriteConcurrencyMode; @@ -25,6 +26,7 @@ import org.apache.hudi.common.model.WriteOperationType; import org.apache.hudi.config.HoodieCleanConfig; import org.apache.hudi.config.HoodieWriteConfig; import org.apache.hudi.index.HoodieIndex; +import org.apache.hudi.utils.TestConfigurations; import org.apache.flink.configuration.Configuration; import org.junit.jupiter.api.Test; @@ -191,4 +193,36 @@ public class TestOptionsResolver { assertFalse(OptionsResolver.needsAsyncClustering(conf)); assertFalse(OptionsResolver.needsScheduleClustering(conf)); } + + @Test + void testEstimateFileGroupCountForPartitionedRLI() { + Configuration conf = TestConfigurations.getDefaultConf(tempFile.getAbsolutePath()); + conf.set(FlinkOptions.METADATA_ENABLED, true); + conf.set(FlinkOptions.INDEX_TYPE, HoodieIndex.IndexType.RECORD_LEVEL_INDEX.name()); + conf.setString(HoodieMetadataConfig.RECORD_LEVEL_INDEX_ENABLE_PROP.key(), "true"); + + // testing default value + assertEquals(1, OptionsResolver.estimateFileGroupCountForRLI(conf)); + + // testing user configured value + conf.setString(HoodieMetadataConfig.RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_PROP.key(), "3"); + conf.setString(HoodieMetadataConfig.RECORD_LEVEL_INDEX_MAX_FILE_GROUP_COUNT_PROP.key(), "3"); + assertEquals(3, OptionsResolver.estimateFileGroupCountForRLI(conf)); + } + + @Test + void testEstimateFileGroupCountForGlobalRLI() { + Configuration conf = TestConfigurations.getDefaultConf(tempFile.getAbsolutePath()); + conf.set(FlinkOptions.METADATA_ENABLED, true); + conf.set(FlinkOptions.INDEX_TYPE, HoodieIndex.IndexType.GLOBAL_RECORD_LEVEL_INDEX.name()); + conf.setString(HoodieMetadataConfig.GLOBAL_RECORD_LEVEL_INDEX_ENABLE_PROP.key(), "true"); + + // testing default value + assertEquals(8, OptionsResolver.estimateFileGroupCountForRLI(conf)); + + // testing user configured value + conf.setString(HoodieMetadataConfig.GLOBAL_RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_PROP.key(), "11"); + conf.setString(HoodieMetadataConfig.GLOBAL_RECORD_LEVEL_INDEX_MAX_FILE_GROUP_COUNT_PROP.key(), "11"); + assertEquals(11, OptionsResolver.estimateFileGroupCountForRLI(conf)); + } } diff --git a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/TestMinibatchBucketAssignFunction.java b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/TestMinibatchBucketAssignFunction.java index a98a95529fc1..0e83ff81f75a 100644 --- a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/TestMinibatchBucketAssignFunction.java +++ b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/TestMinibatchBucketAssignFunction.java @@ -234,14 +234,94 @@ public class TestMinibatchBucketAssignFunction { } } + @Test + public void testDuplicateKeyCrossPartitionUpdate() throws Exception { + HoodieFlinkInternalRow record1 = insertRecord("id1", "par2", 1); + HoodieFlinkInternalRow record2 = insertRecord("id1", "par3", 2); + + testHarness.processElement(new StreamRecord<>(record1)); + testHarness.processElement(new StreamRecord<>(record2)); + testHarness.prepareSnapshotPreBarrier(1L); + + List<HoodieFlinkInternalRow> output = testHarness.extractOutputValues(); + assertEquals(4, output.size(), "Two cross-partition updates should each emit delete and insert records"); + + HoodieFlinkInternalRow deleteForOriginalPartition = output.get(0); + HoodieFlinkInternalRow insertForFirstUpdate = output.get(1); + HoodieFlinkInternalRow deleteForFirstUpdate = output.get(2); + HoodieFlinkInternalRow insertForSecondUpdate = output.get(3); + + assertEquals("par1", deleteForOriginalPartition.getPartitionPath(), + "First update should delete the location prefetched from the metadata table"); + assertEquals("-U", deleteForOriginalPartition.getOperationType()); + assertEquals("U", deleteForOriginalPartition.getInstantTime()); + + assertEquals("par2", insertForFirstUpdate.getPartitionPath()); + assertEquals("I", insertForFirstUpdate.getOperationType()); + assertEquals("U", insertForFirstUpdate.getInstantTime()); + assertTrue(insertForFirstUpdate.getFileId() != null && !insertForFirstUpdate.getFileId().isEmpty(), + "First update should be assigned to a data bucket"); + + assertEquals("par2", deleteForFirstUpdate.getPartitionPath(), + "Duplicate key should re-read the location updated by the preceding record in the same minibatch"); + assertEquals(insertForFirstUpdate.getFileId(), deleteForFirstUpdate.getFileId(), + "The delete for the second update should target the bucket assigned to the first update"); + assertEquals("-U", deleteForFirstUpdate.getOperationType()); + assertEquals("U", deleteForFirstUpdate.getInstantTime()); + + assertEquals("par3", insertForSecondUpdate.getPartitionPath()); + assertEquals("I", insertForSecondUpdate.getOperationType()); + assertEquals("U", insertForSecondUpdate.getInstantTime()); + assertTrue(insertForSecondUpdate.getFileId() != null && !insertForSecondUpdate.getFileId().isEmpty(), + "Second update should be assigned to a data bucket"); + } + + @Test + public void testDuplicateKeyInSameMinibatch() throws Exception { + HoodieFlinkInternalRow record1 = insertRecord("new_duplicate_key", "par_insert", 1); + HoodieFlinkInternalRow record2 = insertRecord("new_duplicate_key", "par_insert", 2); + + testHarness.processElement(new StreamRecord<>(record1)); + testHarness.processElement(new StreamRecord<>(record2)); + testHarness.prepareSnapshotPreBarrier(1L); + + List<HoodieFlinkInternalRow> output = testHarness.extractOutputValues(); + assertEquals(2, output.size(), "Duplicate insert-miss records should both be emitted"); + + HoodieFlinkInternalRow insertRecord = output.get(0); + HoodieFlinkInternalRow duplicateRecord = output.get(1); + + assertEquals("new_duplicate_key", insertRecord.getRecordKey()); + assertEquals("par_insert", insertRecord.getPartitionPath()); + assertEquals("I", insertRecord.getInstantTime(), + "The first record should be assigned as an insert because the prefetched location is missing"); + assertEquals("I", insertRecord.getOperationType()); + assertTrue(insertRecord.getFileId() != null && !insertRecord.getFileId().isEmpty(), + "First insert should be assigned to a data bucket"); + + assertEquals("new_duplicate_key", duplicateRecord.getRecordKey()); + assertEquals("par_insert", duplicateRecord.getPartitionPath()); + assertEquals(insertRecord.getFileId(), duplicateRecord.getFileId(), + "Duplicate key should re-read the location written by the first insert in the same minibatch"); + assertEquals("U", duplicateRecord.getInstantTime(), + "The duplicate record should be assigned as an update to the first record's bucket"); + assertEquals("I", duplicateRecord.getOperationType()); + } + @Test public void testCloseFunction() throws Exception { // Test that close doesn't throw exceptions HoodieFlinkInternalRow record = new HoodieFlinkInternalRow("id1", "par1", "I", insertRow(StringData.fromString("id1"), StringData.fromString("Danny"), 23, TimestampData.fromEpochMillis(1), StringData.fromString("par1"))); testHarness.processElement(new StreamRecord<>(record)); - + // Close should not throw any exceptions testHarness.close(); } + + private static HoodieFlinkInternalRow insertRecord(String recordKey, String partitionPath, long ts) { + return new HoodieFlinkInternalRow(recordKey, partitionPath, "I", + insertRow(StringData.fromString(recordKey), StringData.fromString("Danny"), 23, + TimestampData.fromEpochMillis(ts), StringData.fromString(partitionPath))); + } } diff --git a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/TestRecordIndexPartitioner.java b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/TestRecordIndexPartitioner.java index 59cd3d93fb0b..e8b6c757ba77 100644 --- a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/TestRecordIndexPartitioner.java +++ b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/TestRecordIndexPartitioner.java @@ -97,6 +97,7 @@ public class TestRecordIndexPartitioner { private RecordIndexPartitioner newPartitioner() throws Exception { Configuration conf = TestConfigurations.getDefaultConf(tempFile.getAbsolutePath()); conf.set(FlinkOptions.INDEX_TYPE, HoodieIndex.IndexType.RECORD_LEVEL_INDEX.name()); + conf.setString(HoodieMetadataConfig.RECORD_LEVEL_INDEX_ENABLE_PROP.key(), "true"); conf.setString(HoodieMetadataConfig.RECORD_LEVEL_INDEX_MIN_FILE_GROUP_COUNT_PROP.key(), String.valueOf(FILE_GROUP_COUNT)); StreamerUtil.initTableIfNotExists(conf); return new RecordIndexPartitioner(conf); diff --git a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/index/TestGlobalRecordLevelIndexBackend.java b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/index/TestGlobalRecordLevelIndexBackend.java index 3b3f72a9fcb5..37ac44d6ae06 100644 --- a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/index/TestGlobalRecordLevelIndexBackend.java +++ b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/index/TestGlobalRecordLevelIndexBackend.java @@ -81,10 +81,12 @@ public class TestGlobalRecordLevelIndexBackend { assertEquals("par1", location.getPartitionPath()); assertEquals(firstCommitTime, location.getInstantTime()); + HoodieRecordGlobalLocation location1 = globalRecordLevelIndexBackend.get("id1"); + assertEquals(location1, location); + // get record location with non existed key location = globalRecordLevelIndexBackend.get(Collections.singletonList("new_key")).get("new_key"); assertNull(location); - assertThrows(UnsupportedOperationException.class, () -> globalRecordLevelIndexBackend.get("id1")); // get records locations for multiple record keys Map<String, HoodieRecordGlobalLocation> locations = globalRecordLevelIndexBackend.get(Arrays.asList("id1", "id2", "id3")); diff --git a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/ITTestHoodieDataSource.java b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/ITTestHoodieDataSource.java index 4af3598c0a97..610a23b3b0e7 100644 --- a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/ITTestHoodieDataSource.java +++ b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/ITTestHoodieDataSource.java @@ -1537,8 +1537,8 @@ public class ITTestHoodieDataSource { } @ParameterizedTest - @EnumSource(value = HoodieIndex.IndexType.class, names = {"FLINK_STATE", "GLOBAL_RECORD_LEVEL_INDEX"}) - void testWriteGlobalIndex(HoodieIndex.IndexType indexType) { + @MethodSource("indexAndBooleanParams") + void testWriteGlobalIndex(String indexType, boolean bootstrapEnabled) { // the source generates 4 commits String createSource = TestConfigurations.getFileSourceDDL( "source", "test_source_4.data", 4); @@ -1548,7 +1548,8 @@ public class ITTestHoodieDataSource { .option(FlinkOptions.PATH, tempFile.getAbsolutePath()) .options(getDefaultKeys()) .option(FlinkOptions.INDEX_GLOBAL_ENABLED, true) - .option(FlinkOptions.INDEX_TYPE, indexType.name()) + .option(FlinkOptions.INDEX_TYPE, indexType) + .option(FlinkOptions.INDEX_BOOTSTRAP_ENABLED, bootstrapEnabled) .option(FlinkOptions.PRE_COMBINE, true) .end(); streamTableEnv.executeSql(hoodieTableDDL); @@ -3416,6 +3417,7 @@ public class ITTestHoodieDataSource { .option(FlinkOptions.PATH, tempFile.getAbsolutePath()) .options(getDefaultKeys()) .option(FlinkOptions.INDEX_TYPE, HoodieIndex.IndexType.GLOBAL_RECORD_LEVEL_INDEX.name()) + .option(FlinkOptions.INDEX_BOOTSTRAP_ENABLED, true) .option(FlinkOptions.READ_DATA_SKIPPING_ENABLED, true) .option(FlinkOptions.TABLE_TYPE, MERGE_ON_READ.name()) .end(); @@ -3717,6 +3719,18 @@ public class ITTestHoodieDataSource { return Stream.of(data).map(Arguments::of); } + /** + * Return test params => (index type, boolean). + */ + private static Stream<Arguments> indexAndBooleanParams() { + Object[][] data = + new Object[][] { + {"FLINK_STATE", false}, + {"GLOBAL_RECORD_LEVEL_INDEX", false}, + {"GLOBAL_RECORD_LEVEL_INDEX", true}}; + return Stream.of(data).map(Arguments::of); + } + /** * Return test params => (index type, table type). */ diff --git a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/TestHoodieTableFactory.java b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/TestHoodieTableFactory.java index 6aa6f20581c3..a85d3d29b6ae 100644 --- a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/TestHoodieTableFactory.java +++ b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/TestHoodieTableFactory.java @@ -785,6 +785,23 @@ public class TestHoodieTableFactory { (HoodieTableSink) new HoodieTableFactory().createDynamicTableSink(MockContext.getInstance(this.conf)); Configuration conf2 = tableSink2.getConf(); assertThat(conf2.get(FlinkOptions.PRE_COMBINE), is(false)); + + // Global RLI setup should not enable index bootstrap implicitly. + Configuration globalRLIConf = new Configuration(this.conf); + globalRLIConf.set(FlinkOptions.OPERATION, "upsert"); + globalRLIConf.set(FlinkOptions.INDEX_TYPE, HoodieIndex.IndexType.GLOBAL_RECORD_LEVEL_INDEX.name()); + globalRLIConf.set(FlinkOptions.METADATA_ENABLED, true); + globalRLIConf.set(FlinkOptions.INDEX_GLOBAL_ENABLED, true); + HoodieTableSink globalRLISink = + (HoodieTableSink) new HoodieTableFactory().createDynamicTableSink(MockContext.getInstance(globalRLIConf)); + Configuration globalRLIResolvedConf = globalRLISink.getConf(); + assertThat(globalRLIResolvedConf.get(FlinkOptions.INDEX_BOOTSTRAP_ENABLED), is(false)); + + globalRLIConf.set(FlinkOptions.INDEX_BOOTSTRAP_ENABLED, true); + HoodieTableSink globalRLIWithBootstrapSink = + (HoodieTableSink) new HoodieTableFactory().createDynamicTableSink(MockContext.getInstance(globalRLIConf)); + Configuration globalRLIWithBootstrapResolvedConf = globalRLIWithBootstrapSink.getConf(); + assertThat(globalRLIWithBootstrapResolvedConf.get(FlinkOptions.INDEX_BOOTSTRAP_ENABLED), is(true)); } @Test
