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.FlinkKafkaConsumerAPI,不存在代码中调用的新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
相关产品推荐
相关产品推荐

