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

Spark Scala中基于RDD计算全班学生物理成绩平均分

Calculate Average Physics Score Using Spark Scala RDD (No DataFrames)

Got it, let's work through this problem step by step. You need to compute the average physics score for the given student dataset using only RDDs in Spark Scala—no DataFrames allowed. Here's a straightforward, efficient way to do it:

Step-by-Step Explanation & Code

First, we'll set up our Spark environment, process the raw data, compute the average, and format the output as requested.

import org.apache.spark.sql.SparkSession

object PhysicsAverageCalculator {
  def main(args: Array[String]): Unit = {
    // Initialize SparkSession (modern entry point for Spark)
    val spark = SparkSession.builder()
      .appName("PhysicsAvgRDD")
      .master("local[*]") // Use local mode for testing; remove in cluster deployments
      .getOrCreate()
    val sc = spark.sparkContext

    // Your input dataset
    val studentData = List("ram|30|40", "sam|35|45", "ravi|37|44")

    // Convert the list to an RDD
    val studentRDD = sc.parallelize(studentData)

    // Extract just the physics scores as integers
    val physicsScores = studentRDD.map(line => {
      val columns = line.split("\\|") // Split on | (escaped for regex)
      columns(1).toInt // Physics score is the second column (index 1)
    })

    // Calculate total score and number of students in one pass with aggregate
    val (totalPhysicsScore, studentCount) = physicsScores.aggregate((0, 0))(
      // Combine each score with the accumulator (total, count)
      (acc, score) => (acc._1 + score, acc._2 + 1),
      // Merge accumulators from different partitions
      (acc1, acc2) => (acc1._1 + acc2._1, acc1._2 + acc2._2)
    )

    // Compute average and format output (matches your example: 102/3 = 34)
    val avgScore = totalPhysicsScore.toDouble / studentCount
    val output = s"physics|${avgScore.toInt}"

    // Print the result
    println(output)

    // Clean up Spark resources
    sc.stop()
    spark.stop()
  }
}

Breakdown of Key Parts:

  • Spark Setup: We use SparkSession (the recommended entry point for Spark 2.0+) to get a SparkContext for RDD operations. The local[*] master is for local testing—omit this when running on a cluster.
  • Data Loading: parallelize converts our local list of strings into an RDD distributed across Spark partitions.
  • Score Extraction: The map operation splits each line by | (we escape it because split uses regex) and pulls out the physics score, converting it to an integer for calculations.
  • Efficient Aggregation: aggregate lets us compute both the total score and student count in a single pass over the RDD, which is more efficient than doing two separate operations (like sum and count).
  • Output Formatting: We calculate the average, convert it to an integer (to match your example output), and format it as physics|34.

Expected Output

When you run this code, you'll get exactly the output you specified:

physics|34

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 08:28:51