You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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]

优化说明

  1. 直接使用列名(如col("date"))替代案例类字段拼接字符串,代码更简洁不易出错
  2. 直接用小写字符串"y"、"n"替代"Y".toLowerCase(),减少不必要的方法调用

内容的提问来源于stack exchange,提问作者Saad

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.18 11:16:06