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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.23 20:45:03