PySpark消费Kafka数据时spark-submit提交失败求助
我之前也踩过这个坑——明明日志显示Spark已经加载了指定的Kafka JAR包,却还是抛出"Spark Streaming's Kafka libraries not found"的错误,分享几个我试过的有效排查方向:
1. 优先尝试用--packages替代--jars
本地的assembly JAR包有时候会因为下载不完整、版本隐性不兼容等问题导致加载失败,直接让Spark自动拉取官方匹配的依赖包更靠谱:
spark-submit --packages org.apache.spark:spark-streaming-kafka-0-8_2.11:2.2.0 consumer.py
这个命令会自动从Maven仓库下载对应Spark 2.2.0的Kafka Streaming依赖及所有关联包,避免本地JAR的潜在问题。
2. 核对提交命令的文件名拼写
注意到你提交命令里写的是comsumer.py,但日志里显示加载的是consumer.py(少了一个's')!虽然日志显示文件被添加了,但如果实际执行时文件名拼写错误,可能加载的是旧的、没有正确导入KafkaUtils的脚本版本,务必确保提交命令里的文件名和本地文件完全一致。
3. 检查JAR包的完整性与版本兼容性
虽然你用的spark-streaming-kafka-0-8-assembly_2.11-2.2.0.jar版本和Spark 2.2.0匹配,但可以:
- 重新下载该JAR包,替换本地的旧文件,避免文件损坏
- 确认你的Spark是基于Scala 2.11编译的(Spark 2.2.0默认是2.11版本,这点大概率没问题,但如果你的JAR是Scala 2.12版本就会不兼容)
4. 排查Python代码的导入与API调用
看错误栈是在调用KafkaUtils.createStream时触发的,先确认代码里的导入语句正确:
from pyspark.streaming.kafka import KafkaUtils
另外,createStream是针对Kafka 0.8.x版本的API,如果你的Kafka集群是0.9+版本,建议改用createDirectStream,不过这个错误主要是类路径问题,先排除导入错误。
5. 清理Spark临时缓存
Spark有时候会缓存旧的依赖信息,导致新JAR无法正常加载:
- 删除本地Spark临时目录(默认是
/tmp/spark-*开头的文件夹) - 关闭所有Spark相关进程后,重新提交任务
你可以先从--packages命令和文件名拼写这两点入手排查,这两个是我碰到过的最常见的原因。
内容的提问来源于stack exchange,提问作者Alex Sun

