PySpark 2.4.7发送DataFrame到Kafka主题报找不到数据源错误如何解决
环境版本
- PySpark:2.4.7
- Kafka:2.13_3.2.0
需求说明
实现PySpark读取本地CSV文件,将读取到的DataFrame数据作为生产者发送到指定Kafka主题。参考公开资料编写代码后运行失败,抛出找不到Kafka数据源的错误,需要可运行的代码实现与对应配置说明。
原有问题代码
import findspark findspark.init("/usr/local/spark") from pyspark.sql import SparkSession from pyspark.streaming.kafka import KafkaUtils from pyspark.sql.functions import * import os from kafka import KafkaProducer import csv def spark_session(): ''' Description: To open a spark session. Returns a spark session object. ''' spark = SparkSession \ .builder \ .appName("Test_Kafka_Producer") \ .master("local[*]") \ .getOrCreate() return spark if __name__ == '__main__': spark = spark_session() topic = "Kafkatest" spark_version = '2.4.7' os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages org.apache.spark:spark-sql-kafka-0-10_2.13:{}'.format(spark_version) #producer = KafkaProducer(bootstrap_servers=['localhost:9092'], #value_serializer= lambda x: x.encode('utf-8')) df1 = spark.read.csv("annual-enterprise-survey-2020-financial-year-provisional-size-bands-csv.csv", inferSchema = True, header = True) df1.show(10) print("sending df===========") df1.write \ .format("kafka") \ .option("kafka.bootstrap.servers", "localhost:9092") \ .option("topic", topic) \ .save() print("End------")
运行报错信息
原报错内容:
py4j.protocol.Py4JJavaError: An error occurred while calling o41.save. : org.apache.spark.sql.AnalysisException: Failed to find data source: kafka. Please deploy the application as per the deployment section of "Structured Streaming + Kafka Integration Guide".;
报错翻译:调用save方法时抛出Py4JJavaError,根因是Spark SQL分析异常:无法识别kafka数据源,需按照Structured Streaming与Kafka集成指南的部署要求配置应用依赖。
问题根因
- 依赖配置时机错误:
PYSPARK_SUBMIT_ARGS环境变量必须在SparkSession初始化、Spark相关模块导入之前设置,原有代码先初始化SparkSession再配置依赖参数,Spark启动时根本不会加载Kafka连接器包,直接导致数据源找不到。 - 依赖版本不匹配:PySpark 2.4.7官方预编译版本基于Scala 2.11构建,原有代码配置的依赖包后缀为
_2.13(对应Scala 2.13版本),和Spark运行环境不兼容,即便配置时机正确也无法正常加载。 - 数据结构不符合Kafka写入要求:Spark的Kafka数据源写入时强制要求DataFrame包含名为
value的字段(存储消息内容,可选key字段存储消息键),原有代码直接将原始CSV结构的DataFrame写入,依赖问题修复后也会抛出字段缺失错误。
修复后可运行代码
# 所有环境变量配置必须放在Spark导入、SparkSession初始化之前 import os import findspark findspark.init("/usr/local/spark") # 配置Kafka连接器依赖,注意Scala版本为2.11,匹配Spark 2.4.x系列版本 spark_version = "2.4.7" os.environ['PYSPARK_SUBMIT_ARGS'] = f'--packages org.apache.spark:spark-sql-kafka-0-10_2.11:{spark_version} pyspark-shell' from pyspark.sql import SparkSession from pyspark.sql.functions import to_json, struct, col def init_spark_session(): spark = SparkSession.builder \ .appName("Test_Kafka_Producer") \ .master("local[*]") \ .getOrCreate() # 调整日志级别,屏蔽无关INFO日志 spark.sparkContext.setLogLevel("WARN") return spark if __name__ == '__main__': spark = init_spark_session() target_topic = "Kafkatest" kafka_addr = "localhost:9092" # 读取CSV文件 csv_df = spark.read.csv( "annual-enterprise-survey-2020-financial-year-provisional-size-bands-csv.csv", inferSchema=True, header=True ) csv_df.show(10, truncate=False) print("开始向Kafka写入数据===========") # 转换数据结构:将整行数据序列化为JSON格式,存入必填的value字段 kafka_write_df = csv_df.select( to_json(struct([col(c) for c in csv_df.columns])).alias("value") ) # 写入Kafka kafka_write_df.write \ .format("kafka") \ .option("kafka.bootstrap.servers", kafka_addr) \ .option("topic", target_topic) \ .save() print("数据写入完成------") spark.stop()
运行注意事项
- 首次运行时Spark会自动从中央仓库拉取对应版本的Kafka连接器依赖,需要保证运行环境网络通畅;离线环境需提前下载对应jar包放到Spark安装目录的jars文件夹下。
- 如果代码内配置环境变量仍然加载依赖失败,可以直接在spark-submit提交命令中指定依赖包,不需要在代码内设置
PYSPARK_SUBMIT_ARGS,提交命令示例:
spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.11:2.4.7 你的脚本文件名.py
- 若需要给Kafka消息指定key实现分区路由,只需要在转换kafka_write_df时增加key字段即可,示例:
col("表中某列名").alias("key")。
内容的提问来源于stack exchange,提问作者subh
相关产品推荐
相关产品推荐

