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

Spark+Kafka流处理CSV数据时生产消费代码报错如何解决?

生产者KafkaTimeoutError错误解决方案

  • 第一步确认基础服务正常启动:先启动Zookeeper服务,再启动Kafka服务,确保两个服务无报错运行。对应启动命令如下(Windows环境将.sh替换为.bat即可):
    启动Zookeeper:bin/zookeeper-server-start.sh config/zookeeper.properties
    启动Kafka:bin/kafka-server-start.sh config/server.properties
  • 第二步确认Kafka主题已创建:运行命令创建代码中指定的songTopic主题:
    bin/kafka-topics.sh --create --topic songTopic --bootstrap-server localhost:9092 --replication-factor 1 --partitions 1
  • 第三步修正代码逻辑bug:你的生产者代码中仅读取了songs_df但循环遍历的song_list未定义,需在读取csv后补充转换逻辑,示例如下:
    # 示例将每行数据转为json字符串传输,避免序列化异常
    import json
    song_list = songs_df.apply(lambda x: json.dumps(x.to_dict()), axis=1).tolist()
    
  • 第四步验证Kafka端口可访问:确认本地9092端口未被防火墙拦截,可通过telnet命令测试连通性:telnet localhost 9092

消费者找不到kafka数据源错误解决方案

该报错是因为PySpark默认未携带Structured Streaming与Kafka的集成依赖包,直接修改SparkSession初始化代码即可,需添加依赖包配置:

  • 首先查询你本地安装的PySpark版本:运行pip show pyspark获取版本号,比如返回Version为3.3.0
  • 修改SparkSession构建逻辑,添加对应版本的依赖配置,示例如下(注意将末尾的3.3.0替换为你自己的PySpark版本):
    spark = SparkSession \
            .builder \
            .appName("Reccomendation System") \
            .master("local[*]") \
            .config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0") \
            .getOrCreate()
    
  • 若网络环境不佳无法在线拉取依赖,可提前下载对应版本的spark-sql-kafka jar包到本地,将配置项改为spark.jars指向本地jar包路径即可。

功能验证

问题修复后可在消费者代码末尾添加如下控制台输出逻辑,验证流数据是否正常接收:

query = songs_df.writeStream \
        .outputMode("append") \
        .format("console") \
        .start()
query.awaitTermination()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 05:54:05