diff --git a/hudi-flink/src/main/java/org/apache/hudi/configuration/FlinkOptions.java b/hudi-flink/src/main/java/org/apache/hudi/configuration/FlinkOptions.java index 1be90603605cd..37d3ceeaabf81 100644 --- a/hudi-flink/src/main/java/org/apache/hudi/configuration/FlinkOptions.java +++ b/hudi-flink/src/main/java/org/apache/hudi/configuration/FlinkOptions.java @@ -387,6 +387,12 @@ private FlinkOptions() { .noDefaultValue() .withDescription("Parallelism of tasks that do bucket assign, default is the parallelism of the execution environment"); + public static final ConfigOption HOODIE_ROECORD_MAP_TASKS = ConfigOptions + .key("hoodie.record_map.tasks") + .intType() + .noDefaultValue() + .withDescription("Parallelism of tasks that convert row data to hoodie record, default is the parallelism of the execution environment"); + public static final ConfigOption WRITE_TASKS = ConfigOptions .key("write.tasks") .intType() diff --git a/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java b/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java index ae8b4f21300a2..7eb918ae08bff 100644 --- a/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java +++ b/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java @@ -189,7 +189,7 @@ public static DataStream bootstrap( boolean overwrite) { final boolean globalIndex = conf.getBoolean(FlinkOptions.INDEX_GLOBAL_ENABLED); if (overwrite || OptionsResolver.isBucketIndexType(conf)) { - return rowDataToHoodieRecord(conf, rowType, dataStream); + return rowDataToHoodieRecord(conf, rowType, dataStream, defaultParallelism); } else if (bounded && !globalIndex && OptionsResolver.isPartitionedTable(conf)) { return boundedBootstrap(conf, rowType, defaultParallelism, dataStream); } else { @@ -203,7 +203,7 @@ private static DataStream streamBootstrap( int defaultParallelism, DataStream dataStream, boolean bounded) { - DataStream dataStream1 = rowDataToHoodieRecord(conf, rowType, dataStream); + DataStream dataStream1 = rowDataToHoodieRecord(conf, rowType, dataStream, defaultParallelism); if (conf.getBoolean(FlinkOptions.INDEX_BOOTSTRAP_ENABLED) || bounded) { dataStream1 = dataStream1 @@ -233,7 +233,7 @@ private static DataStream boundedBootstrap( dataStream = dataStream .keyBy(rowDataKeyGen::getPartitionPath); - return rowDataToHoodieRecord(conf, rowType, dataStream) + return rowDataToHoodieRecord(conf, rowType, dataStream, defaultParallelism) .transform( "batch_index_bootstrap", TypeInformation.of(HoodieRecord.class), @@ -245,8 +245,9 @@ private static DataStream boundedBootstrap( /** * Transforms the row data to hoodie records. */ - public static DataStream rowDataToHoodieRecord(Configuration conf, RowType rowType, DataStream dataStream) { - return dataStream.map(RowDataToHoodieFunctions.create(rowType, conf), TypeInformation.of(HoodieRecord.class)); + public static DataStream rowDataToHoodieRecord(Configuration conf, RowType rowType, DataStream dataStream, int defaultParallelism) { + return dataStream.map(RowDataToHoodieFunctions.create(rowType, conf), TypeInformation.of(HoodieRecord.class)) + .setParallelism(conf.getOptional(FlinkOptions.HOODIE_ROECORD_MAP_TASKS).orElse(defaultParallelism)); } /**