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
相关产品推荐
相关产品推荐

