PySpark如何基于行值添加动态前缀生成嵌套JSON
PySpark按ID聚合生成带动态前缀的JSON列解决方案
问题描述
现有PySpark DataFrame结构如下:
| id | Name_type | Name | Car |
|---|---|---|---|
| 1 | First | rob | Nissan |
| 2 | First | joe | Hyundai |
| 1 | Last | dent | Infiniti |
| 2 | Last | Kent | Genesis |
需求:按id分组,将Name_type的值作为前缀添加到Name、Car列名(如First_Name),最终生成包含所有对应键值对的JSON列,期望输出如下:
| id | json_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)
逻辑说明
- 生成带前缀的键值对:使用
create_map为每个目标字段生成键值对,键通过concat_ws拼接Name_type与字段名(如First_Name),值为字段对应的数据; - 分组合并Map:按
id分组后,用collect_list收集该id下所有键值对Map,再通过map_concat合并为一个完整的Map; - 转JSON列:最后用
to_json将合并后的Map转为JSON格式的字符串列,得到目标结果。
内容的提问来源于stack exchange,提问作者sadrobot
相关产品推荐
相关产品推荐

