PySpark按id1、product分组按id2前缀统计distinct id2报错如何解决
报错原因
你编写的代码将sf.when()条件逻辑放在了聚合函数的外层,此时sf.col('id2')属于未被分组/聚合的原始列,不符合SQL分组聚合的语法要求:所有出现在agg()中的非分组字段,必须被聚合函数包裹。
单独写sf.countDistinct('id2')可以正常运行,是因为此时id2被聚合函数countDistinct包裹,满足语法规则。
正确实现方案
你需要将条件判断放在countDistinct内部,仅对符合条件的id2进行去重计数:sf.when()在条件不满足时会默认返回NULL,而countDistinct会自动忽略NULL值,刚好匹配你的需求。
代码示例
import pyspark.sql.functions as sf agg_data = ( df1 .groupby('id1', 'product') .agg( # 统计不带abc-前缀的去重id2数量 sf.countDistinct(sf.when(~sf.col('id2').like('abc-%'), sf.col('id2'))).alias('id2_count_without_abc'), # 统计带abc-前缀的去重id2数量 sf.countDistinct(sf.when(sf.col('id2').like('abc-%'), sf.col('id2'))).alias('id2_count_with_abc') ) )
样例输出
对应你提供的测试数据,运行后得到的结果如下:
| id1 | product | id2_count_without_abc | id2_count_with_abc |
|---|---|---|---|
| 1 | Upload | 2 | 1 |
| 2 | Upload | 0 | 1 |
| 3 | Upload | 2 | 1 |
内容的提问来源于stack exchange,提问作者cs_guy
相关产品推荐
相关产品推荐

