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

