
一、Hadoop MapReduce 核心数据类型Hadoop 提供了一系列专门的数据类型这些类型都实现了WritableComparable接口以便数据可以被序列化进行网络传输、文件存储以及大小比较。1.1 常用数据类型列表BooleanWritable标准布尔型数值ByteWritable单字节数值DoubleWritable双精度浮点数FloatWritable单精度浮点数IntWritable整型数LongWritable长整型数Text使用 UTF-8 格式存储的文本NullWritable当 key,value 中的 key 或 value 为空时使用二、WordCount 示例分析旧 API2.1 完整源代码package org.apache.hadoop.examples; import java.io.IOException; import java.util.Iterator; import java.util.StringTokenizer; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapred.FileInputFormat; import org.apache.hadoop.mapred.FileOutputFormat; import org.apache.hadoop.mapred.JobClient; import org.apache.hadoop.mapred.JobConf; import org.apache.hadoop.mapred.MapReduceBase; import org.apache.hadoop.mapred.Mapper; import org.apache.hadoop.mapred.OutputCollector; import org.apache.hadoop.mapred.Reducer; import org.apache.hadoop.mapred.Reporter; import org.apache.hadoop.mapred.TextInputFormat; import org.apache.hadoop.mapred.TextOutputFormat; public class WordCount { public static class Map extends MapReduceBase implements MapperLongWritable, Text, Text, IntWritable { private final static IntWritable one new IntWritable(1); private Text word new Text(); public void map(LongWritable key, Text value, OutputCollectorText, IntWritable output, Reporter reporter) throws IOException { String line value.toString(); StringTokenizer tokenizer new StringTokenizer(line); while (tokenizer.hasMoreTokens()) { word.set(tokenizer.nextToken()); output.collect(word, one); } } } public static class Reduce extends MapReduceBase implements ReducerText, IntWritable, Text, IntWritable { public void reduce(Text key, IteratorIntWritable values, OutputCollectorText, IntWritable output, Reporter reporter) throws IOException { int sum 0; while (values.hasNext()) { sum values.next().get(); } output.collect(key, new IntWritable(sum)); } } public static void main(String[] args) throws Exception { JobConf conf new JobConf(WordCount.class); conf.setJobName(wordcount); conf.setOutputKeyClass(Text.class); conf.setOutputValueClass(IntWritable.class); conf.setMapperClass(Map.class); conf.setCombinerClass(Reduce.class); conf.setReducerClass(Reduce.class); conf.setInputFormat(TextInputFormat.class); conf.setOutputFormat(TextOutputFormat.class); FileInputFormat.setInputPaths(conf, new Path(args[0])); FileOutputFormat.setOutputPath(conf, new Path(args[1])); JobClient.runJob(conf); } }2.2 主方法Main配置分析2.2.1 Job 初始化main函数调用JobConf类对 MapReduce Job 进行初始化并通过setJobName()方法为 Job 命名。合理的命名有助于在 JobTracker 和 TaskTracker 页面中快速定位和监控 Job。JobConf conf new JobConf(WordCount.class); conf.setJobName(wordcount);2.2.2 输出数据类型设置设置 Job 输出结果 key,value 的数据类型。由于结果是 单词,个数key 设置为Text类型相当于 Java 的Stringvalue 设置为IntWritable类型相当于 Java 的int。conf.setOutputKeyClass(Text.class); conf.setOutputValueClass(IntWritable.class);2.2.3 处理类设置设置 Job 处理的 Map拆分、Combiner中间结果合并以及 Reduce合并的相关处理类。这里使用 Reduce 类进行 Map 产生的中间结果合并以减少网络数据传输压力。conf.setMapperClass(Map.class); conf.setCombinerClass(Reduce.class); conf.setReducerClass(Reduce.class);2.2.4 输入输出格式设置设置输入输出格式和路径。conf.setInputFormat(TextInputFormat.class); conf.setOutputFormat(TextOutputFormat.class); FileInputFormat.setInputPaths(conf, new Path(args[0])); FileOutputFormat.setOutputPath(conf, new Path(args[1]));2.3 InputFormat 与 InputSplit 详解2.3.1 InputSplit 概念InputSplit是 Hadoop 定义的用来传送给每个单独 map 的数据单元。InputSplit 存储的并非数据本身而是一个分片长度和一个记录数据位置的数组。生成 InputSplit 的方法可以通过InputFormat()来设置。当数据传送给 map 时map 会将输入分片传送给 InputFormatInputFormat 则调用getRecordReader()方法生成 RecordReader。RecordReader 再通过createKey()、createValue()方法创建可供 map 处理的 key,value 对。简而言之InputFormat()方法是用来生成可供 map 处理的 key,value 对的。2.3.2 InputFormat 继承体系Hadoop 预定义了多种方法将不同类型的输入数据转化为 map 能够处理的 key,value 对它们都继承自 InputFormatInputFormat |---BaileyBorweinPlouffe.BbpInputFormat |---ComposableInputFormat |---CompositeInputFormat |---DBInputFormat |---DistSum.Machine.AbstractInputFormat |---FileInputFormat |---CombineFileInputFormat |---KeyValueTextInputFormat |---NLineInputFormat |---SequenceFileInputFormat |---TeraInputFormat |---TextInputFormat2.3.3 TextInputFormat默认输入方法TextInputFormat是 Hadoop 默认的输入方法继承自FileInputFormat。在 TextInputFormat 中每个文件或其一部分都会单独地作为 map 的输入。之后每行数据都会生成一条记录每条记录表示成 key,value 形式key每个数据记录在数据分片中的字节偏移量数据类型为LongWritablevalue每行的内容数据类型为Text2.4 OutputFormat 详解每一种输入格式都有一种输出格式与其对应。默认的输出格式是TextOutputFormat这种输出方式与输入类似会将每条记录以一行的形式存入文本文件。不过它的键和值可以是任意形式的因为程序会调用toString()方法将键和值转换为String类型再输出。2.5 Map 类分析2.5.1 Map 类定义public static class Map extends MapReduceBase implements MapperLongWritable, Text, Text, IntWritable { private final static IntWritable one new IntWritable(1); private Text word new Text(); public void map(LongWritable key, Text value, OutputCollectorText, IntWritable output, Reporter reporter) throws IOException { String line value.toString(); StringTokenizer tokenizer new StringTokenizer(line); while (tokenizer.hasMoreTokens()) { word.set(tokenizer.nextToken()); output.collect(word, one); } } }2.5.2 Map 类解析Map 类继承自MapReduceBase并实现了Mapper接口。此接口是一个规范类型它有 4 种形式的参数分别用来指定 map 的输入 key 值类型、输入 value 值类型、输出 key 值类型和输出 value 值类型。在本例中因为使用的是TextInputFormat它的输出 key 值是LongWritable类型输出 value 值是Text类型所以 map 的输入类型为 LongWritable, Text。本例需要输出 word,1 这样的形式因此输出的 key 值类型是Text输出的 value 值类型是IntWritable。实现此接口类还需要实现map方法该方法具体负责对输入进行操作。在本例中map方法对输入的行以空格为单位进行切分然后使用OutputCollect收集输出的 word,1。2.6 Reduce 类分析2.6.1 Reduce 类定义public static class Reduce extends MapReduceBase implements ReducerText, IntWritable, Text, IntWritable { public void reduce(Text key, IteratorIntWritable values, OutputCollectorText, IntWritable output, Reporter reporter) throws IOException { int sum 0; while (values.hasNext()) { sum values.next().get(); } output.collect(key, new IntWritable(sum)); } }2.6.2 Reduce 类解析Reduce 类也是继承自MapReduceBase需要实现Reducer接口。Reduce 类以 map 的输出作为输入因此 Reduce 的输入类型是 Text, IntWritable。而 Reduce 的输出是单词和它的数目因此它的输出类型是 Text, IntWritable。Reduce 类也要实现reduce方法在此方法中reduce 函数将输入的 key 值作为输出的 key 值然后将获得的多个 value 值加起来作为输出的值。三、WordCount 示例分析新 API3.1 完整源代码package org.apache.hadoop.examples; import java.io.IOException; import java.util.StringTokenizer; import org.apache.hadoop.conf.Configuration; 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.Mapper; import org.apache.hadoop.mapreduce.Reducer; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; import org.apache.hadoop.util.GenericOptionsParser; public class WordCount { public static class TokenizerMapper extends MapperObject, Text, Text, IntWritable { private final static IntWritable one new IntWritable(1); private Text word new Text(); public void map(Object key, Text value, Context context) throws IOException, InterruptedException { StringTokenizer itr new StringTokenizer(value.toString()); while (itr.hasMoreTokens()) { word.set(itr.nextToken()); context.write(word, one); } } } public static class IntSumReducer extends ReducerText, IntWritable, Text, IntWritable { private IntWritable result new IntWritable(); public void reduce(Text key, IterableIntWritable values, Context context) throws IOException, InterruptedException { int sum 0; for (IntWritable val : values) { sum val.get(); } result.set(sum); context.write(key, result); } } public static void main(String[] args) throws Exception { Configuration conf new Configuration(); String[] otherArgs new GenericOptionsParser(conf, args).getRemainingArgs(); if (otherArgs.length ! 2) { System.err.println(Usage: wordcount in out); System.exit(2); } Job job new Job(conf, word count); job.setJarByClass(WordCount.class); job.setMapperClass(TokenizerMapper.class); job.setCombinerClass(IntSumReducer.class); job.setReducerClass(IntSumReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); FileInputFormat.addInputPath(job, new Path(otherArgs[0])); FileOutputFormat.setOutputPath(job, new Path(otherArgs[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }3.2 Map 过程分析3.2.1 Map 类定义public static class TokenizerMapper extends MapperObject, Text, Text, IntWritable { private final static IntWritable one new IntWritable(1); private Text word new Text(); public void map(Object key, Text value, Context context) throws IOException, InterruptedException { StringTokenizer itr new StringTokenizer(value.toString()); while (itr.hasMoreTokens()) { word.set(itr.nextToken()); context.write(word, one); } } }3.2.2 Map 过程解析Map 过程需要继承org.apache.hadoop.mapreduce包中的Mapper类并重写其map方法。通过在map方法中添加输出 key 值和 value 值到控制台的代码可以发现map方法中 value 值存储的是文本文件中的一行以回车符为行结束标记而 key 值为该行的首字母相对于文本文件的首地址的偏移量。然后StringTokenizer类将每一行拆分成一个个的单词并将 word,1 作为map方法的结果输出其余的工作都交由MapReduce 框架处理。3.3 Reduce 过程分析3.3.1 Reduce 类定义public static class IntSumReducer extends ReducerText, IntWritable, Text, IntWritable { private IntWritable result new IntWritable(); public void reduce(Text key, IterableIntWritable values, Context context) throws IOException, InterruptedException { int sum 0; for (IntWritable val : values) { sum val.get(); } result.set(sum); context.write(key, result); } }3.4 新 API 与旧 API 对比总结通过对比新旧两个 WordCount 示例可以清晰地看到 Hadoop MapReduce 新 APIorg.apache.hadoop.mapreduce包相对于旧 APIorg.apache.hadoop.mapred包的主要改进和优势1. 统一的 Context 接口旧 APIMapper 和 Reducer 分别使用OutputCollector和Reporter两个参数来输出结果和报告进度。新 API引入了统一的Context接口将输出和进度报告功能合并简化了方法签名提高了代码的简洁性和一致性。2. 更现代的 Job 配置方式旧 API使用JobConf类进行配置通过setXxxClass()方法链式调用。新 API使用Job类配置方法更加面向对象如job.setMapperClass()、job.setReducerClass()等代码可读性更好。3. 类型安全的泛型使用旧 APIMapper 和 Reducer 接口使用原始类型类型安全性较差。新 API全面使用泛型如MapperObject, Text, Text, IntWritable和ReducerText, IntWritable, Text, IntWritable提供了更好的编译时类型检查。4. 改进的输入输出处理旧 API使用IteratorIntWritable作为 Reduce 的输入值迭代器。新 API使用IterableIntWritable支持增强的 for 循环代码更加简洁直观。5. 配置和参数解析的改进旧 API直接使用命令行参数数组。新 API引入GenericOptionsParser和Configuration类支持更灵活的参数解析和配置管理。6. 执行流程的优化旧 API使用JobClient.runJob(conf)提交作业。新 API使用job.waitForCompletion(true)提交并等待作业完成返回布尔值表示成功与否控制更加精细。总结新 API 在设计上更加现代化、类型安全且易于使用减少了样板代码提高了开发效率。虽然旧 API 仍然可用但新项目推荐使用新 API以获得更好的开发体验和代码质量。