3.3 Combiner 本地预聚合


文档摘要

3.3 Combiner 本地预聚合 本节摘要:Combiner 是运行在 Map 任务溢写环节的本地 Reduce:相同键的值先在本任务内归并一次,再落盘、再传输。它不改变最终结果,只压缩中间数据量。合法性条件是归并操作满足"先局部合并再全局合并与一次全局合并等价"。用对了,Shuffle 流量数量级下降;用错了,结果悄悄出错。 先看一组流量数字 统计一个 100GB 日志里每个 URL 的访问次数。假设 200 个 Map 任务,日志里 URL 总数不多(10 万个),但每个任务都要为自己见过的每个 URL 输出一条记录。极端均匀的假设下,每个任务输出 50 万条"URL, 1"(平均每条 20 字节左右),全集群中间数据约 200 × 10MB = 2GB——听起来还能忍。

3.3 Combiner 本地预聚合

本节摘要: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——这也解释了为什么它必须足够廉价(纯内存计算,通常是求和、取最值、本地拼接)。

合法性条件:什么时候能用

判断标准一句话:局部先合并、再全局合并,必须与直接全局合并等价。数学上是结合律加上合并函数与最终函数的可复用性。展开成可操作的检查:

  1. 合并函数与 reduce 函数相同,或功能上是它的"前缀";
  2. 交换合并顺序不影响结果(求和、求最值、求条数、集合并都满足);
  3. 合并不会销毁后续需要的信息。

正反两个经典案例。

正例:求和。 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 不是 SQL 的聚合推下推

一个常被混淆的点:Combiner 只是优化手段,框架不保证它一定执行(输出量小时可能不值得)、也不保证执行次数(溢写几次跑几次)。因此:

  • 不能把业务正确性押在 Combiner 上(例如"去重必须靠 Combiner 完成"是错误设计);
  • 不能在 Combiner 里做有外部副作用的操作(写数据库、发请求),它可能被调零次或多次;
  • 它的唯一合法使命是减量,正确性必须完全由 map 加 reduce 保证。

这与查询引擎里"聚合下推必须语义等价才能下"是同一条原则:优化不能改变可观察的结果

效果与限制的量级感

Combiner 的收益上限由键的重复度决定。每个 Map 任务内部,键的"任务内去重数 ÷ 任务内记录数"越接近 0,压缩越猛(词频、UV 这类高重复场景收益最大);键接近唯一(如按用户 ID 明细输出)则几乎无收益,白付一次归并的 CPU。它也无法解决跨任务的重复——同一键在 200 个任务里仍会各留一条,这部分要靠 Reduce 端或两阶段聚合消化。

中间数据量公式示意 无 Combiner map 任务内记录数 × 单条大小 有 Combiner 任务内去重键数 × 合并值大小 收益比 约等于 平均每键重复次数

顺带一提,Hive 用户其实天天在用 Combiner 思想——映射端聚合(map side aggregation)就是把 group by 的局部合并推到 Map 端,背后挂的正是这个机制。

本节要点回顾

  • 定位:溢写环节的本地 Reduce,排序后相邻同键就地归并,只减流量不改结果。
  • 合法性:先局部合并再全局合并须等价于直接全局合并;求和、计数、最值、集合并安全。
  • 反例:平均数不能直接局部平均,要合并"和与计数"的中间形态;中位数无结合形态,不可用。
  • 不保证执行:次数与有无都由框架定,业务正确性不能依赖它,禁止副作用。
  • 收益公式:压缩比约等于任务内平均每键重复次数,跨任务重复管不着。

至此第二幕落幕:数据已在每个 Map 端排好序、分好区、可 optional 地预聚合。下一章进入全书最重的一章——第三幕,汇流。


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