Spark 2.3解析动态键JSON数据转换为id、value列结构化DataFrame方法
核心处理思路
你的JSON结构属于顶层键为动态ID值的非标结构,只需三步即可得到目标结果:
- 将顶层的多列动态struct列转置为行格式
- 提取struct中内置的
id和value字段 - 展开
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
相关产品推荐
相关产品推荐

