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
相关产品推荐
相关产品推荐

