Skip to content

Commit 5e3a38b

Browse files
committed
[fix] send parallelism default value 1
1 parent a8f4afa commit 5e3a38b

7 files changed

Lines changed: 7 additions & 7 deletions

File tree

cassandra/cassandra-sink/src/main/java/com/dtstack/flink/sql/sink/cassandra/CassandraSink.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -58,7 +58,7 @@ public class CassandraSink implements RetractStreamTableSink<Row>, IStreamSinkGe
5858
protected Integer readTimeoutMillis;
5959
protected Integer connectTimeoutMillis;
6060
protected Integer poolTimeoutMillis;
61-
protected Integer parallelism = -1;
61+
protected Integer parallelism = 1;
6262
protected String registerTableName;
6363

6464
public CassandraSink() {

core/src/main/java/com/dtstack/flink/sql/table/AbstractTableInfo.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -59,7 +59,7 @@ public abstract class AbstractTableInfo implements Serializable {
5959

6060
private List<String> primaryKeys;
6161

62-
private Integer parallelism = -1;
62+
private Integer parallelism = 1;
6363

6464
public String[] getFieldTypes() {
6565
return fieldTypes;

elasticsearch5/elasticsearch5-sink/src/main/java/com/dtstack/flink/sql/sink/elasticsearch/ElasticsearchSink.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -76,7 +76,7 @@ public class ElasticsearchSink implements RetractStreamTableSink<Row>, IStreamSi
7676

7777
private TypeInformation[] fieldTypes;
7878

79-
private int parallelism = -1;
79+
private int parallelism = 1;
8080

8181
protected String registerTableName;
8282

elasticsearch6/elasticsearch6-sink/src/main/java/com/dtstack/flink/sql/sink/elasticsearch/ElasticsearchSink.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -69,7 +69,7 @@ public class ElasticsearchSink implements RetractStreamTableSink<Row>, IStreamSi
6969

7070
private TypeInformation[] fieldTypes;
7171

72-
private int parallelism = -1;
72+
private int parallelism = 1;
7373

7474
protected String registerTableName;
7575

hbase/hbase-sink/src/main/java/com/dtstack/flink/sql/sink/hbase/HbaseSink.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -63,7 +63,7 @@ public class HbaseSink implements RetractStreamTableSink<Row>, IStreamSinkGener<
6363

6464
private String clientPrincipal;
6565
private String clientKeytabFile;
66-
private int parallelism = -1;
66+
private int parallelism = 1;
6767

6868

6969
public HbaseSink() {

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;

rdb/rdb-sink/src/main/java/com/dtstack/flink/sql/sink/rdb/AbstractRdbSink.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -72,7 +72,7 @@ public abstract class AbstractRdbSink implements RetractStreamTableSink<Row>, Se
7272

7373
private TypeInformation[] fieldTypes;
7474

75-
private int parallelism = -1;
75+
private int parallelism = 1;
7676

7777
protected String schema;
7878

0 commit comments

Comments
 (0)