PySpark DataFrame写入JSON文件去success日志格式丢失问题
你出现列名和格式丢失的核心原因是collect()返回的是PySpark Row 对象列表,直接传入json_dumps()时默认只会序列化行的值序列,不会保留列名与值的映射关系,按照以下两个适配不同场景的方案操作即可解决:
方案1:小数据量场景,手动生成JSON文件(适合数据量小于Driver内存上限的情况)
先把每行Row对象转为字典保留列名映射,再做JSON序列化:
import json # 假设待输出的DataFrame名为df dict_rows = [row.asDict() for row in df.collect()] # 输出为整段JSON数组格式 with open("result.json", "w", encoding="utf-8") as f: json.dump(dict_rows, f, indent=2, ensure_ascii=False) # 如果需要和Spark原生输出一致的JSON Lines格式(每行一个独立JSON对象) with open("result.jsonl", "w", encoding="utf-8") as f: for row in dict_rows: f.write(json.dumps(row, ensure_ascii=False) + "\n")
方案2:大数据量场景,用Spark原生输出关闭冗余文件(避免Driver内存溢出,无需手动collect)
直接通过Spark配置关闭_SUCCESS日志文件和元数据文件生成,不用自行处理序列化逻辑:
# 关闭_SUCCESS文件生成 spark.conf.set("mapreduce.fileoutputcommitter.marksuccessfuljobs", "false") # 可选:关闭Parquet等格式的元数据摘要文件生成 spark.conf.set("parquet.enable.summary-metadata", "false") # 如果需要合并为单个JSON文件,可加coalesce(1),仅建议小数据量场景使用 # 直接输出JSON,mode可根据需求调整为append/overwrite等 df.write.mode("overwrite").json("/path/to/output_directory")
内容的提问来源于stack exchange,提问作者Dinesh Kumar Rajendran
相关产品推荐
相关产品推荐

