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

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集成指南的部署要求配置应用依赖。

问题根因
  1. 依赖配置时机错误:PYSPARK_SUBMIT_ARGS环境变量必须在SparkSession初始化、Spark相关模块导入之前设置,原有代码先初始化SparkSession再配置依赖参数,Spark启动时根本不会加载Kafka连接器包,直接导致数据源找不到。
  2. 依赖版本不匹配:PySpark 2.4.7官方预编译版本基于Scala 2.11构建,原有代码配置的依赖包后缀为_2.13(对应Scala 2.13版本),和Spark运行环境不兼容,即便配置时机正确也无法正常加载。
  3. 数据结构不符合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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 23:33:21