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

如何在PySpark DataFrame中展开并透视类JSON结构

处理Google Analytics原始Event数据,提取event_params并转换为宽表结构

步骤说明

  1. 展开event_params数组:用explode把数组类型的event_params拆分成单行单键值对的形式,方便后续处理。
  2. 提取有效value值:GA的value结构里四个字段只有一个非空,用coalesce合并这四个字段,得到统一的param_value列。
  3. 过滤空值:移除param_value为空的记录,避免无效字段干扰。
  4. 透视生成宽表:将key字段转成列名,param_value作为对应列的值,还原成原始行结构。

完整PySpark代码示例

from pyspark.sql import functions as F
from pyspark.sql.types import *

# 假设你的原始DataFrame名为ga_events_df
# 先确认event_params的结构(如果结构不符合,可能需要先转换,比如从字符串解析成Struct/Array)
# 示例结构:ArrayType(StructType([
#    StructField("key", StringType()),
#    StructField("value", StructType([
#        StructField("string_value", StringType()),
#        StructField("int_value", IntegerType()),
#        StructField("float_value", FloatType()),
#        StructField("double_value", DoubleType())
#    ]))
# ]))

# 步骤1:展开数组
exploded_df = ga_events_df.select(
    "event_date", "event_timestamp", "event_name",
    F.explode("event_params").alias("param")
)

# 步骤2:提取key和有效value
parsed_df = exploded_df.select(
    "event_date", "event_timestamp", "event_name",
    F.col("param.key").alias("param_key"),
    F.coalesce(
        F.col("param.value.string_value"),
        F.col("param.value.int_value").cast(StringType()),
        F.col("param.value.float_value").cast(StringType()),
        F.col("param.value.double_value").cast(StringType())
    ).alias("param_value")
)

# 步骤3:过滤空值
filtered_df = parsed_df.filter(F.col("param_value").isNotNull())

# 步骤4:透视生成目标宽表
final_df = filtered_df.groupBy(
    "event_date", "event_timestamp", "event_name"
).pivot("param_key").agg(F.first("param_value"))

# 查看结果
final_df.show()

关键细节说明

  • coalesce的类型转换:因为不同类型的value(int/string等)需要统一类型才能合并,所以把数值类型转成StringType,后续如果需要可以再转回去。
  • 聚合函数选择:用first是因为每个原始行的每个key只会出现一次,确保取到唯一值;如果存在重复key,可根据需求换成max或min。
  • 字符串格式解析:如果你的event_params是字符串格式(不是原生的Array/Struct),需要先解析:
# 定义event_params的Schema
param_schema = ArrayType(StructType([
    StructField("key", StringType()),
    StructField("value", StructType([
        StructField("string_value", StringType()),
        StructField("int_value", IntegerType()),
        StructField("float_value", FloatType()),
        StructField("double_value", DoubleType())
    ]))
]))

# 解析字符串格式的event_params
ga_events_df = ga_events_df.withColumn(
    "event_params", F.from_json(F.col("event_params"), param_schema)
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 11:35:18