Apache Flink中dataStream.sinkTo()无法接收KafkaSink参数的问题
问题
我是Apache Flink新手,尝试从Kafka读取数据、在Flink中处理后写入另一个Kafka主题。已添加相关依赖并编写代码,但finalStream.sinkTo()方法无法接收KafkaSink<String>对象作为参数,它期望的参数类型为Sink<String, ?, ?, ?>,请问我遗漏了什么?
添加的依赖
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-core</artifactId> <version>1.17.1</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka</artifactId> <version>1.17.1</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java_2.11</artifactId> <version>1.14.6</version> </dependency> <dependency> <groupId>junit</groupId> <artifactId>junit</artifactId> <version>3.8.1</version> <scope>test</scope> </dependency>
代码实现
public class App { public static void main( String[] args ) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); KafkaSource<String> kafkaSource = KafkaSource.<String>builder() .setBootstrapServers("localhost:9092") .setTopics("webapp-analytics") .setDeserializer(KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class)) .setStartingOffsets(OffsetsInitializer.latest()) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); DataStreamSource<String> source = env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "Kafka Source"); DataStream<String> finalStream = source.filter(new FilterFunction<String>() { @Override public boolean filter(String s) throws Exception { return s.contains("clicked"); } }); KafkaSink<String> sink = KafkaSink.<String>builder() .setBootstrapServers("localhost:9092") .setRecordSerializer(KafkaRecordSerializationSchema.builder() .setTopic("flink") .setValueSerializationSchema(new SimpleStringSchema()) .build() ) .setDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE) .build(); finalStream.sinkTo(sink); env.execute("Kafka"); } }
解决方案
你的问题核心是Flink依赖版本不兼容,具体问题和修复步骤如下:
版本不一致导致API不匹配
你混用了Flink 1.17.1的核心、Kafka连接器依赖,和1.14.6的流处理核心依赖。KafkaSink是Flink 1.15+才引入的新Sink API,而1.14.6版本的DataStream.sinkTo()方法仅支持旧版Sink接口实现,二者类型无法兼容。修复步骤
- 统一所有Flink相关依赖的版本,将
flink-streaming-java_2.11的版本改为1.17.1。注意:Flink 1.17.x默认适配Scala 2.12/2.13,若坚持使用Scala 2.11,需确认对应版本是否存在;更稳妥的方式是换成不带Scala后缀的flink-streaming-java依赖,避免Scala版本冲突。 - 修正后的依赖示例:
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-core</artifactId> <version>1.17.1</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka</artifactId> <version>1.17.1</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>1.17.1</version> </dependency> <dependency> <groupId>junit</groupId> <artifactId>junit</artifactId> <version>3.8.1</version> <scope>test</scope> </dependency>
- 统一所有Flink相关依赖的版本,将
额外代码优化建议
- KafkaSource构建代码中同时调用了
setDeserializer和setValueOnlyDeserializer,二者互斥,保留其中一个即可,推荐用setValueOnlyDeserializer(new SimpleStringSchema())简化配置。
- KafkaSource构建代码中同时调用了
内容的提问来源于stack exchange,提问作者Ram
相关产品推荐
相关产品推荐

