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

如何将嵌套JSON转换为指定格式的PySpark DataFrame?

嵌套JSON转PySpark DataFrame格式优化

问题场景

需要将以下嵌套JSON转换为PySpark DataFrame,当前代码输出为数组列形式,无法得到预期的行展开结果。

原始JSON

{
"key1": 0.75,
"values":[
    {
        "id": 2313,
        "val1": 350,
        "val2": 6000
    },
    {
        "id": 2477,
        "val1": 340,
        "val2": 6500
    }
]
}

当前代码

import json
from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("CreateDataFrame").getOrCreate()

json_string = json.dumps({
    "key1": 0.75,
    "values":[
        {
            "id": 2313,
            "val1": 350,
            "val2": 6000
        },
        {
            "id": 2477,
            "val1": 340,
            "val2": 6500
        }
    ]
})
df = spark.read.json(spark.sparkContext.parallelize([json_string]))

df = df.select("key1", "values.id", "values.val1", "values.val2")
df.show()

当前输出

+----+-------------+-------------+-------------+
|key1|           id|         val1|         val2|
+----+-------------+-------------+-------------+
|0.75| [2313, 2477]|   [350, 340]| [6000, 6500]|
+----+-------------+-------------+-------------+

预期输出

+----+----+----+----+
|key1|  id|val1|val2|
+----+----+----+----+
|0.75|2313| 350|6000|
|0.75|2477| 340|6500|
+----+----+----+----+

解决方案

问题核心是values为数组类型,直接选取字段会生成数组列。需要用explode函数拆分数组为多行,再提取结构体字段。

优化后代码

import json
from pyspark.sql import SparkSession
from pyspark.sql.functions import explode, col

spark = SparkSession.builder.appName("CreateDataFrame").getOrCreate()

json_string = json.dumps({
    "key1": 0.75,
    "values":[
        {
            "id": 2313,
            "val1": 350,
            "val2": 6000
        },
        {
            "id": 2477,
            "val1": 340,
            "val2": 6500
        }
    ]
})
df = spark.read.json(spark.sparkContext.parallelize([json_string]))

# 拆分values数组为多行,保留key1
df_exploded = df.select("key1", explode(col("values")).alias("value"))
# 从结构体中提取各字段
df_result = df_exploded.select("key1", col("value.id"), col("value.val1"), col("value.val2"))
df_result.show()

关键说明

  • explode(col("values")):将数组中的每个元素拆分为独立行,key1的值会自动复制到每一行。
  • alias("value"):给拆分后的结构体列命名,便于后续字段提取。
  • col("value.id"):从结构体列中提取指定字段,最终得到扁平化的DataFrame。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 14:43:17