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数组:
struct("id", "timestamp", "weight")将每行字段封装为结构体,确保字段顺序和目标JSON一致collect_list(...)把全表的结构体聚合为一个数组,实现多条记录合并为单条数组字段- 直接使用
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
相关产品推荐
相关产品推荐

