PyFlink on Yarn提交Kafka消费作业报NoClassDefFoundError如何解决
问题解决方案
根因说明
报错核心原因是Flink运行时加载不到Kafka客户端的相关类,主要由Jar包版本冲突、参数使用错误、依赖缺失三类问题导致。
解决步骤
第一步:清理并替换正确的Kafka连接器Jar包
你当前/opt/flink/lib目录下的两个Kafka连接器Jar包存在两个问题:- Scala版本不统一,一个是Scala 2.12版本,一个是Scala 2.11版本,会引发类加载冲突
- Jar包命名存在格式错误(版本号分隔符应为
-而非.),可能导致程序无法正常识别加载
操作:删除原有两个Kafka相关Jar包,下载与你Flink 1.13.2版本Scala版本完全匹配的带Kafka客户端依赖的flink-sql-connector-kafka Uber包,放置到所有Flink节点的/opt/flink/lib目录下。如果下载的是不含Kafka依赖的轻量版连接器,需要同步上传兼容版本的kafka-clients Jar包。
第二步:修正作业提交参数
--jarfile参数仅会将Jar包加载到本地提交端的类路径,YARN集群的JobManager、TaskManager节点无法读取该Jar包,需要替换为-yj(--yarnship)参数将Jar包同步分发到集群所有运行节点的类路径。
修正后的提交命令示例:/opt/flink/bin/flink run -m yarn-cluster -yid application_1634021687380_0009 -yj /opt/flink/lib/flink-sql-connector-kafka_<你的Flink对应Scala版本>-1.13.2.jar -pyarch venv.zip -pyexec venv.zip/venv/bin/python -py demo.py第三步:可选优化SLF4J依赖冲突
日志中的SLF4J绑定冲突不影响作业核心运行,如果需要消除该警告,可以删除/usr/hdp/3.1.0.0-78/hadoop/lib/slf4j-log4j12-1.7.25.jar或者临时移除Flink lib目录下的log4j-slf4j-impl Jar包即可。
内容的提问来源于stack exchange,提问作者Ronnie
相关产品推荐
相关产品推荐

