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 aSparkContextfor RDD operations. Thelocal[*]master is for local testing—omit this when running on a cluster. - Data Loading:
parallelizeconverts our local list of strings into an RDD distributed across Spark partitions. - Score Extraction: The
mapoperation splits each line by|(we escape it becausesplituses regex) and pulls out the physics score, converting it to an integer for calculations. - Efficient Aggregation:
aggregatelets 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 (likesumandcount). - 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
相关产品推荐
相关产品推荐

