如何使用MongoDB MapReduce筛选指定学生ID的所有成绩?
Problem Context
You've got a student grades dataset structured like this (each line is a JSON object):
{ "StudentID" : 1, "Subject" : "Maths", "Grade": "Good" }
{ "StudentID" : 1, "Subject" : "Physics", "Grade": "Excellent" }
{ "StudentID" : 2, "Subject" : "Maths", "Grade": "Very Good" }
Your goal is to extract all the grades associated with a specific StudentID using MapReduce. Let's break down how to do this step by step.
1. Core MapReduce Strategy
The approach is straightforward:
- Mapper: Filter and emit only records matching your target StudentID. Use the StudentID as the key, and a combined string of Subject + Grade as the value. This cuts down on unnecessary data transfer early in the process.
- Reducer: Collect all values for the target key and output them in a clean, readable format (since we filtered in the mapper, the reducer only handles the student we care about).
2. Java MapReduce Implementation
If you're using Hadoop's Java API, here's how to build the components:
2.1 Mapper Class
We'll parse each JSON line, check for the target StudentID, and emit relevant key-value pairs. We use the org.json library for simple JSON parsing.
import java.io.IOException; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; import org.json.JSONObject; public class StudentGradeMapper extends Mapper<Object, Text, Text, Text> { private final Text studentKey = new Text(); private final Text gradeValue = new Text(); // Set your target StudentID here private static final int TARGET_STUDENT_ID = 1; @Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { // Parse the input JSON line JSONObject jsonRecord = new JSONObject(value.toString()); int studentId = jsonRecord.getInt("StudentID"); // Only process records for our target student if (studentId == TARGET_STUDENT_ID) { String subject = jsonRecord.getString("Subject"); String grade = jsonRecord.getString("Grade"); studentKey.set(String.valueOf(studentId)); gradeValue.set(subject + ": " + grade); context.write(studentKey, gradeValue); } } }
2.2 Reducer Class
The reducer groups all subject-grade pairs for the target student and outputs them in a structured way.
import java.io.IOException; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; public class StudentGradeReducer extends Reducer<Text, Text, Text, Text> { @Override protected void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { // Print a header for the student context.write(new Text("Grades for Student ID: " + key), new Text("")); // List all subject-grade pairs for (Text gradeEntry : values) { context.write(new Text("- "), gradeEntry); } } }
2.3 Driver Class
This sets up and runs the MapReduce job:
import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; public class StudentGradeDriver { public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "Student Grade Lookup"); job.setJarByClass(StudentGradeDriver.class); job.setMapperClass(StudentGradeMapper.class); job.setReducerClass(StudentGradeReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(Text.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }
3. Hadoop Streaming Alternative (Python)
If you prefer Python for quicker scripting, use Hadoop Streaming with these scripts:
Mapper Script (student_mapper.py)
import sys import json TARGET_ID = 1 for line in sys.stdin: line = line.strip() if not line: continue try: record = json.loads(line) if record["StudentID"] == TARGET_ID: print(f"{TARGET_ID}\t{record['Subject']}: {record['Grade']}") except Exception: # Skip invalid JSON lines continue
Reducer Script (student_reducer.py)
import sys current_student = None grade_list = [] for line in sys.stdin: line = line.strip() if not line: continue student_id, grade_info = line.split("\t", 1) if current_student != student_id: if current_student: # Print the previous student's grades print(f"Grades for Student ID: {current_student}") for grade in grade_list: print(f"- {grade}") current_student = student_id grade_list = [grade_info] else: grade_list.append(grade_info) # Print the final student's grades if current_student: print(f"Grades for Student ID: {current_student}") for grade in grade_list: print(f"- {grade}")
Run the Job
Execute this command in your Hadoop environment:
hadoop jar $HADOOP_HOME/share/hadoop/tools/lib/hadoop-streaming-*.jar \ -files student_mapper.py,student_reducer.py \ -mapper python3 student_mapper.py \ -reducer python3 student_reducer.py \ -input /path/to/your/dataset \ -output /path/to/output/dir
4. Sample Output
For the sample dataset and target ID 1, the output will look like:
Grades for Student ID: 1 - Maths: Good - Physics: Excellent
内容的提问来源于stack exchange,提问作者David Sepp

