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
相关产品推荐
相关产品推荐

