如何从含无固定Schema JSON数组的Spark DataFrame列生成新DataFrame
无固定Schema的JSON数组列解析方案
问题背景
- 需求:从现有DataFrame的某一列生成新DataFrame,该列的值为多个JSON数组
- 痛点:JSON无固定Schema,无法使用需要预先指定Schema的
from_json函数直接解析
输入示例
| Column A | Column B |
|---|---|
| 1 | [{"id":"123","phone":"124"}] |
| 3 | [{"id":"456","phone":"741"}] |
期望输出
| id | phone |
|---|---|
| 123 | 124 |
| 456 | 741 |
解决思路
思路1:扁平化数组+按路径提取字段
先将JSON数组展开为单个JSON字符串,再用get_json_object按JSON路径提取字段,无需预先定义复杂Schema。
以PySpark为例:
from pyspark.sql import functions as F # 1. 将JSON数组转成字符串数组并展开,得到单个JSON对象字符串 df_flatten = df.withColumn( "json_obj", F.explode(F.from_json(F.col("Column B"), F.array_type(F.string_type()))) ) # 2. 按JSON路径提取目标字段 result_df = df_flatten.select( F.get_json_object(F.col("json_obj"), "$.id").alias("id"), F.get_json_object(F.col("json_obj"), "$.phone").alias("phone") ) result_df.show()
该方案适合明确知道需要提取的字段,但JSON整体Schema不固定的场景,灵活且性能较好。
思路2:自动推导Schema解析
如果JSON结构大部分一致,可从样本数据中自动推导Schema,再进行解析。
from pyspark.sql import functions as F # 从非空样本中提取单个JSON对象,自动推导Schema sample_json = df.filter(F.col("Column B").isNotNull()) \ .select(F.col("Column B")).first()[0][0] derived_schema = F.schema_of_json(sample_json) # 解析JSON数组并展开为行 df_parsed = df.withColumn( "json_arr", F.from_json(F.col("Column B"), F.array_type(derived_schema)) ) result_df = df_parsed.select(F.explode("json_arr").alias("data")).select("data.*") result_df.show()
此方案无需手动编写Schema,自动适配样本中的字段结构,个别差异字段会被解析为null,适合结构相对统一的场景。
思路3:RDD原生解析(极端无固定Schema场景)
如果JSON结构完全混乱、字段差异极大,可转成RDD后用Python原生JSON库解析,动态生成列。
import json from pyspark.sql import Row # 转换为RDD,解析每个JSON数组中的对象 rdd = df.rdd.flatMap( lambda row: [json.loads(obj) for obj in json.loads(row["Column B"])] ) # 收集所有可能出现的字段 all_fields = set() for item in rdd.collect(): all_fields.update(item.keys()) # 动态生成DataFrame result_df = rdd.map( lambda x: Row(**{field: x.get(field) for field in all_fields}) ).toDF() result_df.show()
该方案完全不依赖Schema,能适配任意结构的JSON,但性能较低,仅适合小数据量或结构极度不固定的场景。
内容的提问来源于stack exchange,提问作者Praveen
相关产品推荐
相关产品推荐

