请求解析《Hadoop权威指南》中NCDC气象数据集的MapReduce温度最大值代码
Let's walk through this classic example from Hadoop: The Definitive Guide—it's a great introduction to how MapReduce works for processing structured (fixed-width) text data like NCDC's weather datasets. We'll cover the mapper logic in detail, then fill in the gaps for the incomplete driver code.
Overall Program Goal
This MapReduce job calculates the maximum recorded air temperature per year from raw NCDC weather records. The pipeline splits work into two core stages:
- Mapper: Extracts valid year-temperature pairs from raw input lines
- Reducer: Aggregates all temperatures for each year to find the highest value (note: the reducer code isn't included here, but we'll cover its expected logic)
Mapper Class: MaxTemperatureMapper
First, let's unpack the mapper code line by line:
1. Class Definition & Generics
public class MaxTemperatureMapper extends Mapper<LongWritable, Text, Text, IntWritable> {
The Mapper class uses four generics to define input/output types:
LongWritable: Input key (byte offset of the current line in the input file—Hadoop passes this automatically)Text: Input value (the entire line of raw weather data)Text: Output key (the year we're extracting)IntWritable: Output value (the recorded air temperature for that year)
We use Hadoop's Writable types instead of Java's native long, String, or int because they're optimized for serialization/deserialization across cluster nodes—critical for distributed processing.
2. Constant for Missing Data
private static final int MISSING = 9999;
NCDC uses 9999 to mark missing or invalid temperature readings. We'll filter these out later to avoid skewing our results.
3. The map() Method (Core Logic)
This method runs once per line of input data. Here's what each part does:
Convert input to a string:
String line = value.toString();The input
Textobject holds the raw line; we convert it to a JavaStringfor easy manipulation.Extract the year:
String year = line.substring(15, 19);NCDC's fixed-width format stores the 4-digit year at positions 15–18 (0-indexed). The
substring(15,19)call grabs exactly these four characters.Parse the air temperature:
int airTemperature; if (line.charAt(87) == '+') { // parseInt doesn't handle leading plus signs airTemperature = Integer.parseInt(line.substring(88, 92)); } else { airTemperature = Integer.parseInt(line.substring(87, 92)); }The temperature is stored at positions 87–91. A quirk:
Integer.parseInt()throws an error if the number starts with a+, so we check for that and adjust our substring start position if needed.Check data quality:
String quality = line.substring(92, 93);NCDC includes a 1-character quality code at position 92. Only codes
0,1,4,5,9indicate valid, reliable temperature readings.Filter & output valid data:
if (airTemperature != MISSING && quality.matches("[01459]")) { context.write(new Text(year), new IntWritable(airTemperature)); }We skip missing temperatures and low-quality readings. For valid entries, we use the
Contextobject to emit a key-value pair: the year (asText) and the temperature (asIntWritable). TheContextacts as a bridge between the mapper and the rest of the Hadoop framework, passing outputs to the reducer.
Driver Code: MaxTemperature
The provided driver code is incomplete, but this is what a full implementation would look like, and what it does:
Full Driver Logic (Completed)
import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.IntWritable; 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; import org.apache.hadoop.mapreduce.Reducer; import java.io.IOException; public class MaxTemperature { // Reducer class (missing from the original snippet) public static class MaxTemperatureReducer extends Reducer<Text, IntWritable, Text, IntWritable> { @Override public void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { int maxTemp = Integer.MIN_VALUE; for (IntWritable val : values) { maxTemp = Math.max(maxTemp, val.get()); } context.write(key, new IntWritable(maxTemp)); } } public static void main(String[] args) throws Exception { // 1. Initialize a MapReduce Job Job job = Job.getInstance(); job.setJarByClass(MaxTemperature.class); job.setJobName("Max Temperature"); // 2. Configure Mapper and Reducer classes job.setMapperClass(MaxTemperatureMapper.class); job.setReducerClass(MaxTemperatureReducer.class); // 3. Set output key/value types (must match mapper/reducer outputs) job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); // 4. Define input and output paths (from command-line arguments) FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); // 5. Submit the job and wait for completion System.exit(job.waitForCompletion(true) ? 0 : 1); } }
Key Driver Responsibilities
- Job Configuration: Sets up the job name, jar file, and links to the mapper/reducer classes.
- Input/Output Setup: Specifies where raw input data lives and where to write final results (passed as command-line args:
hadoop jar <jar-file> MaxTemperature <input-path> <output-path>). - Reducer Logic: The reducer takes all temperature values for a single year, iterates through them to find the maximum, then emits the year and its highest temperature.
Key Takeaways
- Fixed-Width Data Handling: The mapper relies on knowing exactly where each field lives in the raw text—common for legacy datasets like NCDC's.
- Data Quality: Filtering invalid/missing data is critical to getting accurate results in big data processing.
- Hadoop Writable Types: Always use these instead of native Java types for efficient serialization across the cluster.
内容的提问来源于stack exchange,提问作者Anirudh

