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

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

相关产品推荐
方舟 Agent Plan

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

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