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

如何在Databricks中用PySpark解析嵌套JSON列为key和value列

PySpark处理无固定Schema JSON列转换为目标表

核心思路

由于JSON字段无固定Schema,需先动态推导所有可能的字段结构,再将嵌套JSON扁平化,最后与原表其他列合并得到目标表。以下是可复用的实现代码:

步骤1:创建示例源表(模拟你的数据)

from pyspark.sql import SparkSession
from pyspark.sql.functions import from_json, schema_of_json, collect_set, concat_ws, col, explode, map_keys, map_values
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, MapType

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

# 模拟源表:3列(id, name, json_col),json_col包含不同结构的JSON
data = [
    (1, "Alice", '{"age": 25, "city": "New York"}'),
    (2, "Bob", '{"age": 30, "contact": {"phone": "123456", "email": "bob@test.com"}}'),
    (3, "Charlie", '{"hobbies": ["reading", "hiking"], "score": {"math": 90, "english": 85}}')
]

source_df = spark.createDataFrame(data, ["id", "name", "json_col"])
source_df.show(truncate=False)

步骤2:动态推导统一JSON Schema

收集所有JSON字符串,生成包含所有可能字段的Schema:

# 收集所有JSON字符串并拼接,用于推导Schema
all_jsons = source_df.select(concat_ws(",", collect_set("json_col")).alias("all_jsons")).first()["all_jsons"]
# 生成统一Schema
json_schema = schema_of_json(all_jsons)
# 解析JSON列为Struct类型
parsed_df = source_df.withColumn("parsed_json", from_json(col("json_col"), json_schema))
parsed_df.printSchema()

步骤3:扁平化嵌套Struct字段

编写递归函数,将嵌套的Struct/Map类型字段展开为一级列(列名用.分隔嵌套层级):

def flatten_df(df, prefix=""):
    fields = []
    for field in df.schema.fields:
        col_name = f"{prefix}.{field.name}" if prefix else field.name
        if isinstance(field.dataType, StructType):
            # 递归处理嵌套Struct
            fields.extend(flatten_df(df.select(field.name), col_name).columns)
        elif isinstance(field.dataType, MapType):
            # 处理Map类型,展开为key-value列
            keys = df.select(explode(map_keys(col(field.name)))).distinct().rdd.flatMap(lambda x: x).collect()
            for key in keys:
                fields.append(col(f"{field.name}.{key}").alias(f"{col_name}.{key}"))
        else:
            fields.append(col(col_name).alias(col_name))
    return df.select(fields)

# 扁平化parsed_json列
flattened_json_df = flatten_df(parsed_df.select("parsed_json"))
# 合并原表的id、name列与扁平化后的JSON列
target_df = parsed_df.select("id", "name").crossJoin(flattened_json_df)
target_df.show(truncate=False)

补充说明

  • 若JSON中存在数组类型,可根据业务需求添加explode逻辑展开数组行;
  • 动态Schema推导会自动包含所有出现过的字段,缺失字段会显示为null;
  • 扁平化后的列名保留嵌套层级(如contact.phone),可修改函数中的命名规则调整格式。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 12:05:18