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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 12:52:49