3.2 类型系统与常用 API 本节摘要:map 与 reduce 之间隔着网络与磁盘,所以键值类型必须是可序列化、可比较的二等公民——Hadoop 为此自建了 Writable 类型系统,而不是直接用 Java 原生类型。本节拆解 Writable 与 WritableComparable 两层接口,盘点 Text、LongWritable 等常用类型的内存布局,走查 Job 配置 API 的完整清单,最后动手写一个自定义复合键,并解释 reduce 端值迭代器背后的对象复用陷阱。 为什么 Java 自带类型不够用 Map 端的对象要被序列化进环形缓冲区、可能压缩落盘、再经网络发给 Reduce 端反序列化——一个类型的对象在分布式全程要经历"序列化—比较—排序—复制"四种命运。
本节摘要:map 与 reduce 之间隔着网络与磁盘,所以键值类型必须是可序列化、可比较的二等公民——Hadoop 为此自建了 Writable 类型系统,而不是直接用 Java 原生类型。本节拆解 Writable 与 WritableComparable 两层接口,盘点 Text、LongWritable 等常用类型的内存布局,走查 Job 配置 API 的完整清单,最后动手写一个自定义复合键,并解释 reduce 端值迭代器背后的对象复用陷阱。
Map 端的对象要被序列化进环形缓冲区、可能压缩落盘、再经网络发给 Reduce 端反序列化——一个类型的对象在分布式全程要经历"序列化—比较—排序—复制"四种命运。Java 原生 String、Long 满足不了其中的两条:标准序列化(Serializable)冗余大,每写一个对象都带类描述开销;也没有统一的二进制比较协议,排序前必须完整反序列化成对象。
Hadoop 的回答是自建一套类型契约,分两级:
| 接口 | 承担的命运 | 代表类型 |
|---|---|---|
| Writable | 序列化与反序列化 | Text、BytesWritable、NullWritable |
| WritableComparable | 序列化 + 键的排序比较 | LongWritable、IntWritable、Text |
规则很干脆:map 与 reduce 的输出类型必须实现 Writable;充当键的类型还必须实现 WritableComparable。值只要能"打包上车",键还要能"排队"——因为 Shuffle 要按键排序,Reduce 端要靠比较器把同键记录归到一组。
常用类型速查:
| Hadoop 类型 | 对应 Java 类型 | 一个细节 |
|---|---|---|
| IntWritable / LongWritable | int / long | 定长 4 或 8 字节,比较用原始数值 |
| FloatWritable / DoubleWritable | float / double | 定长,适合度量值 |
| Text | String | 变长 UTF-8,长度前缀加字节体 |
| BooleanWritable | boolean | 定长 1 字节 |
| BytesWritable | byte[] | 任意二进制 payload |
| NullWritable | 无 | 占位符,输出"只有值没有键"时用它 |
Text 值得多看一眼:它内部存的是 UTF-8 字节而非 Java 字符串,toString() 才做转码;它重载了 set(String) 与 set(byte[], off, len),频繁调用时复用同一个 Text 对象、反复 set,比每次 new 省一大截对象分配。这个习惯贯穿所有高性能 MapReduce 代码。
一个作业里类型要在五个位置保持一致,任何一处对不上就是运行时 ClassCastException:
map 输出键类型 → 分区器与排序比较器要按它工作 map 输出值类型 → 序列化进缓冲区 reduce 输入键值 → 必须与 map 输出对齐 reduce 输出键值 → 由 OutputFormat 写文件 作业全局声明 → setOutputKeyClass 一处声明多处校验
当 map 与 reduce 的键类型不同(比如 map 输出 Text、reduce 输出 NullWritable 的"只留值"场景),必须额外调用 setMapOutputKeyClass 与 setMapOutputValueClass 声明中间类型。这是新手报错榜排名第一的来源:
Job job = Job.getInstance(); job.setMapOutputKeyClass(Text.class); // 中间键 job.setMapOutputValueClass(LongWritable.class); job.setOutputKeyClass(NullWritable.class); // 最终键(只写值) job.setOutputValueClass(Text.class);
用户能拧的旋钮都挂在 Job 对象上,按职能分组记忆:
Job job = Job.getInstance(conf, "dept-salary"); job.setJarByClass(DeptSalarySum.class); // 定位 jar 供集群分发 // 逻辑插件:全都可以换成自定义类 job.setMapperClass(DeptMapper.class); job.setReducerClass(SumReducer.class); job.setCombinerClass(SumReducer.class); // 本地预聚合,3.3 节主角 job.setPartitionerClass(HashPartitioner.class); job.setGroupingComparatorClass(...); // 4.4 节主角 job.setSortComparatorClass(...); // 数量旋钮 job.setNumReduceTasks(2); // 输入输出 FileInputFormat.addInputPath(job, new Path("/input/dept")); FileInputFormat.setInputDirRecursive(job, true); FileOutputFormat.setOutputPath(job, new Path("/output/dept")); LazyOutputFormat.setOutputFormatClass(job, TextOutputFormat.class); // 提交并等待 boolean ok = job.waitForCompletion(true); // true 表示打印进度到控制台
三对 InputFormat 与 OutputFormat 的选择逻辑值得单独记:InputFormat 决定"怎么切、怎么读"(2.1 节),OutputFormat 决定"怎么写、写哪里"。最常用的组合是 TextInputFormat 进、TextOutputFormat 出;链式作业之间则用 SequenceFileOutputFormat 接 SequenceFileInputFormat,中间产物以二进制键值对落盘,省去反复的文本解析。
Reducer 一侧还有三个"换演员"的钩子:setup 在处理任何键之前调用一次,适合加载词典到内存;reduce 逐键调用;cleanup 在所有键处理完后调用一次,适合把攒了半天的缓冲刷出去。Mapper 同样有这三段生命周期。一个典型用法:
public static class DictMapper extends Mapper<LongWritable, Text, Text, LongWritable> { private Set<String> stopWords = new HashSet<>(); @Override protected void setup(Context context) { // 把分布式缓存里的停用词表读进内存 只做一次 try (BufferedReader r = new BufferedReader( new FileReader("stopwords.txt"))) { String w; while ((w = r.readLine()) != null) stopWords.add(w.trim()); } } @Override protected void map(LongWritable k, Text line, Context context) { for (String t : line.toString().split("\\s+")) { if (t.isEmpty() || stopWords.contains(t)) continue; context.write(new Text(t), new LongWritable(1)); } } }
配合提交前的 job.addCacheFile(new Path("/dict/stopwords.txt").toUri()),每个任务节点本地就能读到小文件,不占 map 的每行开销。setup 与 cleanup 的价值就在"把每条记录只该摊一次的成本搬出循环"。
复合键是类型系统的实战考题。假设要做"按年份和气温排序",中间键要同时携带两个字段且按"年升序、温度降序"比较。手写如下:
public class YearTempKey implements WritableComparable<YearTempKey> { private int year; private int temp; public void set(int year, int temp) { this.year = year; this.temp = temp; } @Override public void write(DataOutput out) throws IOException { out.writeInt(year); out.writeInt(temp); } @Override public void readFields(DataInput in) throws IOException { year = in.readInt(); temp = in.readInt(); } @Override public int compareTo(YearTempKey o) { if (year != o.year) return Integer.compare(year, o.year); return Integer.compare(o.temp, temp); // 温度降序 反着比 } @Override public int hashCode() { return year * 31 + temp; // 与 compareTo 保持同键同 hash } }
写自定义类型有三条军规:write 与 readFields 的字段顺序必须严格一致;compareTo 决定 Shuffle 排序,写错顺序结果就错;用作分区键时 hashCode 必须与 equals 语义一致,否则同键数据会散到不同 Reduce。字段全为定长原始类型时,还能进一步实现 write 与 compareTo 直接操作字节的 RawComparator(YearTempKey.Comparator),排序时零反序列化——这正是 4.4 节排序优化的话题。
reduce 签名里的 Iterable<VALUEIN> 藏着全 API 最阴的坑:迭代器吐出的很可能是同一个对象,每走一步框架只重填它的字段。想先把所有值攒进 List 再处理,攒到最后是一串指向同一对象的引用:
// 错误写法 攒了 1000 个引用 全指向最后一个值 List<Long> all = new ArrayList<>(); for (LongWritable v : values) all.add(v.get()); // 正确写法之一 直接拷出原始值 List<Long> all = new ArrayList<>(); for (LongWritable v : values) all.add(Long.valueOf(v.get()));
好在 v.get() 取的是原始 long 的快照,上面第一种"错误"其实被 get 救了;真正翻车的是 all.add(v) 这种存对象引用的写法。成因是框架为省对象分配,用单例 value 对象流式重放——next() 每次读新字节进同一对象。推演一下输出就能看到症状:输入三个值 7、8、9,错误版输出 9、9、9,正确版输出 7、8、9。记住口诀:进 reduce 先拷值,别收藏迭代器。
另一个同源知识点:reduce 端的键对象在遍历 values 期间保持不变,但离开循环后不保证仍有效,需要留键就 new Text(key) 拷一份。
下一节回到执行机制主线:Combiner 如何在 map 端就地消化掉一部分 reduce 的工作。