Spark DataFrame多行JSON字符串解析:统一通用Schema解决方案
解决Spark DataFrame中JSON字段结构不统一的Schema解析问题
我在Spark DataFrame的payload列中存储了结构不固定的JSON字符串,示例如下:
df = spark.createDataFrame( [ ["""{"header":{"title":"ABC123","name":"test"},"results":{"data_A":[{"key":"value"},{"key":"value"}]}, "data_B" : {"a":"1"}}"""], ["""{"header":{"title":"ABC123","name":"test"},"results":{"data_A":{"key":"value"}}, "data_B" : null}"""], ["""{"header":{"title":"ABC123","name":"test"},"results":{"data_A":[{"key":"value"},{"key":"value"}]}, "data_B" : [{"a":"1"}, {"a":"2"}]}"""] ], ['payload'] )
可以看到data_A和data_B字段的数据结构不固定:有时是多值数组,有时是单个对象,有时为null。尝试过常规的统一Schema方案,但会因为无法将单个值存入ArrayType而产生null值,无法满足需求。
需要将payload列解析为具有通用Schema的Struct类型,目标Schema如下:
payload:struct ---data_B:array ------element:struct ---------a:string ---header:struct ------name:string ------title:string ---results:struct ------data_A:array ---------element:struct ------------key:string ------elt:array ---------element:struct ------------key:string ---------------test:struct ------------------test2:struct ---------------------elt:array ------------------------element:struct ---------------------------A:long ---------------------------B:long ------test:struct ---------a:struct ------------elt:array ---------------element:array ------------------element:struct ---------------------ab:long ---------b:string
核心思路
要解决结构不统一的问题,需先将JSON解析为动态结构,再对data_A、data_B这类字段做类型归一化处理:将单个对象包装为数组,null可转换为空数组(或保留null,按需调整),最后转换为目标Schema。
步骤1:定义目标Schema
将目标Schema转换为Spark的StructType定义:
from pyspark.sql.types import ( StructType, StructField, StringType, ArrayType, LongType ) # 定义嵌套结构 data_b_element_schema = StructType([StructField("a", StringType())]) data_a_element_schema = StructType([StructField("key", StringType())]) results_elt_test_test2_elt_element_schema = StructType([ StructField("A", LongType()), StructField("B", LongType()) ]) results_elt_test_test2_schema = StructType([ StructField("elt", ArrayType(results_elt_test_test2_elt_element_schema)) ]) results_elt_element_schema = StructType([ StructField("key", StringType()), StructField("test", results_elt_test_test2_schema) ]) results_test_a_element_element_schema = StructType([StructField("ab", LongType())]) results_test_a_element_schema = ArrayType(results_test_a_element_element_schema) results_test_a_schema = StructType([ StructField("elt", ArrayType(results_test_a_element_schema)) ]) results_test_schema = StructType([ StructField("a", results_test_a_schema), StructField("b", StringType()) ]) results_schema = StructType([ StructField("data_A", ArrayType(data_a_element_schema)), StructField("elt", ArrayType(results_elt_element_schema)), StructField("test", results_test_schema) ]) header_schema = StructType([ StructField("title", StringType()), StructField("name", StringType()) ]) target_schema = StructType([ StructField("header", header_schema), StructField("results", results_schema), StructField("data_B", ArrayType(data_b_element_schema)) ])
步骤2:完整解析与归一化代码
from pyspark.sql import functions as F # 1. 解析payload为动态JSON对象 parsed_df = df.withColumn("parsed_payload", F.from_json(F.col("payload"), StructType([]))) # 2. 对data_A做类型归一化:数组保留,单对象转数组,null转空数组 normalized_df = parsed_df.withColumn( "normalized_data_A", F.when( F.isarray(F.col("parsed_payload.results.data_A")), F.col("parsed_payload.results.data_A") ).when( F.col("parsed_payload.results.data_A").isNotNull(), F.array(F.col("parsed_payload.results.data_A")) ).otherwise(F.array().cast(ArrayType(data_a_element_schema))) ) # 3. 对data_B做同样的归一化处理 normalized_df = normalized_df.withColumn( "normalized_data_B", F.when( F.isarray(F.col("parsed_payload.data_B")), F.col("parsed_payload.data_B") ).when( F.col("parsed_payload.data_B").isNotNull(), F.array(F.col("parsed_payload.data_B")) ).otherwise(F.array().cast(ArrayType(data_b_element_schema))) ) # 4. 组装成目标Struct结构并转换为指定Schema final_df = normalized_df.withColumn( "payload", F.struct( F.col("parsed_payload.header").cast(header_schema), F.struct( F.col("normalized_data_A").cast(ArrayType(data_a_element_schema)), F.col("parsed_payload.results.elt").cast(ArrayType(results_elt_element_schema)), F.col("parsed_payload.results.test").cast(results_test_schema) ).alias("results"), F.col("normalized_data_B").cast(ArrayType(data_b_element_schema)) ).cast(target_schema) ).select("payload") # 验证结果 final_df.printSchema() final_df.show(truncate=False)
说明
- 对于
results中的elt、test等可选字段,原JSON中不存在时会自动填充null,符合目标Schema定义 - 若需保留
data_B的null而非转为空数组,可将otherwise(F.array().cast(...))改为otherwise(F.lit(None).cast(...)) - 该方案通过先解析动态结构再归一化类型,避免了固定Schema解析时的类型不匹配问题
内容的提问来源于stack exchange,提问作者pooltje
相关产品推荐
相关产品推荐

