Python中Spark Streaming DataFrame转JSON格式时遇报错求助
问题根源
报错的直接原因是DataFrameWriter.json()方法必须传入输出路径参数,原代码里items.write.json()没传这个必填参数,触发了TypeError。另外,流处理的每个批次如果写入同一个路径,还会出现数据覆盖的问题,得一起解决。
修复方法
1. 最基础的修复:补上输出路径
直接给json()方法指定存储路径,比如DBFS或者云存储路径:
query = (items.writeStream.outputMode("append") .foreachBatch(lambda items, epoch_id: items.write.json("/dbfs/your/output/path/json_data")) .start())
2. 更合理的写法:按批次拆分路径
为了避免不同批次的数据互相覆盖,用批次ID(epoch_id)作为路径的一部分,每个批次生成独立的目录:
query = (items.writeStream.outputMode("append") .foreachBatch(lambda items, epoch_id: items.write.json(f"/dbfs/your/output/path/json_batch_{epoch_id}")) .start())
3. 完整修正后的代码
整合所有逻辑,加上追加模式配置(避免目录已存在报错):
from pyspark.sql import functions as F import boto3 import sys table_name = "dev.emp.master_events" df = ( spark.readStream.format("delta") .option("readChangeFeed", "true") .option("startingVersion", 2) .table(table_name) ) items = df.select('*') # 按批次写入独立路径,启用append模式防止目录存在报错 query = (items.writeStream.outputMode("append") .foreachBatch(lambda items, epoch_id: items.write.mode("append").json(f"/dbfs/your/target/path/batch_{epoch_id}")) .start()) # 可选:让程序等待流处理完成,避免脚本直接退出 query.awaitTermination()
重要提醒
json()方法的path参数是必填项,必须指定数据写入的分布式存储路径(比如DBFS、S3、ADLS等)- 流处理批次写入时,务必用epoch_id或时间戳区分路径,否则后续批次会覆盖之前的数据
- 如果要往同一个目录追加数据,必须设置
.mode("append"),Spark默认的error模式会因为目录已存在直接报错
内容的提问来源于stack exchange,提问作者Veerendra
相关产品推荐
相关产品推荐

