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

PySpark如何基于行值添加动态前缀生成嵌套JSON

PySpark按ID聚合生成带动态前缀的JSON列解决方案

问题描述

现有PySpark DataFrame结构如下:

idName_typeNameCar
1FirstrobNissan
2FirstjoeHyundai
1LastdentInfiniti
2LastKentGenesis

需求:按id分组,将Name_type的值作为前缀添加到Name、Car列名(如First_Name),最终生成包含所有对应键值对的JSON列,期望输出如下:

idjson_column
1{"First_Name":"rob","First_Car":"Nissan","Last_Name":"dent","Last_Car":"Infiniti"}
2{"First_Name":"joe","First_Car":"Hyundai","Last_Name":"Kent","Last_Car":"Genesis"}

目前已实现单条记录生成JSON列的代码,但无法完成动态前缀添加与分组聚合:

column_set = ['Name','Car']
df = df.withColumn("json_data", to_json(struct([df[x] for x in column_set])))

解决代码及说明

完整实现代码

from pyspark.sql import functions as F
from pyspark.sql.types import StringType

# 定义需要处理的字段列表
column_set = ['Name', 'Car']

# 生成带Name_type前缀的键值对Map
df = df.withColumn(
    "key_value_map",
    F.create_map(
        *[F.concat_ws("_", F.col("Name_type"), F.lit(col)).cast(StringType()), F.col(col) 
          for col in column_set]
    )
)

# 按id分组合并Map并转成JSON
result_df = df.groupBy("id").agg(
    F.to_json(F.map_concat(F.collect_list("key_value_map"))).alias("json_column")
)

# 查看结果
result_df.show(truncate=False)

逻辑说明

  1. 生成带前缀的键值对:使用create_map为每个目标字段生成键值对,键通过concat_ws拼接Name_type与字段名(如First_Name),值为字段对应的数据;
  2. 分组合并Map:按id分组后,用collect_list收集该id下所有键值对Map,再通过map_concat合并为一个完整的Map;
  3. 转JSON列:最后用to_json将合并后的Map转为JSON格式的字符串列,得到目标结果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 22:40:24