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

Spark Java如何过滤Month为空行后聚合Units与Dollars?

Spark Java中跳过Month为null的行进行Rollup聚合实现方案

你的核心需求是在Rollup聚合Units和Dollars时,排除Month字段为null的行,当前结果不符合预期是因为这些null行被计入了聚合计算。以下两种方案可以解决问题:

方案一:聚合前过滤无效行

在执行Rollup操作前,先过滤掉Month为null的行,后续聚合将仅处理有效数据,这是性能最优的方式。修改代码如下:

import static org.apache.spark.sql.functions.col;

input.filter(col("Month").isNotNull())
    .rollup(JavaConversions.asScalaBuffer(allSelectedCols).seq())
    .agg(ratingCols.get(0),
         JavaConverters.asScalaIteratorConverter(ratingCols.subList(1, ratingCols.size()).iterator())
             .asScala().toSeq())
    .sort(JavaConversions.asScalaBuffer(orderedCols).seq());

方案二:条件聚合(保留原数据集结构)

如果需要保留原数据集的所有行,仅在聚合时忽略Month为null的行,可以在sum函数中添加条件判断,只计算有效行的数值。修改ratingCols的定义:

方式1:使用SQL表达式

ratingCols.add(expr("sum(CASE WHEN Month IS NOT NULL THEN Units ELSE 0 END)").as("Units"));
ratingCols.add(
    expr("CAST(ROUND(SUM(CASE WHEN Month IS NOT NULL THEN CAST(Dollars AS DOUBLE) ELSE 0 END),3)  as decimal(36,3))").as("Dollars"));

方式2:使用Spark内置when函数(更简洁)

import static org.apache.spark.sql.functions.when;
import static org.apache.spark.sql.functions.sum;

ratingCols.add(sum(when(col("Month").isNotNull(), col("Units")).otherwise(0)).as("Units"));
ratingCols.add(
    expr("CAST(ROUND(SUM(WHEN Month IS NOT NULL THEN CAST(Dollars AS DOUBLE) ELSE 0 END),3)  as decimal(36,3))").as("Dollars"));

两种方案对比:

  • 方案一直接过滤无效行,减少后续计算的数据量,性能更优;
  • 方案二保留原数据集的完整结构,适合需要基于全量数据做其他操作的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 09:33:30