Flink单词计数示例运行报错java.lang.ClassNotFoundException求助
问题分析与解决
核心原因
org.apache.flink.connector.base.source.reader.RecordEmitter 类属于 flink-connector-base 依赖包,你的POM中缺少这个依赖,导致运行时无法找到该类。此外你的POM还存在版本兼容问题:
flink-connector-jdbc版本3.2.0与Flink 1.20.0不匹配,官方要求连接器版本需与Flink核心版本对应- 使用
systemscope引入本地JAR会导致依赖无法被正确打包到Job包中,运行时极易出现缺失
修复后的POM配置
替换你的<dependencies>部分为以下内容:
<dependencies> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-core</artifactId> <version>1.20.0</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>1.20.0</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-walkthrough-common</artifactId> <version>1.20.0</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-clients</artifactId> <version>1.20.0</version> </dependency> <!-- 添加缺失的flink-connector-base依赖 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-base</artifactId> <version>1.20.0</version> </dependency> <!-- 保持与Flink 1.20匹配的kafka连接器版本 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka</artifactId> <version>3.3.0-1.20</version> </dependency> <!-- 使用与Flink 1.20匹配的JDBC连接器版本,移除system scope --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-jdbc</artifactId> <version>3.3.0-1.20</version> </dependency> <dependency> <groupId>log4j</groupId> <artifactId>log4j</artifactId> <version>1.2.17</version> </dependency> <dependency> <groupId>mysql</groupId> <artifactId>mysql-connector-java</artifactId> <version>8.0.33</version> </dependency> </dependencies>
额外操作建议
- 执行
mvn clean install清理并重新构建项目,确保依赖被正确下载和打包 - 如果提交到Flink集群运行,建议配置
maven-shade-plugin将所有依赖打包成一个fat jar,避免集群环境中依赖缺失 - 所有Flink官方连接器版本必须与Flink核心版本保持一致,彻底规避版本冲突问题
内容的提问来源于stack exchange,提问作者Piyush Jha
相关产品推荐
相关产品推荐

