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

Spark 2.3解析动态键JSON数据转换为id、value列结构化DataFrame方法

核心处理思路

你的JSON结构属于顶层键为动态ID值的非标结构,只需三步即可得到目标结果:

  1. 将顶层的多列动态struct列转置为行格式
  2. 提取struct中内置的id和value字段
  3. 展开value数组,将数组内每个元素拆分为独立行

Scala 实现代码

import org.apache.spark.sql.functions.{col, explode, expr}

// 替换为你的JSON文件路径,加载原始数据
val rawDf = spark.read.json("/path/to/your/json/file")

// 动态生成stack表达式,适配任意数量的顶层ID列,无需手动写死列名
val stackExpr = expr(s"stack(${rawDf.columns.length}, ${rawDf.columns.map(c => s"`$c`").mkString(", ")}) as struct_val")

val resultDf = rawDf
  .select(stackExpr)
  // 提取struct中的id和value数组
  .select(
    col("struct_val.id").alias("id"),
    col("struct_val.value").alias("value_arr")
  )
  // 展开数组得到单行单value
  .withColumn("value", explode(col("value_arr")))
  .select("id", "value")

// 验证输出
resultDf.show()

PySpark 实现代码

from pyspark.sql.functions import col, explode, expr

# 替换为你的JSON文件路径,加载原始数据
raw_df = spark.read.json("/path/to/your/json/file")

# 动态生成stack表达式,适配任意数量的顶层ID列
stack_expr = expr(f"stack({len(raw_df.columns)}, {', '.join([f'`{c}`' for c in raw_df.columns])}) as struct_val")

result_df = raw_df \
    .select(stack_expr) \
    .select(
        col("struct_val.id").alias("id"),
        col("struct_val.value").alias("value_arr")
    ) \
    .withColumn("value", explode(col("value_arr"))) \
    .select("id", "value")

# 验证输出
result_df.show()

输出结果

运行后会得到和预期完全一致的DataFrame:

+----+-----+
|  id|value|
+----+-----+
|3200|  cat|
|3200|  dog|
|2000| bird|
|2500|kitty|
|3650|horse|
+----+-----+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 15:27:00