如何将PySpark DataFrame按分区写入JSON且仅保留Payload内容
解决Spark分区写入JSON时仅输出payload列内容的问题
问题背景
按country_code、state_code、size列分区写入S3时,需要输出的JSON仅包含payload列的MapType内容,而非外层包裹payload键的结构。此前误以为文本格式方案不适用,实际通过文本格式存储JSON字符串可满足需求,且符合JSON Lines规范。
现有代码
( df .write .format('json') .partitionBy("country_code", "state_code", "size") .mode("append") .save('/mnt/dev/test') )
当前输出
{ "payload": { "100": { "cumulative_ttl_sold": 11, "cumulative_ttl_returned": 1 } } }
期望输出
{ "100": { "cumulative_ttl_sold": 11, "cumulative_ttl_returned": 1 } }
解决方案
使用to_json函数将payload列转为JSON字符串,再以文本格式写入,即可得到期望的输出结构。代码如下:
import pyspark.sql.functions as F ( df # 将payload列转换为JSON字符串 .select(F.to_json(F.col("payload")).alias("value")) .write .format("text") .partitionBy("country_code", "state_code", "size") .mode("append") .save('/mnt/dev/test') )
说明
- 此方法输出的文本文件每行是一个独立的JSON对象,符合JSON Lines格式,Spark、Pandas等多数数据处理工具均支持解析该格式。
- 若需生成单个包含所有对象的JSON数组,可通过后续合并文件并添加数组首尾符号实现,但分区场景下更推荐JSON Lines格式,便于并行处理。
内容的提问来源于stack exchange,提问作者satoshi
相关产品推荐
相关产品推荐

