Spark DataFrame嵌套数组值提取:拆分k、value为列求助
Spark DataFrame 嵌套JSON扁平化实现方案
假设你的输入DataFrame结构如下(以Scala为例,Python逻辑完全一致):
// 示例输入数据 val inputDF = spark.createDataFrame(Seq( (1, """[{"t_id": "t1", "properties": [{"k": "name", "value": "Alice"}, {"k": "age", "value": "30"}]}, {"t_id": "t2", "properties": [{"k": "city", "value": "New York"}]}]"""), (2, """[{"t_id": "t3", "properties": [{"k": "name", "value": "Bob"}, {"k": "gender", "value": "Male"}]}]""") )).toDF("id", "info")
步骤1:解析JSON列并展开t对象数组
先把字符串类型的info列解析为数组结构,再用explode将每个t对象拆分为单独行,保留原id关联:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ // 定义t对象的Schema val tSchema = ArrayType(StructType(Seq( StructField("t_id", StringType), StructField("properties", ArrayType(StructType(Seq( StructField("k", StringType), StructField("value", StringType) ))) ))) // 解析info列并展开t对象 val explodedTDF = inputDF .withColumn("info", from_json(col("info"), tSchema)) .select(col("id"), explode(col("info")).alias("t")) .select(col("id"), col("t.t_id"), col("t.properties"))
步骤2:展开properties数组
将每个t对象下的properties数组拆分为单独的k-v行:
val explodedPropsDF = explodedTDF .select(col("id"), col("t_id"), explode(col("properties")).alias("prop")) .select(col("id"), col("t_id"), col("prop.k"), col("prop.value"))
步骤3:Pivot转置k为列名
通过pivot将k值转为列名,聚合value得到最终扁平化结果:
val flattenedDF = explodedPropsDF .groupBy("id", "t_id") .pivot("k") .agg(first("value"))
Python版本实现
如果使用PySpark,逻辑完全一致,代码如下:
from pyspark.sql import functions as F from pyspark.sql.types import * # 示例输入数据 inputDF = spark.createDataFrame([ (1, """[{"t_id": "t1", "properties": [{"k": "name", "value": "Alice"}, {"k": "age", "value": "30"}]}, {"t_id": "t2", "properties": [{"k": "city", "value": "New York"}]}]"""), (2, """[{"t_id": "t3", "properties": [{"k": "name", "value": "Bob"}, {"k": "gender", "value": "Male"}]}]""") ], ["id", "info"]) # 定义Schema t_schema = ArrayType(StructType([ StructField("t_id", StringType()), StructField("properties", ArrayType(StructType([ StructField("k", StringType()), StructField("value", StringType()) ]))) ])) # 解析并展开t对象 exploded_t_df = inputDF \ .withColumn("info", F.from_json(F.col("info"), t_schema)) \ .select("id", F.explode("info").alias("t")) \ .select("id", "t.t_id", "t.properties") # 展开properties exploded_props_df = exploded_t_df \ .select("id", "t_id", F.explode("properties").alias("prop")) \ .select("id", "t_id", "prop.k", "prop.value") # Pivot得到扁平化结果 flattened_df = exploded_props_df \ .groupBy("id", "t_id") \ .pivot("k") \ .agg(F.first("value")) flattened_df.show()
关键注意事项
- 如果
value字段存在多数据类型,建议先统一转为字符串再执行pivot,避免类型冲突 - 若同一
id+t_id下存在相同k的多个值,可根据业务需求替换first为collect_list/concat_ws等聚合函数 - 如果
info本身就是数组类型(而非JSON字符串),直接跳过from_json步骤执行explode即可
内容的提问来源于stack exchange,提问作者Yousuf Sultan
相关产品推荐
相关产品推荐

