PySpark基于UserID合并行、按Property值动态生成新列的实现方法
PySpark 用户属性长表转宽表实现方案
需求说明
业务侧开发用户表单用于追踪单个用户会话的特定属性及关联值,原始数据为长表结构,需要转换为每个UserID对应单行的宽表结构,转换规则如下:
- 以
Property字段的所有取值作为宽表的新列名 - 列值填充规则:对应记录的
isDisabled字段为False时填充对应Value值,为True时填充空值 - 最终结果中每个UserID仅保留一行数据
原始数据初始化
测试数据及DataFrame初始化代码如下:
data =[ [100,"Prop1","Value1",False], [100,"Prop2","Value2",False], [200,"Prop1","Value3",False], [200,"Prop2","Value4",False], [200,"Prop3","Value5",True], [300,"Prop3","Value6",False], [300,"Prop4","Value7",False] ] columns = ["UserID","Property","Value","isDisabled"] df_in = spark.createDataFrame(data,columns)
实现代码
实现逻辑分两步:首先预处理属性值,将禁用状态的属性值置空;再通过pivot算子完成行转列聚合。
from pyspark.sql import functions as F # 步骤1:预处理值字段,禁用属性置为null df_mid = df_in.withColumn( "valid_value", F.when(~F.col("isDisabled"), F.col("Value")).otherwise(None) ) # 步骤2:按UserID分组,透视Property字段生成宽表 # 若需要固定列顺序、提升性能,可在pivot中传入指定列列表,例:.pivot("Property", ["Prop1","Prop2","Prop3","Prop4"]) df_res = df_mid.groupBy("UserID")\ .pivot("Property")\ .agg(F.first("valid_value", ignorenulls=True)) # 输出结果验证 df_res.show()
运行结果
执行后输出结果如下,完全符合预期:
+------+------+------+------+------+ |UserID|Prop1 |Prop2 |Prop3 |Prop4 | +------+------+------+------+------+ |100 |Value1|Value2|null |null | |200 |Value3|Value4|null |null | |300 |null |null |Value6|Value7| +------+------+------+------+------+
补充说明
- 当前代码默认同一个UserID下同一个Property仅存在一条记录,因此使用
first聚合即可满足需求,性能最优 - 若同用户同属性存在多条有效记录,可根据业务规则替换聚合逻辑,例如取最新值、多值拼接等
- 生产环境建议在
pivot方法中传入明确的列名列表,避免Spark全表扫描枚举Property值,大幅提升运行效率
内容的提问来源于stack exchange,提问作者josephstx
相关产品推荐
相关产品推荐

