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

如何用Spark的from_json解析任意JSON并实现数据拆分?

解决Spark中解析含随机键的JSON列问题(无需UDF)

我完全理解你的困扰——当JSON列的顶层键是随机字符串时,没法提前定义StructType,直接用MapType又会被from_json拒绝,用UDF还怕拖慢性能。别担心,我们可以用Spark原生函数完美解决这个问题,分两种版本适配不同Spark版本:

方法一:Spark 3.1+(推荐,更简洁)

Spark 3.1及以上提供了json_object_keys函数,可以直接提取JSON对象的所有键,再配合transform逐个获取对应的值:

步骤1:提取JSON中的所有键并转换为结构化数组

from pyspark.sql import functions as F

# 假设你已经读入原始数据到DataFrame df
# df = spark.read.csv("example.csv", header=True)

# 提取JSON对象的所有随机键
df = df.withColumn("keys", F.json_object_keys(F.col("value")))

# 遍历每个键,提取对应的name和profession,生成Struct数组
df = df.withColumn(
    "people",
    F.transform(
        F.col("keys"),
        lambda key: F.struct(
            F.get_json_object(F.col("value"), f"$.{key}.name").alias("name"),
            F.get_json_object(F.col("value"), f"$.{key}.profession").alias("profession")
        )
    )
)

步骤2:生成合并数组的结果

result_agg = df.select(
    "ix",
    F.transform(F.col("people"), lambda x: x.name).alias("names"),
    F.transform(F.col("people"), lambda x: x.profession).alias("professions")
)

result_agg.show(truncate=False)

输出和你期望的第一版一致:

+---+------------+-------------------+
|ix |names       |professions        |
+---+------------+-------------------+
|1  |[bob]       |[engineer]         |
|2  |[sarah,matt]|[scientist,doctor] |
+---+------------+-------------------+

步骤3:拆分成每行一个人的结果

result_flattened = df.select(
    "ix",
    F.explode(F.col("people")).alias("person")
).select(
    "ix",
    "person.name",
    "person.profession"
)

result_flattened.show()

输出就是你要的第二版拆分结果:

+---+-----+----------+
|ix |name |profession|
+---+-----+----------+
|1  |bob  |engineer  |
|2  |sarah|scientist |
|2  |matt |doctor    |
+---+-----+----------+

方法二:Spark 2.4+(兼容低版本)

如果你的Spark版本低于3.1,没有json_object_keys,我们可以用正则表达式把外层JSON对象转换成键值对数组,再用from_json解析:

步骤1:转换JSON格式为数组结构

from pyspark.sql import functions as F
from pyspark.sql.types import ArrayType, StructType, StructField, StringType

# 把原始JSON对象转换成键值对数组的JSON字符串
df = df.withColumn(
    "value_array_str",
    # 把外层{}替换为[]
    F.regexp_replace(F.col("value"), r'^\{(.*)\}$', r'[$1]')
).withColumn(
    "value_array_str",
    # 把每个"key":value格式替换为{"key":"key","value":value}
    F.regexp_replace(F.col("value_array_str"), r'\"([^\"]+)\":', r'{"key":"$1","value":')
)

步骤2:解析转换后的JSON数组

# 定义数组的Schema
array_schema = ArrayType(StructType([
    StructField("key", StringType()),
    StructField("value", StructType([
        StructField("name", StringType()),
        StructField("profession", StringType())
    ]))
]))

df = df.withColumn("people_array", F.from_json(F.col("value_array_str"), array_schema))

步骤3:生成聚合和拆分结果

和方法一的步骤2、3完全一致:

# 聚合结果
result_agg = df.select(
    "ix",
    F.transform(F.col("people_array.value"), lambda x: x.name).alias("names"),
    F.transform(F.col("people_array.value"), lambda x: x.profession).alias("professions")
)

# 拆分结果
result_flattened = df.select(
    "ix",
    F.explode(F.col("people_array.value")).alias("person")
).select(
    "ix",
    "person.name",
    "person.profession"
)

为什么这个方法比UDF好?

这些操作都是Spark的原生高阶函数和内置函数,会被优化器转换成高效的执行计划,完全避免了UDF带来的序列化/反序列化开销,性能和原生Spark操作一致。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 09:47:35