如何使用Scala的reduceByKey计算RDD(String,List[Int])的List[Int]平均值
Scala Spark 使用reduceByKey计算RDD中List[Int]平均值方案
完整实现代码
import org.apache.spark.rdd.RDD import org.apache.spark.SparkContext import org.apache.spark.SparkConf object ListAvgCalculation { def main(args: Array[String]): Unit = { // 初始化Spark环境 val conf = new SparkConf().setAppName("AvgCalc").setMaster("local[*]") val sc = new SparkContext(conf) // 模拟题目给出的源数据集 RDD[(String, List[Int])] val sourceRdd: RDD[(String, List[Int])] = sc.parallelize(Seq( ("David", List(60,70,80)), ("John", List(70,80,90)) )) // 核心计算逻辑 val resultRdd: RDD[(String, Int)] = sourceRdd // 将每个List转换为(总得分, 元素个数)的二元组,不修改Key .mapValues(list => (list.sum, list.size)) // 按Key聚合,累加相同Key的总得分和元素总个数 .reduceByKey((prev, curr) => (prev._1 + curr._1, prev._2 + curr._2)) // 计算平均值,整数除法和题目示例输出格式匹配 .mapValues { case (totalScore, totalCount) => totalScore / totalCount } // 验证输出 resultRdd.foreach(println) // 输出结果: // (David,70) // (John,80) sc.stop() } }
逻辑说明
mapValues预处理:仅对Value(也就是List[Int])做转换,把每个List拆解成「所有元素总和、元素个数」的二元组,避免聚合阶段重复遍历List,提升执行效率reduceByKey聚合:针对相同Key的二元组,分别累加总和数值、元素个数数值,得到该Key对应的全量数据的总得分和总样本数- 结果计算:最后再次用
mapValues计算平均值,若需要保留小数位,可以将除法逻辑改为totalScore.toDouble / totalCount即可得到浮点型平均值
内容的提问来源于stack exchange,提问作者bigdata
相关产品推荐
相关产品推荐

