PySpark读取Kafka数据时createStream调用错误求助
问题分析
你遇到的核心问题是Spark Streaming与Kafka版本不兼容或依赖包加载失败,报错里的java.lang.ClassNotFoundException: kafka.common.TopicAndPartition是关键——说明Spark需要的Kafka类没有被正确加载,主要有两个原因:
- 你在代码里设置
PYSPARK_SUBMIT_ARGS的时机太晚,PySpark启动时已经读取了环境变量,导致指定的依赖包根本没被加载。 - 你用的是旧版
spark-streaming-kafka-0-8包,这个包仅兼容Kafka 0.8.x版本,如果你的Kafka服务器是0.9+版本,必然会出现类不匹配的问题。
解决方案
1. 修正依赖包的加载方式
不要在代码里设置环境变量,而是在运行脚本前通过命令行配置,或者直接用spark-submit提交时指定依赖:
方式一:运行前设置环境变量
export PYSPARK_SUBMIT_ARGS='--packages org.apache.spark:spark-streaming-kafka-0-8_2.10:2.2.1 pyspark-shell' python your_script.py
方式二:用spark-submit直接提交
spark-submit --packages org.apache.spark:spark-streaming-kafka-0-8_2.10:2.2.1 your_script.py
2. 适配Kafka版本(推荐)
如果你的Kafka是0.10.x及以上版本,建议直接换成新版的spark-streaming-kafka-0-10包——它支持更现代的Kafka API,兼容性更好,而且不需要依赖ZooKeeper连接:
提交命令
spark-submit --packages org.apache.spark:spark-streaming-kafka-0-10_2.10:2.2.1 your_script.py
对应的新版代码
import sys from pyspark import SparkContext from pyspark.streaming import StreamingContext from pyspark.streaming.kafka import KafkaUtils if __name__ == "__main__": sc = SparkContext(appName="PythonStreamingKafkaWordCount") ssc = StreamingContext(sc, 60) # 新版API用bootstrap.servers替代ZK地址 kafkaParams = {"bootstrap.servers": "localhost:9092", "group.id": "console-consumer-68081"} topics = ["near_line"] # 用createDirectStream替代旧的createStream kvs = KafkaUtils.createDirectStream(ssc, topics, kafkaParams) lines = kvs.map(lambda x: x[1]) counts = lines.flatMap(lambda line: line.split(" ")) \ .map(lambda word: (word, 1)) \ .reduceByKey(lambda a, b: a + b) counts.pprint() ssc.start() ssc.awaitTermination()
3. 基础验证步骤
- 确认ZooKeeper(旧版API)或Kafka Broker(新版API)正常运行,地址无误。
- 检查
near_line主题是否存在:- 旧版Kafka:
kafka-topics.sh --list --zookeeper localhost:2181 - 新版Kafka:
kafka-topics.sh --list --bootstrap-server localhost:9092
- 旧版Kafka:
- 确认Spark的Scala版本(你的包后缀是
_2.10)和实际使用的Scala版本一致,避免依赖包不兼容。
为什么原来的代码不行?
你在代码里设置os.environ['PYSPARK_SUBMIT_ARGS']是无效的——PySpark的JVM进程在你导入SparkContext之前就已经启动了,这时候设置环境变量根本无法影响依赖包的加载,必须在启动PySpark前配置好。
内容的提问来源于stack exchange,提问作者Nayana Madhu
相关产品推荐
相关产品推荐

