3.2 类型系统与常用 API


文档摘要

3.2 类型系统与常用 API 本节摘要:map 与 reduce 之间隔着网络与磁盘,所以键值类型必须是可序列化、可比较的二等公民——Hadoop 为此自建了 Writable 类型系统,而不是直接用 Java 原生类型。本节拆解 Writable 与 WritableComparable 两层接口,盘点 Text、LongWritable 等常用类型的内存布局,走查 Job 配置 API 的完整清单,最后动手写一个自定义复合键,并解释 reduce 端值迭代器背后的对象复用陷阱。 为什么 Java 自带类型不够用 Map 端的对象要被序列化进环形缓冲区、可能压缩落盘、再经网络发给 Reduce 端反序列化——一个类型的对象在分布式全程要经历"序列化—比较—排序—复制"四种命运。

3.2 类型系统与常用 API

本节摘要:map 与 reduce 之间隔着网络与磁盘,所以键值类型必须是可序列化、可比较的二等公民——Hadoop 为此自建了 Writable 类型系统,而不是直接用 Java 原生类型。本节拆解 Writable 与 WritableComparable 两层接口,盘点 Text、LongWritable 等常用类型的内存布局,走查 Job 配置 API 的完整清单,最后动手写一个自定义复合键,并解释 reduce 端值迭代器背后的对象复用陷阱。

为什么 Java 自带类型不够用

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 的"只留值"场景),必须额外调用 setMapOutputKeyClasssetMapOutputValueClass 声明中间类型。这是新手报错榜排名第一的来源:

Job job = Job.getInstance(); job.setMapOutputKeyClass(Text.class); // 中间键 job.setMapOutputValueClass(LongWritable.class); job.setOutputKeyClass(NullWritable.class); // 最终键(只写值) job.setOutputValueClass(Text.class);

Job 配置 API 全景

用户能拧的旋钮都挂在 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 的价值就在"把每条记录只该摊一次的成本搬出循环"。

动手写一个自定义 WritableComparable

复合键是类型系统的实战考题。假设要做"按年份和气温排序",中间键要同时携带两个字段且按"年升序、温度降序"比较。手写如下:

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。字段全为定长原始类型时,还能进一步实现 writecompareTo 直接操作字节的 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) 拷一份。

本节自查

  1. 为什么充当键的类型必须实现 WritableComparable 而值只需 Writable;
  2. 写出 map 输出类型与 reduce 输出类型不一致时必须补的两行配置;
  3. 手推对象复用陷阱:输入值序列 1、2、3,分别写出"存引用"与"存 get 值"两种写法的列表内容;
  4. 给自定义 YearTempKey 再加一个 station 字段并保持年、温度、站的排序语义,检查 write 顺序与 compareTo 是否同步改。

下一节回到执行机制主线:Combiner 如何在 map 端就地消化掉一部分 reduce 的工作。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U