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

PySpark实现列名映射为key、列值映射为value的JSON结构转换

问题描述

现有样例数据,执行指定PySpark代码后,输出结构与预期不符:

原执行代码

from pyspark.sql.functions import *

df_result1 = df_data.groupBy(col("Name").alias("Name"), col("Country").alias("Country"), col("Property").alias("Property")).agg(
    (collect_list(struct(col("No"), col("Place"))).alias("details")),
)

df_result2 = df_result1.groupBy(col("Name").alias("Name"), col("Country").alias("Country")).agg(collect_list(struct(col("Property"), col("details"))).alias("Main"))
final = df_result2.toJSON().collect()

实际输出

{
  "Name": "David",
  "Country": "Dubai",
  "Main": [
    {
      "Property": "House",
      "details": [
        {
          "No": "1",
          "Place": "JLT"
        }
      ]
    }
  ]
}

期望输出

{
  "Name": "David",
  "Country": "Dubai",
  "Main": [
    {
      "Property": "House",
      "details": [
        {
          "key": "No",
          "value": "1"
        },
        {
          "key": "Place",
          "value": "JLT"
        }
      ]
    }
  ]
}

需修改代码将details中的对象转换为key-value格式的数组。


解决方案

核心是调整生成details的逻辑,将原有的键值对对象转为key-value结构体数组,以下是两种实现方式:

方法1:使用Map相关函数(Spark 3.0+)

利用create_map和map_entries快速生成目标结构,代码更简洁:

from pyspark.sql.functions import *

# 分组并生成key-value格式的details数组
df_result1 = df_data.groupBy("Name", "Country", "Property").agg(
    collect_list(
        map_entries(create_map(lit("No"), col("No"), lit("Place"), col("Place")))
    ).alias("details")
)

# 保持原有上层分组逻辑
df_result2 = df_result1.groupBy("Name", "Country").agg(
    collect_list(struct("Property", "details")).alias("Main")
)

final = df_result2.toJSON().collect()

逻辑说明

  • create_map(lit("No"), col("No"), lit("Place"), col("Place")):将当前行的No、Place字段转为Map类型,键为字段名,值为字段对应内容。
  • map_entries(...):将Map转换为key-value结构体数组,每个结构体包含key和value字段,完全匹配期望格式。

方法2:手动构造结构体数组(兼容低版本Spark)

如果Spark版本低于3.0,不支持map_entries,可以直接构造目标结构体数组:

from pyspark.sql.functions import *

df_result1 = df_data.groupBy("Name", "Country", "Property").agg(
    collect_list(
        array(
            struct(lit("No").alias("key"), col("No").alias("value")),
            struct(lit("Place").alias("key"), col("Place").alias("value"))
        )
    ).alias("details")
)

df_result2 = df_result1.groupBy("Name", "Country").agg(
    collect_list(struct("Property", "details")).alias("Main")
)

final = df_result2.toJSON().collect()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 06:49:52