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上下文,很容易导致依赖加载异常。
修复步骤
- 对齐组件版本
执行以下命令把本地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版本对齐。
- 调整代码顺序,清理冗余逻辑
把环境变量设置、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()
- 离线包兜底方案
如果按上述步骤操作后还是无法加载依赖,直接手动下载对应版本的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
相关产品推荐
相关产品推荐

