3.4 MapReduce 高级特性与优化 本节摘要:MapReduce 的优化几乎全部作用于 3.3 节拆解的 Shuffle 链:Combiner 在 Map 端预聚合砍传输量、中间结果压缩省网络与磁盘、推测执行对冲慢节点。本节给出三项手段的适用条件与真实配置,并覆盖 DistributedCache 广播连接、计数器体系与数据倾斜的加盐解法。 图 3-4 优化三板斧的作用点:对准 Shuffle 链的三个环节 图 3-4 优化三板斧的作用点:对准 Shuffle 链的三个环节 优化一:Combiner——Map 端的预 Reduce 回看词频统计:某 Map 分片里 hadoop 出现了 10 万次,就会向网络发送 10 万个 (hadoop,1),而 Reduce 端只关心它们的和。
本节摘要:MapReduce 的优化几乎全部作用于 3.3 节拆解的 Shuffle 链:Combiner 在 Map 端预聚合砍传输量、中间结果压缩省网络与磁盘、推测执行对冲慢节点。本节给出三项手段的适用条件与真实配置,并覆盖 DistributedCache 广播连接、计数器体系与数据倾斜的加盐解法。

回看词频统计:某 Map 分片里 hadoop 出现了 10 万次,就会向网络发送 10 万个 (hadoop,1),而 Reduce 端只关心它们的和。Combiner 允许在每次溢写前先做一次本地聚合,把 10 万条压成 1 条 (hadoop,100000)。
job.setCombinerClass(SumReducer.class); // 词频场景直接复用Reducer
效果有多显著?做个保守估算:1000 个 Map、每个输出 1GB 中间结果,未压缩传输 1000GB;Combiner 压缩比哪怕只有 10:1,跨网络量降到 100GB,Reduce 拉取时间缩短近一个量级。3.3 节的 shuffle bytes 计数器可以直接验证这个差值。
但 Combiner 有一道不可逾越的数学门槛:聚合函数必须满足结合律,且与交换律配合后可任意分段执行。求和、最大值、计数可以;求平均数直接用 Combiner 会出错——(a,10,20) 先局部平均 15、(b,15,25) 局部平均 20,再平均 17.5,而正确答案是 17.5 吗?(10+20+15+25)/4=17.5 恰好对,但把局部样本数不均的情况换成先求局部均值再全局均值就会错。标准解法是输出 (sum, count) 二元组,Reduce 端再相除——把不可结合的量改写成可结合的量,是使用 Combiner 的通用心法。
还有一个细节:Combiner 的调用次数不确定(可能每个溢写一次、合并时再一次),因此它只能作为优化、不能承担正确性职责。逻辑正确性必须在"没有 Combiner"时也成立。
一行配置,网络与磁盘双省:
<property> <name>mapreduce.map.output.compress</name> <value>true</value> </property> <property> <name>mapreduce.map.output.compress.codec</name> <value>org.apache.hadoop.io.compress.SnappyCodec</value> </property>
编码器选择的权衡表:
| 编码器 | 压缩比 | 速度 | 可分割 | 适用 |
|---|---|---|---|---|
| Snappy | 低 | 极快 | 否 | 中间结果:速度优先 |
| LZO | 中 | 快 | 需索引 | 大日志归档 |
| Gzip | 高 | 慢 | 否 | 冷归档、最终输出 |
| Bzip2 | 最高 | 最慢 | 是 | 超冷归档且需分片 |
中间结果生命周期短,选 Snappy 这类"快而糙"的;最终落 HDFS 的产出若会被下游反复读,值得 Gzip 或含索引的 LZO。可分割性决定一个压缩文件能否被多个 Map 并行处理:Gzip 不可分割,一个 10GB 的 Gzip 文件只能一个 Map 串行解——中间结果无所谓,落盘输出就要慎重。
集群里常有"拖后腿"节点:坏了一半的磁盘、被邻居挤占的 CPU、失灵的风扇导致降频。一个 Map 任务分到这种机器上,99% 的任务都完成了,就等它。推测执行的策略:作业接近尾声时,对显著慢于均值的任务,在另一台机器启动一个同样的副本任务,谁先完成用谁的结果,另一个被杀。
<property><name>mapreduce.map.speculative</name><value>true</value></property> <property><name>mapreduce.reduce.speculative</name><value>true</value></property>
代价是重复计算消耗资源,因此两类场景必须关闭:Reduce 逻辑非幂等(写外部副作用,如向数据库写数);集群本身资源紧张(重复任务挤占真实任务)。判断"要不要关"的原则就一条:你的 reduce 重跑一遍,结果和副作用是否完全一致。
3.2 节的 Reduce 端连接要求维表装得进内存;更大的维表可以走另一条路——Map 端连接:把维表文件随作业分发到每个 Map 任务的本地工作目录,Map 启动时整个载入内存,逐条事实直接查表输出,连接在 Map 阶段完成,完全跳过 Shuffle:
// 分发维表文件到各任务节点 job.addCacheFile(new URI("/dim/user_dim.csv#user_dim")); // Mapper的setup阶段加载: @Override protected void setup(Context ctx) throws IOException { dim = new HashMap<>(); try (BufferedReader r = new BufferedReader( new FileReader("user_dim"))) { // 符号链接 指向本地副本 String line; while ((line = r.readLine()) != null) { String[] f = line.split(","); dim.put(f[0], f[1] + "," + f[2]); } } }
适用条件:维表能装进单个 Map 任务的内存(几 GB 是上限经验值),且事实表远大于维表。维表频繁更新时要配增量刷新管道,否则广播的是过期数据——这是广播连接最隐蔽的坑。
3.1 节埋的雷在此拆除。症状:99 个 Reduce 早早完成,1 个 Reduce 跑了三小时(计数器里某 Reduce 的 records 数是别的百倍)。解法是给热点 key 加随机前缀、分两轮聚合:
// 第一轮:热点key加盐打散 到100个Reduce各自局部聚合 map: ctx.write(new Text(salt + "_" + hotKey), value); // salt = rand(0,99) reduce1: 局部求和 → 输出 (salt_hotKey, partSum) // 第二轮:去掉盐 全局汇总 map2: ctx.write(new Text(realKey), partSum); reduce2: 求和 → 最终结果
代价是多一轮作业。变通做法是只对 top N 热点 key 加盐(其余 key 走原路径),需要先跑一次抽样统计热点分布——工程上通常由调度层(第 5 章 Hive)自动完成,手写 MR 时需自己判断。
优化离不开度量。除了框架内置计数器,自定义计数器是作业对账的标准工具:
enum LogCounters { DIRTY_LINE, SPIDER_HIT, NORMAL } // Mapper里累加: ctx.getCounter(LogCounters.DIRTY_LINE).increment(1); // 作业结束日志: // Counters: CLEAN_LOGS // DIRTY_LINE=1,204,552 // SPIDER_HIT=8,391,007 // NORMAL=190,404,441
ETL 管道的铁律是"输入行数 = 各类输出之和",用计数器对账比用输出文件行数可靠——它能区分"过滤掉的"与"处理失败的"。
| 症状(看哪) | 诊断 | 首选手段 |
|---|---|---|
| Spilled records 远超 Map output | 缓冲小或无Combiner | 加 io.sort.mb、上 Combiner |
| shuffle bytes 巨大 | 中间结果未压缩 | Snappy 压缩 + Combiner |
| 个别 Reduce 长尾 | 数据倾斜 | 热点加盐两轮聚合 |
| 尾部单任务慢 | 慢节点 | 确认推测执行开启 |
| 大维表连接慢 | 走了Shuffle | DistributedCache 广播连接 |
| 小文件海量分片 | Map数爆炸 | CombineFileInputFormat |
计算层走完。下一章把镜头拉高:支撑这一切计算的资源视角——YARN。