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
相关产品推荐
相关产品推荐

