3.3 Combiner 本地预聚合 本节摘要:Combiner 是运行在 Map 任务溢写环节的本地 Reduce:相同键的值先在本任务内归并一次,再落盘、再传输。它不改变最终结果,只压缩中间数据量。合法性条件是归并操作满足"先局部合并再全局合并与一次全局合并等价"。用对了,Shuffle 流量数量级下降;用错了,结果悄悄出错。 先看一组流量数字 统计一个 100GB 日志里每个 URL 的访问次数。假设 200 个 Map 任务,日志里 URL 总数不多(10 万个),但每个任务都要为自己见过的每个 URL 输出一条记录。极端均匀的假设下,每个任务输出 50 万条"URL, 1"(平均每条 20 字节左右),全集群中间数据约 200 × 10MB = 2GB——听起来还能忍。
本节摘要:Combiner 是运行在 Map 任务溢写环节的本地 Reduce:相同键的值先在本任务内归并一次,再落盘、再传输。它不改变最终结果,只压缩中间数据量。合法性条件是归并操作满足"先局部合并再全局合并与一次全局合并等价"。用对了,Shuffle 流量数量级下降;用错了,结果悄悄出错。
统计一个 100GB 日志里每个 URL 的访问次数。假设 200 个 Map 任务,日志里 URL 总数不多(10 万个),但每个任务都要为自己见过的每个 URL 输出一条记录。极端均匀的假设下,每个任务输出 50 万条"URL, 1"(平均每条 20 字节左右),全集群中间数据约 200 × 10MB = 2GB——听起来还能忍。
换成统计每个用户的访问次数,用户数一千万呢?每个 Map 任务输出千万条记录,中间数据 200 × 数百MB,数十GB 的 Shuffle 传输,而最终 reduce 后的结果可能只有几百MB。中间结果比最终结果膨胀百倍,这是词频类计算的常态:值先以"1"的形式满天飞,聚拢后才变成真正的数。
Combiner 的思路朴素:既然每个 Map 任务的输出里就有大量同键记录,为什么不在本地先把它们加起来?加了 Combiner 后,每个任务对每个键最多输出一条(本任务的局部和),中间数据从"记录数 × 单条大小"塌缩成"去重键数 × 单条大小",上面第二个例子的 Shuffle 流量直接砍掉两个数量级。
第 3.1 节讲过溢写流程:环形缓冲区写满八成 → 锁定区排序 → 溢写。Combiner 就插在"排序之后、写盘之前":排序让相同键的记录相邻,相邻就能就地归并。每次溢写都会执行一遍,任务收尾归并多个分段时还会再执行。所以它不是"任务结束后跑一次",而是穿插在溢写循环里的本地 Reduce——这也解释了为什么它必须足够廉价(纯内存计算,通常是求和、取最值、本地拼接)。
判断标准一句话:局部先合并、再全局合并,必须与直接全局合并等价。数学上是结合律加上合并函数与最终函数的可复用性。展开成可操作的检查:
正反两个经典案例。
正例:求和。 WordCount 的 reduce 是求和,局部先加再全局加,结果不变。Combiner 直接复用 reducer 类:
job.setReducerClass(IntSumReducer.class); job.setCombinerClass(IntSumReducer.class); // 同一个类 求和满足结合律
反例:求平均数。 算每个班级的平均分,reduce 端需要"总分除以人数"。如果 Combiner 直接算局部平均,三个任务分别输出 70、80、90,reduce 端无论怎么合并都不等于真实平均。正确做法是让 Combiner 合并"分数和与人数"这个可结合的中间形态,reduce 端再相除:
// map 输出 班级 与 分数 // Combiner 与 Reducer 合并成 班级 与 总分和人数 的复合值 public class AvgCombiner extends Reducer<Text, Text, Text, Text> { protected void reduce(Text cls, Iterable<Text> vals, Context ctx) throws IOException, InterruptedException { long sum = 0, cnt = 0; for (Text v : vals) { sum += Long.parseLong(v.toString()); cnt++; } ctx.write(cls, new Text(sum + "," + cnt)); // 保留后续所需信息 } } // reduce 端 拿到若干 总分和人数 再累加后相除 才是真正的平均
这个案例的教训有普遍意义:不是"业务求什么 Combiner 就合并什么",而是"找到一个满足结合律的中间表示"。求平均的结合形态是和与计数,求中位数的结合形态难以构造(必须保留全部样本),所以中位数类计算用不了 Combiner,只能靠改变键的粒度或两阶段作业绕行。
一个常被混淆的点:Combiner 只是优化手段,框架不保证它一定执行(输出量小时可能不值得)、也不保证执行次数(溢写几次跑几次)。因此:
这与查询引擎里"聚合下推必须语义等价才能下"是同一条原则:优化不能改变可观察的结果。
Combiner 的收益上限由键的重复度决定。每个 Map 任务内部,键的"任务内去重数 ÷ 任务内记录数"越接近 0,压缩越猛(词频、UV 这类高重复场景收益最大);键接近唯一(如按用户 ID 明细输出)则几乎无收益,白付一次归并的 CPU。它也无法解决跨任务的重复——同一键在 200 个任务里仍会各留一条,这部分要靠 Reduce 端或两阶段聚合消化。
中间数据量公式示意 无 Combiner map 任务内记录数 × 单条大小 有 Combiner 任务内去重键数 × 合并值大小 收益比 约等于 平均每键重复次数
顺带一提,Hive 用户其实天天在用 Combiner 思想——映射端聚合(map side aggregation)就是把 group by 的局部合并推到 Map 端,背后挂的正是这个机制。
至此第二幕落幕:数据已在每个 Map 端排好序、分好区、可 optional 地预聚合。下一章进入全书最重的一章——第三幕,汇流。