You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark读取Kafka数据时createStream调用错误求助

问题分析

你遇到的核心问题是Spark Streaming与Kafka版本不兼容或依赖包加载失败,报错里的java.lang.ClassNotFoundException: kafka.common.TopicAndPartition是关键——说明Spark需要的Kafka类没有被正确加载,主要有两个原因:

  1. 你在代码里设置PYSPARK_SUBMIT_ARGS的时机太晚,PySpark启动时已经读取了环境变量,导致指定的依赖包根本没被加载。
  2. 你用的是旧版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
  • 确认Spark的Scala版本(你的包后缀是_2.10)和实际使用的Scala版本一致,避免依赖包不兼容。
为什么原来的代码不行?

你在代码里设置os.environ['PYSPARK_SUBMIT_ARGS']是无效的——PySpark的JVM进程在你导入SparkContext之前就已经启动了,这时候设置环境变量根本无法影响依赖包的加载,必须在启动PySpark前配置好。

内容的提问来源于stack exchange,提问作者Nayana Madhu

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 08:28:07