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

Spark中列平均值计算:新手求助基于RDD求各科目平均分

计算每个科目的平均分数(Spark RDD实现)

嗨,作为Spark新手能想到用RDD来处理这个成绩统计需求,非常棒!我来一步步带你实现这个功能,用Python和Scala两种常用的Spark语言都给你示例,你可以根据自己的环境选择。

核心思路

要计算每个科目的平均分,我们需要:

  • 把原始数据转换成**(科目,分数)**的键值对RDD
  • 对每个科目,统计总分和参加人数
  • 用总分除以人数得到平均分

Python 实现步骤

1. 初始化Spark上下文并创建RDD

首先我们需要把你的原始数据加载成RDD:

from pyspark import SparkContext

# 本地模式初始化Spark上下文(生产环境根据集群配置调整)
sc = SparkContext("local", "SubjectAverageCalculator")

# 你的原始成绩数据
raw_data = ["Mathew,science,45", "Mathew,history,55", "Mark,maths,23", "Mark,science,76"]
# 并行化创建RDD
score_rdd = sc.parallelize(raw_data)

2. 转换为(科目,分数)键值对

把每行字符串按逗号分割,提取科目和分数(注意要把分数从字符串转成整数,方便后续计算):

# 分割每行数据,取第2个元素(科目)和第3个元素(分数)
subject_score_pairs = score_rdd.map(lambda line: (line.split(",")[1], int(line.split(",")[2])))

3. 计算每个科目的总分与人数,再求平均

我们用mapValues把每个分数转换成(分数,1)的元组(1代表该科目的一个计数),然后用reduceByKey累加同一个科目的总分和人数,最后计算平均分:

# 转换为(科目,(分数,1)),再累加总分和人数,最后计算平均
subject_average = subject_score_pairs.mapValues(lambda s: (s, 1)) \
                                    .reduceByKey(lambda a, b: (a[0] + b[0], a[1] + b[1])) \
                                    .mapValues(lambda total_count: total_count[0] / total_count[1])

# 输出结果
print("每个科目的平均分:", subject_average.collect())

执行后你会得到结果:[('science', 60.5), ('history', 55.0), ('maths', 23.0)]


Scala 实现步骤

如果你用Scala开发,代码逻辑是一样的:

import org.apache.spark.SparkContext
import org.apache.spark.SparkConf

object SubjectAverage {
  def main(args: Array[String]): Unit = {
    val conf = new SparkConf().setAppName("SubjectAverage").setMaster("local")
    val sc = new SparkContext(conf)

    val rawData = Array("Mathew,science,45", "Mathew,history,55", "Mark,maths,23", "Mark,science,76")
    val scoreRDD = sc.parallelize(rawData)

    val subjectScorePairs = scoreRDD.map(line => {
      val parts = line.split(",")
      (parts(1), parts(2).toInt)
    })

    val subjectAverage = subjectScorePairs.mapValues(s => (s, 1))
      .reduceByKey((a, b) => (a._1 + b._1, a._2 + b._2))
      .mapValues(totalCount => totalCount._1.toDouble / totalCount._2)

    println("每个科目的平均分:" + subjectAverage.collect().mkString(", "))
    sc.stop()
  }
}

小提示

  • 如果你的数据是从文件读取的,只需要把sc.parallelize(raw_data)换成sc.textFile("你的文件路径")即可
  • 生产环境中记得不要用local模式,要根据你的Spark集群配置设置SparkContext
  • 除了reduceByKey,你也可以用combineByKey来实现,逻辑是类似的,不过reduceByKey对于这个场景更简洁

内容的提问来源于stack exchange,提问作者iamdeep

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 06:12:32