如何在PySpark DataFrame中分组统计单列唯一值并生成新列
PySpark统计客户购买手机品牌数量实现方案
错误原因说明
PySpark调用groupBy后返回的是GroupedData对象,和Pandas的GroupBy对象API设计完全不同:
- 不能直接通过
.列名的方式访问分组内的列 - 不能直接调用
select方法,必须先通过agg执行聚合得到普通DataFrame后才能操作列
实现方案
方案1:和Pandas逻辑对齐(先聚合再关联)
和你原有Pandas的merge逻辑完全等价:
from pyspark.sql import functions as F # 按用户msisdn分组,统计唯一手机品牌数量 brand_count_df = fd_subsprofile.groupBy("msisdn") \ .agg(F.countDistinct("handset_brand").alias("hpbrand_change_num")) # 左关联回原表 firstvalue = firstvalue.join(brand_count_df, on="msisdn", how="left")
方案2:无需关联(符合你不用merge的要求)
直接用窗口函数在原表上新增统计列,不需要额外关联操作,性能更优:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义按用户msisdn分区的窗口 msisdn_window = Window.partitionBy("msisdn") # 直接新增统计列 fd_subsprofile = fd_subsprofile.withColumn( "hpbrand_change_num", F.countDistinct("handset_brand").over(msisdn_window) )
内容的提问来源于stack exchange,提问作者Nabih Bawazir
相关产品推荐
相关产品推荐

