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

PySpark多记录合并为单记录并导出指定格式JSON文件

高效实现Spark DataFrame转指定JSON格式

解决方案

采用Spark内置的分布式聚合函数,全程无需将数据加载到Driver内存,完美适配大数据场景,且避开groupBy的限制:

完整代码

from pyspark.sql.functions import struct, collect_list, lit

data = [('01-01-2022', 123, 'abc123'), ('02-02-2022', 456, 'def456'), ('03-03-2022', 789, 'ghi789')]
columns = ["timestamp", "weight", "id"]

df = spark.createDataFrame(data, columns)

# 转换为目标结构
df_convert = df.agg(
    collect_list(struct("id", "timestamp", "weight")).alias("summaries")  # 调整字段顺序匹配目标JSON
).withColumn("status", lit(200))

# 导出JSON
df_convert.write.format('json').mode("overwrite").save("MyDocuments/write_path")

方案细节说明

  • 生成summaries数组:

    1. struct("id", "timestamp", "weight")将每行字段封装为结构体,确保字段顺序和目标JSON一致
    2. collect_list(...)把全表的结构体聚合为一个数组,实现多条记录合并为单条数组字段
    3. 直接使用agg()不指定groupBy,默认对全表做全局聚合,完全避开重复记录的问题
  • 添加status常量列:
    lit(200)生成值为200的常量列,精准匹配需求中的status字段

  • 性能优势:
    全程基于Spark分布式引擎执行,不会像collect()那样把全量数据拉到Driver内存,处理大数据集时性能远超内存方案,且无额外shuffle开销(Spark会自动优化全局聚合的执行逻辑)

转换结果验证

转换后的DataFrame结构完全符合要求:

+--------------------------------------------------------+-------+
|                                               summaries| status|
+--------------------------------------------------------+-------+
|[{"id":"abc123","timestamp":"01-01-2022","weight":123},
{"id":"def456","timestamp":"02-02-2022","weight":456},
{"id":"ghi789","timestamp":"03-03-2022","weight":789}] |    200|
+--------------------------------------------------------+--------+

导出的JSON格式与需求完全一致。

内容的提问来源于stack exchange,提问作者Meeldurb

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 19:25:43