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)
实现逻辑说明
- 嵌套结构体展开:通过
.运算符直接访问header下的子字段,重命名后即可生成一级列 - 数组炸开:使用
explode函数将properties.property数组的每个元素拆分为单独的行 - 行转列:通过
pivot函数将property下的name字段值转为独立列,用first聚合函数取对应的value值完成映射
内容的提问来源于stack exchange,提问作者Sanjay Gaddam
相关产品推荐
相关产品推荐

