PySpark实现DataFrame转嵌套JSON并合并为单行列的需求
解决方案
假设你的原始PySpark DataFrame名为df,可以通过以下代码实现需求:
from pyspark.sql import functions as F # 按transaction_date分组,将每条记录的核心字段打包为结构体并收集成数组 grouped_df = df.groupBy("transaction_date").agg( F.collect_list( F.struct( F.col("id"), F.col("latest_quote"), F.col("policy_price") ) ).alias("trans_opp") ) # 将日期与嵌套数组组合为JSON格式的单列 result_df = grouped_df.select( F.to_json( F.struct( F.col("transaction_date"), F.col("trans_opp") ) ).alias("unique_column") ) # 输出结果 result_df.show(truncate=False)
代码说明:
- 分组构建嵌套数组:通过
groupBy按交易日期聚合,用struct把每条记录的id、报价日期、保单价格打包成结构体,再用collect_list将同日期的结构体收集为数组,对应目标格式中的trans_opp字段。 - 生成JSON单列:用
struct将交易日期和嵌套数组组合成完整结构,再通过to_json转换为JSON字符串,最后重命名为unique_column,得到仅含一行的结果DataFrame。
若原始数据包含多个不同的transaction_date,结果会按每个日期生成一行;如果只有单个日期(如示例数据),则结果仅保留一行,完全匹配需求格式。
内容的提问来源于stack exchange,提问作者daniel____
相关产品推荐
相关产品推荐

