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

Pyspark如何将固定嵌套结构的JSON解析为结构化DataFrame

PySpark 多层嵌套JSON快速扁平化解决方案

第一步:自动推导JSON Schema(解决自定义Schema报错问题)

无需手动编写完整Schema,直接让Spark读取JSON数据时自动推导结构,适配层级固定的场景:

# 若你的JSON数据已经是HTTP请求返回的字符串格式RDD
json_rdd = sc.parallelize([your_http_response_json_str])
# 读取时自动推导Schema,multiline参数根据JSON是否为多行格式调整
df = spark.read.option("multiline", "true").json(json_rdd)
# 打印推导后的Schema确认结构正确性
df.printSchema()

第二步:通用递归扁平化函数(支持结构体嵌套+数组字段展开)

该函数可自动将所有嵌套结构体拍平为下划线连接的扁平列,同时自动炸开数组字段为多行,无需手动指定每一层字段:

from pyspark.sql.functions import col, explode_outer
from pyspark.sql.types import StructType, ArrayType

def flatten_df(nested_df):
    stack = [((), nested_df)]
    while stack:
        parents, df = stack.pop()
        flat_cols = []
        nested_cols = []
        array_cols = []
        
        for col_name, col_type in df.dtypes:
            # 结构体类型暂存待后续展开
            if col_type.startswith('struct'):
                nested_cols.append(col_name)
            # 数组类型先炸开再继续展开
            elif col_type.startswith('array'):
                array_cols.append(col_name)
            # 普通类型直接保留,列名用下划线拼接父级路径
            else:
                flat_cols.append(col(*parents, col_name).alias("_".join(parents + (col_name,))))
        
        # 炸开数组字段,若不需要保留数组为null的行可替换为explode
        for array_col in array_cols:
            df = df.withColumn(array_col, explode_outer(col(array_col)))
        
        # 生成当前层扁平列结果
        current_flat_df = df.select(flat_cols) if flat_cols else None
        
        # 嵌套结构体压入栈递归处理
        for nested_col in nested_cols:
            projected_df = df.select(nested_col + ".*")
            stack.append((parents + (nested_col,), projected_df))
        
        # 拼接所有扁平列
        if current_flat_df is not None:
            if 'result_df' not in locals():
                result_df = current_flat_df
            else:
                result_df = result_df.crossJoin(current_flat_df)
    return result_df

使用示例

# 代入自动推导Schema后的嵌套DF直接调用即可
flat_df = flatten_df(df)
# 查看扁平化后的结果
flat_df.show()

可选调整方案

  • 如果不需要炸开数组字段、仅需展开结构体嵌套,删除函数中处理array_cols的代码块即可
  • 列名的拼接规则可以自行修改,比如把"_".join替换为其他分隔符
  • 若数组内仍为多层嵌套结构,函数会自动递归处理无需额外配置

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 12:21:02