Skip to content

Commit b335aa2

Browse files
author
gituser
committed
Merge branch 'hotfix_1.10_4.0.x_32575' into 1.10_release_4.0.x
2 parents edfe0f3 + fd58fee commit b335aa2

File tree

2 files changed

+2
-2
lines changed

2 files changed

+2
-2
lines changed

kafka-base/kafka-base-sink/src/main/java/com/dtstack/flink/sql/sink/kafka/AbstractKafkaSink.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -58,7 +58,7 @@ public abstract class AbstractKafkaSink implements RetractStreamTableSink<Row>,
5858
protected String[] partitionKeys;
5959
protected String sinkOperatorName;
6060
protected Properties properties;
61-
protected int parallelism = 1;
61+
protected int parallelism = -1;
6262
protected String topic;
6363
protected String tableName;
6464
protected String updateMode;

kafka-base/kafka-base-sink/src/main/java/com/dtstack/flink/sql/sink/kafka/table/KafkaSinkParser.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -57,7 +57,7 @@ public AbstractTableInfo getTableInfo(String tableName, String fieldsInfo, Map<S
5757
kafkaSinkTableInfo.setPartitionKeys(MathUtil.getString(props.get(KafkaSinkTableInfo.PARTITION_KEY.toLowerCase())));
5858
kafkaSinkTableInfo.setUpdateMode(MathUtil.getString(props.getOrDefault(KafkaSinkTableInfo.UPDATE_KEY.toLowerCase(), EUpdateMode.APPEND.name())));
5959

60-
Integer parallelism = MathUtil.getIntegerVal(props.getOrDefault(KafkaSinkTableInfo.PARALLELISM_KEY.toLowerCase(), 1));
60+
Integer parallelism = MathUtil.getIntegerVal(props.get(KafkaSinkTableInfo.PARALLELISM_KEY.toLowerCase()));
6161
kafkaSinkTableInfo.setParallelism(parallelism);
6262

6363
for (String key : props.keySet()) {

0 commit comments

Comments
 (0)