Skip to content

Commit aa5e8c5

Browse files
author
dapeng
committed
去除无用的引用
1 parent ea131b6 commit aa5e8c5

File tree

4 files changed

+0
-22
lines changed

4 files changed

+0
-22
lines changed

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

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,8 +3,6 @@
33

44
import org.apache.flink.streaming.connectors.kafka.partitioner.FlinkFixedPartitioner;
55

6-
import java.util.Arrays;
7-
import java.util.Random;
86
import org.apache.flink.util.Preconditions;
97

108
public class CustomerFlinkPartition<T> extends FlinkFixedPartitioner<T> {

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

Lines changed: 0 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -137,12 +137,6 @@ public TableSink<Tuple2<Boolean, Row>> configure(String[] fieldNames, TypeInform
137137
this.fieldTypes = fieldTypes;
138138
return this;
139139
}
140-
private FlinkKafkaPartitioner getFlinkPartitioner(KafkaSinkTableInfo kafkaSinkTableInfo){
141-
if("true".equalsIgnoreCase(kafkaSinkTableInfo.getEnableKeyPartition())){
142-
return new CustomerFlinkPartition<>();
143-
}
144-
return new FlinkFixedPartitioner<>();
145-
}
146140

147141
private String[] getPartitionKeys(KafkaSinkTableInfo kafkaSinkTableInfo){
148142
if(StringUtils.isNotBlank(kafkaSinkTableInfo.getPartitionKeys())){

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

Lines changed: 0 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -144,13 +144,6 @@ public TableSink<Tuple2<Boolean, Row>> configure(String[] fieldNames, TypeInform
144144
return this;
145145
}
146146

147-
private FlinkKafkaPartitioner getFlinkPartitioner(KafkaSinkTableInfo kafkaSinkTableInfo){
148-
if("true".equalsIgnoreCase(kafkaSinkTableInfo.getEnableKeyPartition())){
149-
return new CustomerFlinkPartition<>();
150-
}
151-
return new FlinkFixedPartitioner<>();
152-
}
153-
154147
private String[] getPartitionKeys(KafkaSinkTableInfo kafkaSinkTableInfo){
155148
if(StringUtils.isNotBlank(kafkaSinkTableInfo.getPartitionKeys())){
156149
return StringUtils.split(kafkaSinkTableInfo.getPartitionKeys(), ',');

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

Lines changed: 0 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -144,13 +144,6 @@ public TypeInformation<?>[] getFieldTypes() {
144144
return this;
145145
}
146146

147-
private FlinkKafkaPartitioner getFlinkPartitioner(KafkaSinkTableInfo kafkaSinkTableInfo){
148-
if("true".equalsIgnoreCase(kafkaSinkTableInfo.getEnableKeyPartition())){
149-
return new CustomerFlinkPartition<>();
150-
}
151-
return new FlinkFixedPartitioner<>();
152-
}
153-
154147
private String[] getPartitionKeys(KafkaSinkTableInfo kafkaSinkTableInfo){
155148
if(StringUtils.isNotBlank(kafkaSinkTableInfo.getPartitionKeys())){
156149
return StringUtils.split(kafkaSinkTableInfo.getPartitionKeys(), ',');

0 commit comments

Comments
 (0)