Skip to content

Commit 72ba34e

Browse files
committed
kafka11 parallelism
1 parent acd9335 commit 72ba34e

File tree

1 file changed

+1
-1
lines changed
  • kafka11/kafka11-sink/src/main/java/com/dtstack/flink/sql/sink/kafka

1 file changed

+1
-1
lines changed

kafka11/kafka11-sink/src/main/java/com/dtstack/flink/sql/sink/kafka/KafkaSink.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -98,7 +98,7 @@ public void emitDataStream(DataStream<Tuple2<Boolean, Row>> dataStream) {
9898

9999
DataStream<Row> ds = dataStream.map((Tuple2<Boolean, Row> record) -> {
100100
return record.f1;
101-
}).returns(getOutputType().getTypeAt(1));
101+
}).returns(getOutputType().getTypeAt(1)).setParallelism(parallelism);
102102

103103
kafkaTableSink.emitDataStream(ds);
104104
}

0 commit comments

Comments
 (0)