Spark处理空集合报错:如何使reduce操作在空集合时返回0
解决Spark空集合reduce报错的问题
这个问题我之前也碰到过!当你的filter条件没匹配到任何数据时,最终的RDD就是空的,直接调用reduce就会触发empty collection异常——因为reduce要求集合至少有一个元素才能执行累加逻辑。这里有几个靠谱的解决办法:
方法1:用RDD的fold替代reduce
fold允许你指定一个初始值,当集合为空时就直接返回这个初始值,完美适配你的需求:
val result = df .filter(col("channel_pk") === "abc") .groupBy("member_PK") .agg(sum(col("price") * col("quantityOrdered")) as "totalSum") .select("totalSum") .rdd.map(_(0).asInstanceOf[Double]) .fold(0.0)(_ + _) // 初始值设为0.0,空集合时直接返回0.0
原理很简单:fold先把初始值作为累加器起点,遍历RDD元素进行累加;如果RDD是空的,就直接返回初始值,不会抛出异常。
方法2:全程用DataFrame API处理,避免转RDD
其实完全不用转成RDD来求和,直接在DataFrame层面再做一次聚合更简洁,而且Spark的DataFrame API本身就处理了空集合的情况:
val result = df .filter(col("channel_pk") === "abc") .groupBy("member_PK") .agg(sum(col("price") * col("quantityOrdered")) as "totalSum") .agg(coalesce(sum("totalSum"), lit(0.0))) // 空集合时sum返回null,用coalesce替换为0.0 .head()(0).asInstanceOf[Double]
这里coalesce(sum("totalSum"), lit(0.0))的作用是:如果前面的分组聚合没有产生任何数据,sum("totalSum")会返回null,此时就用0.0替代。最后head()可以安全获取唯一的结果,因为最后一次聚合后只会有一行数据。
方法3:先判断RDD是否为空(不推荐,性能损耗)
如果你一定要保留RDD的处理方式,可以先判断RDD是否为空,再决定执行reduce还是返回0:
val rdd = df .filter(col("channel_pk") === "abc") .groupBy("member_PK") .agg(sum(col("price") * col("quantityOrdered")) as "totalSum") .select("totalSum") .rdd.map(_(0).asInstanceOf[Double]) val result = if (rdd.isEmpty()) 0.0 else rdd.reduce(_ + _)
不过要注意:isEmpty()是一个Action操作,会触发一次Spark job;之后的reduce又会触发一次,相当于做了两次计算,性能不如前两种方法,除非特殊场景,不建议使用。
个人最推荐方法2,全程用DataFrame API不仅代码更简洁,Spark还能做更多的查询优化,运行效率也更高。
内容的提问来源于stack exchange,提问作者Markus
相关产品推荐
相关产品推荐

