Hadoop:如何查看每个中间结果由哪个Mapper生成?
Absolutely feasible! Tracking mapper origins for intermediate results is a fantastic way to debug your MapReduce jobs and understand how work gets distributed across your cluster. Let me walk you through the most straightforward methods to achieve this.
1. Core Concept: Attach Mapper Task IDs to Output
Every mapper task in Hadoop has a unique TaskID that identifies it. By modifying your mapper code to append this ID to the output key-value pairs, you can directly trace each intermediate result back to its source mapper.
2. Modify the Mapper Code (Using WordCount as Example)
Let's use the built-in WordCount sample to demonstrate. We'll tweak the mapper to include the task ID in its output, then adjust the reducer to make the final results readable.
Updated Mapper Class
import org.apache.hadoop.mapreduce.Mapper; import org.apache.hadoop.mapreduce.TaskID; import org.apache.hadoop.io.Text; public class TrackedWordCountMapper extends Mapper<Object, Text, Text, Text> { private TaskID mapperTaskId; @Override protected void setup(Context context) { // Fetch the unique ID of the current mapper task mapperTaskId = context.getTaskAttemptID().getTaskID(); } @Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { StringTokenizer wordTokenizer = new StringTokenizer(value.toString()); while (wordTokenizer.hasMoreTokens()) { String word = wordTokenizer.nextToken(); // Output format: (word, "count:mapper_id") context.write(new Text(word), new Text("1:" + mapperTaskId.toString())); } } }
Updated Reducer Class (Optional but Recommended)
To aggregate counts and clearly show each mapper's contribution, adjust the reducer like this:
import org.apache.hadoop.mapreduce.Reducer; import org.apache.hadoop.io.Text; import java.io.IOException; import java.util.HashMap; import java.util.Map; public class TrackedWordCountReducer extends Reducer<Text, Text, Text, Text> { @Override protected void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { int totalCount = 0; Map<String, Integer> mapperContributions = new HashMap<>(); for (Text val : values) { String[] parts = val.toString().split(":"); int count = Integer.parseInt(parts[0]); String mapperId = parts[1]; totalCount += count; mapperContributions.put(mapperId, mapperContributions.getOrDefault(mapperId, 0) + count); } // Build a human-readable result StringBuilder result = new StringBuilder(totalCount + " | "); for (Map.Entry<String, Integer> entry : mapperContributions.entrySet()) { result.append(entry.getKey()).append(":").append(entry.getValue()).append(", "); } // Trim trailing comma and space if (result.length() > 3) { result.setLength(result.length() - 2); } context.write(key, new Text(result.toString())); } }
3. Package and Run the Modified Job
- Compile your updated classes and package them into a JAR file (e.g.,
tracked-wordcount.jar). - Submit the job to your cluster with this command:
hadoop jar tracked-wordcount.jar com.yourpackage.TrackedWordCount /path/to/input /path/to/output
(Replace com.yourpackage with your actual class package name)
4. View the Results
Once the job finishes, check the output to see mapper origins:
hdfs dfs -cat /path/to/output/part-r-00000
You'll get output like this:
hello 5 | mapper_1234:2, mapper_1235:3
world 3 | mapper_1234:3
Alternative: Check Mapper Logs (No Code Changes)
If you don't want to modify code, use the Hadoop Web UI to view mapper logs:
- Open the ResourceManager UI (default port: 8088)
- Locate your job in "Running Applications" or "Completed Applications"
- Click into the job, then select "Map Tasks"
- For each mapper, click "Logs" and look for intermediate output lines (you may need to adjust log levels first)
This method works for quick checks but is less clean than modifying the code.
内容的提问来源于stack exchange,提问作者ParisaN

