4.3 Reduce 与输出落地 本节摘要:Reduce 任务收齐自己分区的数据后,经历"拷贝—归并—分组—回调"四段,最终经 OutputFormat 写成 part 文件落地。本节走查 reduce 回调里的分组语义(排序比较器与分组比较器的分工)、Reduce 端的内存与磁盘参数、TextOutputFormat 与 SequenceFileOutputFormat 的取舍、多目录输出 MultipleOutputs 与空输出目录的坑,最后把第三幕六步串成一张全景账单。 Reduce 任务的一生 一个 Reduce 任务被调度后依次做四件事: 第一段,拷贝。向所有已完成的 Map 任务发起 HTTP 请求,只拉属于自己的分区段。
本节摘要:Reduce 任务收齐自己分区的数据后,经历"拷贝—归并—分组—回调"四段,最终经 OutputFormat 写成 part 文件落地。本节走查 reduce 回调里的分组语义(排序比较器与分组比较器的分工)、Reduce 端的内存与磁盘参数、TextOutputFormat 与 SequenceFileOutputFormat 的取舍、多目录输出 MultipleOutputs 与空输出目录的坑,最后把第三幕六步串成一张全景账单。
一个 Reduce 任务被调度后依次做四件事:
第一段,拷贝。向所有已完成的 Map 任务发起 HTTP 请求,只拉属于自己的分区段。并行拷贝线程默认 5 个(参数可调), Map 产出得早就早拉,不必等同批——这解释了 3.1 节的观察:Reduce 的"copy 阶段"往往在最后一个 Map 结束前就已开始。
第二段,归并。拉来的小段文件先填内存缓冲(默认约 70% 任务堆内存的份额),溢出则边排序边写盘,最终多路归并成按键有序的大文件。这里的归并流数受归并因子控制(默认 10),中间层数为对数级别。
第三段,分组。归并结果按键有序,框架把"相等"的连续键划为一组。判断相等用的是分组比较器——默认与排序比较器相同,但可以单独设置得更宽松(4.4 节的主角)。每组触发一次 reduce 回调。
第四段,回调写出。reduce 写出的键值对交给 OutputFormat 的 RecordWriter,落成最终文件。
Map 全部输出 │ 只拷贝本分区 ▼ copy 若干段 ──► 内存缓冲 ──溢出──► 磁盘段 │ 多路归并 ▼ 全局有序键值流 ──按分组比较器切组──► reduce 逐组回调 │ context.write ▼ part-r-00000 等
reduce 收到的契约是"一个键 + 它名下的所有值",但更精确的表述是"排序上相邻、且分组比较器判等的一段流"。默认情况下分组比较器就是键的排序比较器,于是"判等"与"排序下相邻"重合,直觉成立。一旦自定义了分组比较器(比如复合键只按年份分组),同一次回调收到的"键"其实是该组第一行的完整键——年份相同、温度可能不同。这个语义是二级排序(6.2 节 TopN 模式)的全部地基。
reduce 回调的三段生命周期在此再压成一张表:
| 时机 | 方法 | 典型用途 |
|---|---|---|
| 组前一次 | setup | 加载字典、初始化连接 |
| 每组一次 | reduce | 业务归并 |
| 全部组后 | cleanup | 刷缓冲、关连接、输出收尾统计 |
值得推演的细节:reduce 处理完一个组后,框架流式地移动到下一组,不会把所有组驻留内存——所以"每组最大值"这类模式天然可扩展,而"全局 TopN"要么靠 Reduce 数为 1(放弃并行),要么靠每个 Reduce 局部 TopN 再汇聚(树形归并,6.2 节展开)。
reduce 写出的每对键值都流经 RecordWriter,路怎么修由 OutputFormat 决定:
TextOutputFormat(默认):每行一条,键与值之间用制表符分隔。人可读、下游 grep 方便,但类型信息全丢——LongWritable 的 1 与 Text 的 "1" 写出来一模一样,链式作业再读就得重新解析。
// 默认等价于 job.setOutputFormatClass(TextOutputFormat.class); // 可调分隔符 conf.set("mapreduce.output.textoutputformat.separator", ",");
SequenceFileOutputFormat:二进制键值对容器,带同步标记,可切分、可压缩(BLOCK 级压缩配 snappy 是中间产物的黄金组合)。作业链 A→B→C 的中间环节用 SequenceFile,B 省掉文本解析,C 拿到的类型即 A 写出的类型。与 2.3 节的压缩结论衔接:中间数据用快压缩省网络,最终数据看下游谁来读。
NullOutputFormat:什么都不写,作业只为副作用而跑(比如把结果写进外部数据库)。
一条配套规则:输出目录必须事先不存在,否则作业直接拒绝执行——这是防止覆盖旧结果的保守设计。重复实验时的惯例是先 hadoop fs -rm -r 再提交:
hadoop fs -rm -r /output/dept hadoop jar dept-sum.jar DeptSalarySum /input/dept /output/dept hadoop fs -ls /output/dept
Found 3 items -rw-r--r-- 1 u u 0 ... _SUCCESS -rw-r--r-- 1 u u 12 ... part-r-00000 -rw-r--r-- 1 u u 12 ... part-r-00001
_SUCCESS 是空标记文件,下游调度系统靠它判断作业成功;part-r-XXXXX 按 Reduce 编号编号,Reduce 数为 0 时文件名是 part-m-XXXXX(Map 直接写出)。
真实作业常要"按类别分仓":清洗作业把好数据写 main 目录、坏数据写 quarantine 目录。MultipleOutputs 把 context 的一个输出通道裂成多个命名通道:
public static class CleanReducer extends Reducer<Text, Text, Text, Text> { private MultipleOutputs<Text, Text> mos; @Override protected void setup(Context context) { mos = new MultipleOutputs<>(context); } @Override protected void reduce(Text key, Iterable<Text> values, Context context) { for (Text v : values) { String[] cols = v.toString().split(","); boolean bad = cols.length != 5; // 好数据进 main 坏数据进 quarantine 文件名带分区后缀 mos.write(bad ? "quarantine" : "main", key, v, bad ? "dirty-" : "clean-"); } } @Override protected void cleanup(Context context) throws IOException { mos.close(); // 不 close 会丢数据且报错 这是头号坑 } }
// 注册 MultipleOutputs.addNamedOutput(job, "main", TextOutputFormat.class, Text.class, Text.class); MultipleOutputs.addNamedOutput(job, "quarantine", TextOutputFormat.class, Text.class, Text.class);
产物形如 main-r-00000、quarantine-r-00000,每个命名通道各自按 Reduce 编号分片。两条纪律:mos 必须在 cleanup 里 close;文件名若自己拼了分区后缀就不要再叠加默认后缀,否则出现 clean--r-00000 这类双重后缀。
另一个常见抱怨"空分区也生成零字节 part 文件"由 LazyOutputFormat 解决——它延迟到第一条记录真正写出时才建文件:
LazyOutputFormat.setOutputFormatClass(job, TextOutputFormat.class);
小文件泛滥的集群上这一行能省下成千上万个空文件的 NameNode 元数据条目。
至此第三幕机制讲完,把 4.1 到本节的六步串成一张账单,标出每步的可调项:
| 步骤 | 归属 | 主要 IO | 关键参数 |
|---|---|---|---|
| 环形缓冲 | Map | 内存 | 缓冲大小、溢写阈值 |
| 溢写排序 | Map | 内存→磁盘写 | 排序实现、Combiner、压缩 |
| 分段归并 | Map | 磁盘读写 | 归并因子 |
| 分区拷贝 | Reduce | 网络读、磁盘写 | 拷贝线程数 |
| 多路归并 | Reduce | 磁盘读写 | 归并因子、内存份额 |
| 分组回调 | Reduce | 内存 | 分组比较器 |
数一数纯 IO 次数:写盘、读盘、网络、写盘、读盘,共五次搬运(4.1 节的旧账)。Reduce 侧能省的是第 4、5 步——如果 Map 端 Combiner 已经把同键记录压扁,或者干脆 Reduce 数为 0,这笔账立刻短一截。机制层面的账算清了,工程层面的账(参数怎么调、倾斜怎么破)留到 6.3 节的调优清单。
下一节深入比较器:排序与分组两把尺子如何配合出二级排序。