This is an automated email from the ASF dual-hosted git repository. martijnvisser pushed a commit to branch main in repository https://gitbox.apache.org/repos/asf/flink-connector-kafka.git
commit 18135901964760ce3311c35fd6202b59f6c87ced Author: JunRuiLee <jrlee....@gmail.com> AuthorDate: Mon Jan 30 00:49:43 2023 +0800 [FLINK-30684][runtime] Mark vertices which use the default parallelism For these vertices, the parallelism can be changed at runtime. This closes #21772. --- .../flink/streaming/connectors/kafka/shuffle/FlinkKafkaShuffle.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/shuffle/FlinkKafkaShuffle.java b/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/shuffle/FlinkKafkaShuffle.java index d11f899..ae9af29 100644 --- a/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/shuffle/FlinkKafkaShuffle.java +++ b/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/shuffle/FlinkKafkaShuffle.java @@ -377,7 +377,8 @@ public class FlinkKafkaShuffle { inputStream.getTransformation(), "kafka_shuffle", shuffleSinkOperator, - inputStream.getExecutionEnvironment().getParallelism()); + inputStream.getExecutionEnvironment().getParallelism(), + false); inputStream.getExecutionEnvironment().addOperator(transformation); transformation.setParallelism(producerParallelism); }