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

如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 06:35:48