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

PySpark读取Kafka报Failed to find data source: kafka错误求解

问题根因

你遇到这个报错是三个核心问题导致的:

  • 版本不匹配:本地安装的pyspark==3.1.2和部署的Spark==3.2.1版本不一致,你配置的Kafka依赖包是适配Spark3.2.1的版本,和本地pyspark版本不兼容,无法正常加载。
  • 环境变量生效顺序错误:PYSPARK_SUBMIT_ARGS环境变量必须在导入任何pyspark模块、初始化Spark上下文之前设置才会生效,你的代码里先导入了pyspark相关对象才设置该变量,等于Kafka依赖配置根本没被Spark识别。
  • 上下文初始化逻辑冲突:你同时混用了旧版DStream API的StreamingContext和新版Structured Streaming的readStream接口,还重复初始化Spark上下文,很容易导致依赖加载异常。
修复步骤
  1. 对齐组件版本
    执行以下命令把本地pyspark版本升级到和集群一致的3.2.1:
pip uninstall pyspark -y
pip install pyspark==3.2.1

注意:Kafka连接器的Scala版本、Spark版本必须和你部署的Spark、Scala版本完全对应,你当前Scala版本为2.12、Spark版本为3.2.1,使用spark-sql-kafka-0-10_2.12:3.2.1是正确的,前提是pyspark版本对齐。

  1. 调整代码顺序,清理冗余逻辑
    把环境变量设置、findspark初始化放到所有pyspark导入之前,删除重复的上下文初始化代码,不需要使用旧版StreamingContext(Structured Streaming接口不需要依赖该对象),修复后的可运行参考代码如下:
import os
# 所有pyspark相关导入前必须完成环境变量配置,末尾必须加pyspark-shell参数
os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.2.1 pyspark-shell'
import findspark
findspark.init()

from pyspark.sql import SparkSession

def ReadingDataToKafka():
    spark = SparkSession.builder \
        .appName("KafkaWordCount") \
        .getOrCreate()
    spark.sparkContext.setLogLevel("ERROR")
    checkpoint_path = "file:///tmp/ZHYCargeProject"

    # 读取Kafka流数据
    df = spark.readStream \
        .format("kafka") \
        .option("kafka.bootstrap.servers", "bigdataweb01:9092") \
        .option("subscribe", "sex") \
        .load()

    # 以下补充你的业务处理逻辑,示例为控制台输出测试
    query = df.writeStream \
        .outputMode("append") \
        .format("console") \
        .option("checkpointLocation", checkpoint_path) \
        .start()
    query.awaitTermination()

if __name__ == "__main__":
    ReadingDataToKafka()
  1. 离线包兜底方案
    如果按上述步骤操作后还是无法加载依赖,直接手动下载对应版本的spark-sql-kafka-0-10_2.12-3.2.1.jar、kafka-clients、spark-token-provider-kafka-0-10_2.12-3.2.1.jar、commons-pool2-2.11.1.jar四个包,放到本地pyspark安装目录的jars文件夹和集群Spark的jars目录下,提交任务时去掉--packages参数即可直接识别Kafka数据源。
验证说明

第一次运行时控制台会输出依赖包拉取日志,等待拉取完成后,如果能正常消费到Kafka数据、不再抛出找不到数据源的报错,即为修复成功。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 20:57:21