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

PySpark如何将带键值对的JSON字典数组展开为独立列

Spark嵌套JSON结构DataFrame扁平化处理方案

你提供的源数据包含header嵌套结构体、properties.property嵌套数组两类嵌套结构,可通过以下两种语言的Spark API实现扁平化,得到你需要的结果:


Scala版本实现

import org.apache.spark.sql.functions._

// 读取JSON数据
val rawDf = spark.read.json("你的数据源路径")

// 展开header字段,同时炸开property数组
val explodeDf = rawDf.select(
  col("header.message-id").alias("message-id"),
  col("header.reply-to").alias("reply-to"),
  col("header.timestamp").alias("timestamp"),
  explode(col("properties.property")).alias("property")
)

// 行转列把property的name转为独立列
val flattenDf = explodeDf.groupBy("message-id", "reply-to", "timestamp")
  .pivot("property.name")
  .agg(first("property.value"))

// 输出结果
flattenDf.show(false)

PySpark版本实现

from pyspark.sql import functions as F

# 读取JSON数据
raw_df = spark.read.json("你的数据源路径")

// 展开header字段+炸开property数组
explode_df = raw_df.select(
    F.col("header.message-id").alias("message-id"),
    F.col("header.reply-to").alias("reply-to"),
    F.col("header.timestamp").alias("timestamp"),
    F.explode(F.col("properties.property")).alias("property")
)

// 行转列生成最终扁平化表
flatten_df = explode_df.groupBy("message-id", "reply-to", "timestamp")\
    .pivot("property.name")\
    .agg(F.first("property.value"))

flatten_df.show(truncate=False)

实现逻辑说明

  1. 嵌套结构体展开:通过.运算符直接访问header下的子字段,重命名后即可生成一级列
  2. 数组炸开:使用explode函数将properties.property数组的每个元素拆分为单独的行
  3. 行转列:通过pivot函数将property下的name字段值转为独立列,用first聚合函数取对应的value值完成映射

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 03:15:01