You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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依赖版本不兼容,具体问题和修复步骤如下:

  1. 版本不一致导致API不匹配
    你混用了Flink 1.17.1的核心、Kafka连接器依赖,和1.14.6的流处理核心依赖。KafkaSink是Flink 1.15+才引入的新Sink API,而1.14.6版本的DataStream.sinkTo()方法仅支持旧版Sink接口实现,二者类型无法兼容。

  2. 修复步骤

    • 统一所有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>
      
  3. 额外代码优化建议

    • KafkaSource构建代码中同时调用了setDeserializer和setValueOnlyDeserializer,二者互斥,保留其中一个即可,推荐用setValueOnlyDeserializer(new SimpleStringSchema())简化配置。

内容的提问来源于stack exchange,提问作者Ram

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.08 22:34:54