如何用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" } }
预期输出
| ID | key | value |
|---|---|---|
| 1 | key2.subkey1 | 2.1 |
| 2 | key2.subkey2 | 2.2 |
| 3 | key2.subkey3.child3.subchild3 | 2.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
相关产品推荐
相关产品推荐

