Spark新手求助:不使用Spark SQL计算RDD单列值的总和
解决Spark RDD单列数值求和问题(无需Spark SQL)
嘿,我来帮你搞定这个问题!先理清你之前尝试的方法哪里出问题了,再给你几个高效安全的分布式聚合方案:
先说说你之前的错误原因:
- 第一个方法
val rdd2 = rdd1.map(agg(sum())):agg和sum()(来自Spark SQL的functions)都是DataFrame/DataSet专属API,RDD根本没有这些方法;而且map是用来逐个处理元素的,不是做全局聚合的,所以肯定报错。 - 第二个方法
val rdd2 = rdd1.map(c => sum(col(c))):同样踩了API混淆的坑——col()是Spark SQL的列操作方法,RDD里的元素是原始数值(比如Int/Long),不是DataFrame的列对象;另外map也不适合用来做聚合操作。 - 第三个方法
val rdd2 = rdd1.collect().sum():首先Scala里数组的sum是不带括号的(正确写法是rdd1.collect().sum),其次这个方法非常不推荐——collect()会把RDD所有数据拉到Driver节点内存里,数据量大的话直接就内存溢出了,完全浪费了Spark分布式计算的优势。
正确的解决方案(分布式聚合,高效安全)
假设你的RDD是数值类型(比如RDD[Int]或RDD[Long]),直接用RDD原生的聚合方法就行:
方法1:用reduce函数(最基础的RDD聚合方式)
reduce会在分布式节点上两两合并元素,最终得到全局总和:
// 完整写法 val totalSum = rdd1.reduce((a, b) => a + b) // 简洁的下划线写法 val totalSum = rdd1.reduce(_ + _)
方法2:直接用sum()便捷方法
Spark为数值类型RDD封装了现成的sum()方法,本质是简化版的reduce操作:
val totalSum = rdd1.sum()
额外注意:如果你的RDD元素是字符串类型
如果原始RDD里的数值是以字符串形式存储的(比如从文本文件读取的),需要先转成数值类型再聚合:
// 假设原始RDD是字符串类型 val stringRdd = sc.parallelize(Array("5000", "6000", "7000", "8000", "9000")) // 转成Int类型 val numericRdd = stringRdd.map(_.toInt) // 求和 val totalSum = numericRdd.reduce(_ + _)
内容的提问来源于stack exchange,提问作者Sankar
相关产品推荐
相关产品推荐

