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

如何用PySpark/Python动态展平JSON,仅提取key2的嵌套键值对

使用PySpark/Python展平JSON的指定嵌套节点(仅处理key2)

我们需要将动态结构的JSON中指定的key2部分展平为三列:ID、以点连接的嵌套键路径key、对应的值value,JSON可能存在多层嵌套结构。

示例输入JSON

{
    "key1": {
        "subkey1":"1.1",
        "subkey2":"1.2"
        },
    "key2": {
        "subkey1":"2.1",
        "subkey2":"2.2",
        "subkey3": {"child3": { "subchild3":"2.3.3.3" } }
        },
    "key3": {
        "subkey1":"3.1",
        "subkey2":"3.2"
        }
}

预期输出

IDkeyvalue
1key2.subkey12.1
2key2.subkey22.2
3key2.subkey3.child3.subchild32.3.3.3

方法一:Python原生实现

通过递归函数遍历key2的嵌套结构,收集完整键路径和对应值,最后转换为DataFrame并添加ID列。

import json
import pandas as pd

def flatten_json(nested_json, parent_key='', sep='.'):
    items = []
    for k, v in nested_json.items():
        new_key = f"{parent_key}{sep}{k}" if parent_key else k
        if isinstance(v, dict):
            items.extend(flatten_json(v, new_key, sep=sep).items())
        else:
            items.append((new_key, v))
    return dict(items)

# 加载示例JSON数据
with open('input.json', 'r') as f:
    data = json.load(f)

# 仅展平key2部分
flattened_data = flatten_json(data['key2'], parent_key='key2')

# 转换为DataFrame并添加ID列
df = pd.DataFrame(list(flattened_data.items()), columns=['key', 'value'])
df['ID'] = df.index + 1

# 调整列顺序
df = df[['ID', 'key', 'value']]
print(df)

运行输出:

ID                            key     value
0   1                  key2.subkey1       2.1
1   2                  key2.subkey2       2.2
2   3  key2.subkey3.child3.subchild3  2.3.3.3

方法二:PySpark实现

利用PySpark的递归列处理和stack函数展平嵌套结构,仅针对key2部分操作。

from pyspark.sql import SparkSession
from pyspark.sql.functions import col
from pyspark.sql.types import StructType

def flatten_nested_cols(df, parent_key='', sep='.'):
    fields = df.schema.fields
    for field in fields:
        col_name = f"{parent_key}{sep}{field.name}" if parent_key else field.name
        if isinstance(field.dataType, StructType):
            nested_df = df.select(field.name + ".*")
            flattened_df = flatten_nested_cols(nested_df, col_name, sep)
            df = df.drop(field.name).join(flattened_df, how='cross')
        else:
            df = df.withColumnRenamed(field.name, col_name)
    return df

# 初始化SparkSession
spark = SparkSession.builder.appName("FlattenJSON").getOrCreate()

# 示例JSON数据
json_data = '''
{
    "key1": {
        "subkey1":"1.1",
        "subkey2":"1.2"
        },
    "key2": {
        "subkey1":"2.1",
        "subkey2":"2.2",
        "subkey3": {"child3": { "subchild3":"2.3.3.3" } }
        },
    "key3": {
        "subkey1":"3.1",
        "subkey2":"3.2"
        }
}
'''

# 创建初始DataFrame
df = spark.read.json(spark.sparkContext.parallelize([json_data]))

# 提取key2部分并展平
key2_df = df.select("key2.*")
flattened_key2 = flatten_nested_cols(key2_df, parent_key='key2')

# 将宽表转换为长表
key_columns = flattened_key2.columns
stack_expr = f"stack({len(key_columns)}, {', '.join([f'{repr(k)}, {k}' for k in key_columns])}) as (key, value)"
long_df = flattened_key2.selectExpr(stack_expr)

# 添加ID列
final_df = long_df.withColumn("ID", col("key").rdd.zipWithIndex().map(lambda x: x[1]+1).toDF().withColumnRenamed("_1", "ID"))

# 调整列顺序并显示结果
final_df.select("ID", "key", "value").show(truncate=False)

运行输出:

+---+--------------------------------+---------+
|ID |key                             |value    |
+---+--------------------------------+---------+
|1  |key2.subkey1                    |2.1      |
|2  |key2.subkey2                    |2.2      |
|3  |key2.subkey3.child3.subchild3   |2.3.3.3  |
+---+--------------------------------+---------+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 10:10:31