Spark处理Map结构:计算RDD[ROW]年度成绩总和并过滤缺失记录
没问题,我来帮你搞定这个需求!下面是基于Spark RDD操作的具体解决方案,完全贴合你的需求:
计算年度成绩总和的实现方案
核心思路是先过滤掉无有效年度成绩的记录,再提取有效成绩进行求和,下面分两种常见的Row结构场景给出代码示例:
1. 当Row包含命名字段(如字段名为annual_score)
如果你的RDD[Row]是从DataFrame转换而来,或者有明确的字段命名,可以通过字段名直接访问:
import org.apache.spark.sql.Row // 第一步:过滤掉年度成绩为null或不存在的记录 val validScoreRows = yourRowRDD.filter(row => { // 安全获取字段,判断是否存在有效数值 row.getAs[Option[Double]]("annual_score").isDefined }) // 第二步:提取成绩并计算总和 val totalAnnualScore = validScoreRows .map(row => row.getAs[Double]("annual_score")) .reduce(_ + _) // 输出结果 println(s"年度成绩总和:$totalAnnualScore")
2. 当Row按索引存储数据(如年度成绩在索引2的位置)
如果Row没有命名字段,仅按位置存储数据,可以通过索引访问并校验:
import org.apache.spark.sql.Row // 过滤掉成绩为null或非数值类型的无效记录 val validScoreRows = yourRowRDD.filter(row => { !row.isNullAt(2) && row.get(2).isInstanceOf[Double] }) // 提取成绩并求和 val totalAnnualScore = validScoreRows .map(row => row.getDouble(2)) .reduce(_ + _) println(s"年度成绩总和:$totalAnnualScore")
额外扩展:按年度分组求和
如果需要按不同年度分别计算总和,可以调整为键值对操作:
// 假设Row包含"year"和"annual_score"两个字段 val groupedTotalScores = yourRowRDD .filter(row => row.getAs[Option[Double]]("annual_score").isDefined) .map(row => (row.getAs[String]("year"), row.getAs[Double]("annual_score"))) .reduceByKey(_ + _) // 遍历输出每个年度的成绩总和 groupedTotalScores.foreach { case (year, total) => println(s"$year年度成绩总和:$total") }
提示:如果你的年度成绩是整数类型,只需将代码中的
Double替换为Int即可。
内容的提问来源于stack exchange,提问作者Khumar
相关产品推荐
相关产品推荐

