Spark GroupBy聚合函数使用问题:ABC转CustomClass报错
Spark分组聚合报错:expression 'flag'不在group by或聚合函数中
问题描述
我在将ABC案例类转换为CustomClass时遇到问题,需求如下:
CustomClass的count:按a、b分组,且日期满足1年过滤条件的总行数t30Ycount和t30Ncount:按a、b分组后,分别满足30天日期过滤且flag为Y、30天日期过滤且flag为N的行数
执行代码时出现报错:
[scalatest] org.apache.spark.sql.AnalysisException: expression 'flag' is neither present in the group by, nor is it an aggregate function. Add to group by or wrap in first() (or first_value) if you don't care which value you get.;
相关代码:
case class ABC(a: String , b: Long, flag: String, date: Timestamp) extends Product {} info .filter(col(s"${abc.date}") > oneYear) .groupBy(col(s"${abc.a}"), col(s"${abc.b}")) .agg( // 分组后的总行数 count("a").as(s"${customClass.count}"), // 分组内flag为Y且日期符合30天过滤的行数 when(lower(col(s"${abc.flag}")) === "Y".toLowerCase && col(s"${abc.date}") > thirtyDays, count("a")).otherwise(lit(0)).as(s"${customClass.t30Ycount}"), // 分组内flag为N且日期符合30天过滤的行数 when(lower(col(s"${abc.flag}")) === "N".toLowerCase && col(s"${abc.date}") > thirtyDays, count("a")).otherwise(lit(0)).as(s"${customClass.t30Ncount}") ).as[CustomClass] case class CustomClass(a: String , b: Long , count: Long , t30Ycount: Long , t30NCount: Long )
解决方案
错误原因
报错核心是:聚合操作(agg)中when直接引用了flag和date列,但这两个字段既不在groupBy分组键中,也未被聚合函数包裹。Spark要求分组聚合时,所有非分组键字段必须经过聚合函数处理。另外原写法逻辑错误:when里的count("a")是对整个分组的计数,而非符合条件行的计数。
修正后的代码
import org.apache.spark.sql.functions.{col, count, lower, lit, sum, when} case class ABC(a: String , b: Long, flag: String, date: Timestamp) extends Product {} case class CustomClass(a: String , b: Long , count: Long , t30Ycount: Long , t30NCount: Long ) info .filter(col("date") > oneYear) .groupBy(col("a"), col("b")) .agg( // 分组后的总行数 count("a").as("count"), // 统计分组内flag为Y且日期在30天内的行数 sum( when(lower(col("flag")) === "y" && col("date") > thirtyDays, lit(1)) .otherwise(lit(0)) ).as("t30Ycount"), // 统计分组内flag为N且日期在30天内的行数 sum( when(lower(col("flag")) === "n" && col("date") > thirtyDays, lit(1)) .otherwise(lit(0)) ).as("t30Ncount") ) .as[CustomClass]
优化说明
- 直接使用列名(如
col("date"))替代案例类字段拼接字符串,代码更简洁不易出错 - 直接用小写字符串
"y"、"n"替代"Y".toLowerCase(),减少不必要的方法调用
内容的提问来源于stack exchange,提问作者Saad
相关产品推荐
相关产品推荐

