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

Flink任务部署报org.apache.flink.connector.kafka.source.KafkaSource类缺失异常

问题根因分析
  • 版本不匹配:org.apache.flink.connector.kafka.source.KafkaSource 是Flink 1.14版本才正式加入官方Kafka连接器的新API,你当前使用的1.13.0版本Kafka连接器仅提供旧版的org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer API,不存在代码中调用的新KafkaSource类,这是报错的核心原因。
  • 打包依赖缺失:即使版本匹配,若提交的Flink作业Jar只打包了项目自身代码、未包含连接器依赖,且Flink集群lib目录下也没有预放置对应版本的Kafka连接器Jar,也会触发同类类缺失异常。
解决方案

你可以根据自身需求选择任意一种方案修复:

方案1:升级Flink相关依赖版本到1.14+(适配新KafkaSource API)

如果要继续使用代码里的新KafkaSource写法,把所有Flink相关依赖的版本统一升级到1.14.0及以上稳定版本即可,pom依赖修改示例:

<!-- 所有Flink依赖版本需和集群运行版本完全一致 -->
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-connector-kafka_2.12</artifactId>
    <version>1.14.0</version>
</dependency>
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-avro</artifactId>
    <version>1.14.0</version>
</dependency>
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-clients_2.12</artifactId>
    <version>1.14.0</version>
</dependency>

方案2:修改代码适配1.13版本旧API(不升级依赖版本)

如果要保留当前1.13.0的依赖版本,把代码里的KafkaSource调用替换为旧版FlinkKafkaConsumer写法即可,代码示例:

Properties kafkaProps = new Properties();
kafkaProps.setProperty("bootstrap.servers", "Kafka集群地址");
kafkaProps.setProperty("group.id", "消费组ID");
// 替换原KafkaSource实现为旧版FlinkKafkaConsumer
FlinkKafkaConsumer<String> kafkaConsumer = new FlinkKafkaConsumer<>(
    "待消费Topic名称",
    new SimpleStringSchema(),
    kafkaProps
);
DataStream<String> kafkaStream = env.addSource(kafkaConsumer);
打包注意事项
  • 提交到集群运行的作业Jar需要通过maven-shade-plugin或maven-assembly-plugin打包为包含所有第三方依赖的Fat Jar,避免集群缺失依赖。
  • 如果你选择将Kafka连接器Jar预先放到Flink集群所有节点的lib目录下并重启集群,也可以不把连接器依赖打入作业Jar,注意集群放置的Jar版本要和项目依赖版本完全一致。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 14:18:02