Skip to content

Commit 44af091

Browse files
committed
test kafka bugfix
1 parent c57c84f commit 44af091

File tree

1 file changed

+1
-2
lines changed

1 file changed

+1
-2
lines changed

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

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,6 @@
2121
import com.dtstack.flink.sql.sink.kafka.table.KafkaSinkTableInfo;
2222
import org.apache.flink.api.common.typeinfo.TypeInformation;
2323
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
24-
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer;
2524
import org.apache.flink.streaming.connectors.kafka.partitioner.FlinkKafkaPartitioner;
2625
import org.apache.flink.types.Row;
2726

@@ -38,6 +37,6 @@ public class KafkaProducerFactory extends AbstractKafkaProducerFactory {
3837

3938
@Override
4039
public RichSinkFunction<Row> createKafkaProducer(KafkaSinkTableInfo kafkaSinkTableInfo, TypeInformation<Row> typeInformation, Properties properties, Optional<FlinkKafkaPartitioner<Row>> partitioner) {
41-
return new FlinkKafkaProducer<Row>(kafkaSinkTableInfo.getTopic(), createSerializationMetricWrapper(kafkaSinkTableInfo, typeInformation), properties, partitioner);
40+
return new KafkaProducer(kafkaSinkTableInfo.getTopic(), createSerializationMetricWrapper(kafkaSinkTableInfo, typeInformation), properties, partitioner);
4241
}
4342
}

0 commit comments

Comments
 (0)