PySpark中groupby分组后同时计算sum与countDistinct的实现方法
PySpark单次聚合同时实现sum和countDistinct的解决方案
问题根因
- PySpark的
agg方法传入字典参数时,仅支持简单聚合函数的字符串名称,countDistinct没有对应的直接字符串映射,需要用count(DISTINCT 列名)的SQL表达式格式替代。 - 单独使用去重计数表达式列表报错,是因为没有对列表做解包操作,
agg接收的是多个独立的Column参数,直接传列表会被识别为单个非Column类型参数,触发AssertionError: all exprs should be Column报错。
方案1:统一使用Column表达式列表(推荐)
该方式类型安全,也方便自定义聚合后的列名:
# 首先导入需要的聚合函数 from pyspark.sql.functions import sum, countDistinct sum_cols = ['a', 'b'] count_cols = ['id'] # 构造sum类聚合表达式,指定别名避免列名混乱 sum_exprs = [sum(x).alias(f"sum_{x}") for x in sum_cols] # 构造去重计数类聚合表达式 count_exprs = [countDistinct(x).alias(f"distinct_{x}") for x in count_cols] # 解包两个表达式列表,一次性传入agg方法 df_aggregated = df.groupBy("month", "product").agg(*sum_exprs, *count_exprs)
方案2:兼容字典传参写法
如果习惯用字典定义聚合规则,可以调整去重计数的表达式格式:
sum_cols = ['a', 'b'] count_cols = ['id'] exprs1 = {x: "sum" for x in sum_cols} # 用SQL格式的count(DISTINCT)替代countDistinct字符串 exprs2 = {x: f"count(DISTINCT {x})" for x in count_cols} exprs = {**exprs1, **exprs2} df_aggregated = df.groupBy("month", "product").agg(exprs)
该方式生成的列名默认是sum(a)、count(DISTINCT id)格式,可后续通过withColumnRenamed调整别名。
内容的提问来源于stack exchange,提问作者TiTo
相关产品推荐
相关产品推荐

