LinMingQiang commented on code in PR #5997: URL: https://github.com/apache/hudi/pull/5997#discussion_r912315747
########## hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java: ########## @@ -330,17 +330,23 @@ public static DataStream<Object> hoodieStreamWrite(Configuration conf, int defau .setParallelism(conf.getInteger(FlinkOptions.WRITE_TASKS)); } else { WriteOperatorFactory<HoodieRecord> operatorFactory = StreamWriteOperator.getFactory(conf); - return dataStream - // Key-by record key, to avoid multiple subtasks write to a bucket at the same time - .keyBy(HoodieRecord::getRecordKey) - .transform( - "bucket_assigner", - TypeInformation.of(HoodieRecord.class), - new KeyedProcessOperator<>(new BucketAssignFunction<>(conf))) - .uid("uid_bucket_assigner_" + conf.getString(FlinkOptions.TABLE_NAME)) - .setParallelism(conf.getOptional(FlinkOptions.BUCKET_ASSIGN_TASKS).orElse(defaultParallelism)) - // shuffle by fileId(bucket id) - .keyBy(record -> record.getCurrentLocation().getFileId()) + + DataStream<HoodieRecord> bucketDataStream = dataStream + // Key-by record key, to avoid multiple subtasks write to a bucket at the same time + .keyBy(HoodieRecord::getRecordKey) + .transform( + "bucket_assigner", + TypeInformation.of(HoodieRecord.class), + new KeyedProcessOperator<>(new BucketAssignFunction<>(conf))) + .uid("uid_bucket_assigner_" + conf.getString(FlinkOptions.TABLE_NAME)) + .setParallelism(conf.getOptional(FlinkOptions.BUCKET_ASSIGN_TASKS).orElse(defaultParallelism)); + + bucketDataStream = conf.getOptional(FlinkOptions.BUCKET_ASSIGN_TASKS).orElse(defaultParallelism) == + conf.getInteger(FlinkOptions.WRITE_TASKS) ? bucketDataStream : bucketDataStream + // shuffle by fileId(bucket id) + .keyBy(record -> record.getCurrentLocation().getFileId()); Review Comment: `bucket index` does not need to be repartitioned by `fileId`. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: commits-unsubscr...@hudi.apache.org For queries about this service, please contact Infrastructure at: us...@infra.apache.org