如何在PySpark DataFrame中展开并透视类JSON结构
处理Google Analytics原始Event数据,提取event_params并转换为宽表结构
步骤说明
- 展开event_params数组:用
explode把数组类型的event_params拆分成单行单键值对的形式,方便后续处理。 - 提取有效value值:GA的
value结构里四个字段只有一个非空,用coalesce合并这四个字段,得到统一的param_value列。 - 过滤空值:移除
param_value为空的记录,避免无效字段干扰。 - 透视生成宽表:将
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
相关产品推荐
相关产品推荐

