周一早上,连锁超市的运营同事来找你:昨晚的促销覆盖了几千家门店,他想在午会之前拿到每种商品的销量、销售额和客单价。订单没有整齐地待在一台机器里,而是散在许多数据节点上,总量已经大到无法一次装进单机内存。
最直接的办法,是把全部订单拉到一台中央服务器再计算。可真正开始传输后,你会发现计算还没忙起来,网络先被塞满了;中央服务器既要接收数据,又要排序和汇总,很快也成了新的瓶颈。给中央服务器换一台更贵的机器,只是把问题往后推,并没有改变单点处理的结构。
MapReduce 换了一个思路:我们不急着搬动原始数据,而是把同一段计算送到数据所在的节点。每个节点先处理手边那一份订单,只把体积小得多的中间结果交给后续环节。这样一来,几十台机器可以同时开工,某台机器失败时,对应的小任务也能被重新安排。

MapReduce 最值得记住的并不是两个函数的名字,而是它对分布式工作的切分方式:能独立处理的记录尽量并行,必须碰头的数据按键聚拢,任务调度、网络传输和失败重试交给执行框架管理。
一个 MapReduce 作业通常会把输入和输出都看成键值对。你需要描述两件核心工作:map 如何把一条输入记录变成零个、一个或多个中间键值对;reduce 如何把同一个中间键对应的一组值压缩成结果。至于输入怎样切块、任务分给哪台机器、同键数据怎样汇合、失败任务怎样重跑,通常由框架负责。
这个分工很重要。我们不是把一段普通循环原封不动地丢到集群上,而是在给框架一份可以并行调度的计算说明。只有先把问题改写成“逐条转换”和“按键归并”,框架才知道哪些工作互不依赖,哪些数据最终必须见面。
下面这张交互卡片把完整路径拆成五站。你可以按顺序查看,也可以直接跳到某个阶段,观察数据的形态怎样变化。
回到销售统计。原始订单可以长这样:
{
"orderId": "O-1048",
"storeId": "S-23",
"items": [
{ "product": "手机", "quantity": 2, "unitPrice": 6000 },
{ "product": "耳机", "quantity": 1, "unitPrice": 500 }
]
}如果目标是计算每种商品的销量和销售额,map 会展开订单明细,并把商品编号放在键的位置。值里保留后续求和需要的数量和金额:
function map(order, emit) {
for (const item of order.items) {
emit(item.product, {
quantity: item.quantity,
amount: item.quantity * item.unitPrice,
})
}
}这条订单会产生两条中间记录:
手机 -> { quantity: 2, amount: 12000 }
耳机 -> { quantity: 1, amount: 500 }reduce 收到的已经不是一张订单,而是某个商品的全部局部统计。它把同键值相加,再输出商品级结果:
function reduce(product, values) {
const result = { quantity: 0, amount: 0 }
for (const value of values) {
result.quantity += value.quantity
result.amount += value.amount
}
return result
}
这里有一个容易忽略的边界:MapReduce 保证某个键的值会聚在一起,却不会让一个 reduce 同时看见所有键。你可以很自然地算“每种商品的总销售额”,但如果要从全部商品中选出全站销售额最高的十种商品,通常还要再接一个阶段。
映射任务可以在数据本地独立运行,真正需要跨机器协调的地方,是把同键记录送到同一个归并任务。这个过程通常包含三件事:分区器决定目的地,洗牌负责搬运各节点产生的对应分区,排序与分组把相同键组织成一组。
假设一共有三个归并任务,默认的思路常常是对键做哈希,再对分区数取余:
function getPartition(key, reducerCount) {
return positiveHash(key) % reducerCount
}这个函数最重要的约束是确定性。同一个商品键无论从哪个映射节点产生,都必须算出同一个分区号。否则“手机”的销售额会散落在多个归并结果里,我们拿到的只是几个互不知情的小计。
增加归并任务,通常能提高并行度,也能把单个任务的输入压小。但每个任务都要启动、调度、拉取数据并写出文件,数量过多会让管理开销和碎片文件一起上涨。反过来,如果归并任务太少,单个任务会处理过多数据,集群里的其他机器也可能闲着。
更麻烦的是数据倾斜。促销期间,“爆款手机”可能占全部明细的一半。即使哈希分区从数学上看很均匀,这一个热点键也必须进入同一个归并任务,于是一个任务排起长队,其他任务早早结束。作业的完成时间最终取决于最慢的那一个,而不是平均速度。

处理热点键时,我们可以把一次归并改成两次。第一阶段给热点键附加小范围随机盐,例如把 爆款手机 暂时拆成 爆款手机#0 到 爆款手机#7,让八个局部键分散计算;第二阶段去掉盐,再把八个小计合起来。它用多一个阶段换来更均衡的负载,适合求和、计数这类容易再次合并的统计。
“相同键必须到同一个归并任务”与“负载要尽量均匀”之间可能发生冲突。面对热点键,不能只盯着分区数量;还要看每个键的数据量,并判断这个聚合是否允许先拆开、再合并。
映射任务输出以后,如果立刻把每一条中间记录送上网络,传输量可能非常大。假设某节点本地就产生了一百万条 手机 -> 1,最终目标只是计数,那么先在节点内把它们合成 手机 -> 1000000,网络上就只需要搬一条小计。这一步常被称为合并器。
合并器像一次本地预归并,却不能被当成“保证执行的迷你 Reduce”。执行框架可能不调用它,也可能调用一次或多次,还可能在不同批次上调用。因此,开启或关闭合并器,最终答案都必须一样。

哪些操作适合合并?求和、计数、最大值和最小值通常很自然,因为局部结果还能使用同样的规则继续合并。哪些操作危险?直接计算平均数就是最经典的陷阱。
假设甲门店有 2 笔订单,平均购买 3 台;乙门店有 8 笔订单,平均购买 5 台。把两个平均数再平均得到 4 台,但真实平均是 (2×3 + 8×5) ÷ 10 = 4.6 台。错因不是除法本身,而是局部平均丢掉了各组的笔数。
试着调整下面的“节点数”和“每个节点的记录数”。它会估算求和场景中,合并器能把多少中间记录留在本地。
平均数不是不能并行算,而是不能只保留“平均值”。我们可以在每个局部结果里保留 {sum, count},合并时分别相加,最后才做一次除法:
function map(order, emit) {
emit(order.product, { sum: order.quantity, count: 1 })
}
function combine(product, states) {
return states.reduce(
(total, state) => ({
sum: total.sum + state.sum,
count: total.count + state.count,
}),
{ sum: 0, count:

这套思路还能推广:
判断一个计算能否高效分布式执行时,先别急着写代码。问问自己:局部结果需要保留什么状态,两个局部状态能不能继续合并,合并顺序改变会不会影响答案。能回答这三个问题,合并器和多阶段设计通常就清楚了。
运营同事接着提出新问题:他不只要每月销售额,还要找出每个月同比增长最快的十种商品。第一阶段可以按“商品 + 年月”聚合销售额,但这个阶段的每个归并任务只知道自己的商品月份键,无法直接在所有商品之间做全局排名。
我们可以把任务拆成两段:
第一阶段从订单中发出以“商品、年份、月份”组成的复合键,归并后得到每种商品每个月的销售额。输出不再是海量明细,而是一份紧凑的月度汇总。
第二阶段读取月度汇总,把同一商品、同一月份的相邻年份数据放到一起,计算同比变化。需要全局前十时,可以先做分区内前十,再做一次很小的全局归并。
如果季度报表和年度报表也要复用月度结果,就把第一阶段输出保存为有明确口径、可追踪版本的中间数据。这样其他作业不必反复扫描原始订单。

多阶段的代价也很直接:每增加一段,通常就会多一次调度、读写和故障边界。中间结果还涉及命名、版本、过期清理和数据口径。拆分不是越细越好;我们要在“单个阶段是否可表达”“是否出现数据倾斜”“中间结果能否复用”和“额外读写成本”之间取舍。
MapReduce 最舒服的任务,是同键数据的聚合。要把订单与商品档案做关联时,我们得想办法让同一商品的订单记录和商品记录在归并端相遇。小表可以复制到每个映射节点,在本地完成查找;两张大表则可以都以商品编号为键发出,在归并端做连接。前一种方式减少洗牌却要求小表放得下,后一种方式更通用但网络和倾斜成本更高。
这也解释了为什么真实的数据处理平台会在 MapReduce 之上提供更高层的查询和数据流接口。开发者描述过滤、分组、连接和窗口,优化器再决定如何组合底层阶段。理解 MapReduce 仍然有价值,因为一旦作业出现慢节点、热点键或巨大洗牌,你能看懂成本是从哪里冒出来的。
昨天的报表已经算完,今天新增十万张订单。若指标是销量总和,我们不必重新扫描几年的历史订单:先对今日新增部分做一次 MapReduce,再把新增小计与昨日结果相加即可。这个思路适合只追加的数据,也适合能用增量状态表达的统计。
可一旦出现退货、改价或订单删除,简单相加就不够了。系统至少要知道哪条旧记录受到了影响,并决定发出一条可抵消的负增量,还是重算对应商品、日期或分区。对于最大值也要格外小心:如果被删除的记录恰好是当前最大值,仅从“旧最大值”这个结果无法推回第二大值。
增量作业还要守住三条工程边界:
下面的小工具把业务变化和聚合类型放到一起。选择条件后,它会给出更稳妥的更新策略,并说明为什么。
分布式集群里,机器故障和慢任务不是例外,而是日常。映射任务读取固定输入并产生中间结果,失败后通常可以换一台机器重跑;归并任务也可以重新拉取中间分区再计算。框架通过重试隐藏了不少故障细节,但前提是同一任务执行两次,不会把业务世界永久改变两次。
如果你在 map 里直接调用支付接口、发送短信,或给外部数据库做一次不可重复的加一,任务重试就可能造成重复扣款、重复通知或重复计数。更稳妥的做法,是让映射与归并尽量成为确定性的纯计算,把输出交给有提交协议、幂等键或事务保护的写入环节。
慢任务也会拖住整项作业。有些框架会为明显落后的任务启动一个备份副本,谁先完成就采用谁的结果。这能缓解机器偶发变慢,却同样要求输出提交可控;否则两个副本都产生外部副作用,结果就不再可靠。
不要把“框架会重试”理解成“任何代码都能安全重跑”。MapReduce 能恢复的是受管理的任务与数据流;外部接口调用、非幂等写入和依赖当前时间或随机数的逻辑,都需要你自己设计去重、幂等或提交边界。
MapReduce 并不专属于 NoSQL。它最初解决的是大规模分布式数据处理问题,后来被一些文档数据库和大数据系统吸收,用于在数据节点附近执行自定义聚合。它的优势是表达能力强、任务容易并行和重试;代价是中间键值对、洗牌、排序与多阶段读写可能很重,自定义代码也不容易被查询优化器改写。
在现代 MongoDB 中,旧式 map-reduce 接口已经不再是新业务的优先选择,常见的过滤、分组、连接和结果写回应优先使用聚合管道。聚合管道的阶段含义更明确,优化空间也更大。只有当我们阅读旧系统、迁移历史作业,或研究分布式聚合的底层思路时,才会直接遇到数据库内的 map-reduce 代码。
这并不意味着这一章过时了。相反,当你使用聚合管道、批处理引擎或流式计算系统时,仍会反复遇到相同的问题:数据是否能本地处理,哪些记录必须按键聚拢,中间状态能否组合,热点键会不会拖慢一个分区,失败重试是否会产生重复。MapReduce 把这些问题压缩成了一套足够清楚的骨架。
先写清最终问题的粒度。你要的是“每种商品每天的销售额”,还是“全站每天的前十商品”?粒度决定中间键,也决定一个归并任务能看到哪些数据。
再写出一条原始记录经过映射后的具体键值对。若你无法用两三条样例演算,就先不要急着部署到集群。
检查键的分布。除了键的种类数,还要估算最大键的数据量;发现热点后,判断能否加盐拆分再二次汇总。
设计可组合状态。求平均保留总和与个数,前 N 名保留候选,去重考虑集合大小和可接受误差。合并器是否执行,不应改变最终语义。
一家内容平台保存了用户播放事件,每条记录包含用户编号、视频编号、播放秒数、事件时间和设备类型。现在要统计每个视频每天的总播放时长。
MapReduce 把一个庞大的分布式任务拆成了清晰的责任链:输入切分让工作可以并行,映射把记录转换成中间键值对,分区保证同键数据去往同一归并任务,洗牌与排序完成跨节点聚拢,归并再把一组值压成结果。
真正决定作业质量的,往往不是 map 和 reduce 函数各写了几行,而是中间键是否选对、数据是否倾斜、状态能否安全组合、阶段是否拆得合理,以及失败重试后结果是否仍然只计算一次。把平均数改写成“总和 + 个数”,把热点键拆成局部汇总再二次合并,把多目标分析变成可复用的阶段流水,这些都是同一种思维:先找到能独立处理的部分,再只让真正需要见面的数据碰头。
当数据持续增长时,我们还可以只处理增量,但必须同时处理重复、遗漏、修改删除和口径版本。掌握这些边界以后,你看到的不再只是一个历史上的计算框架,而是一套仍然活跃在现代数据平台里的分布式聚合方法。
最后补上工程边界:输入范围、失败重试、幂等输出、异常记录、指标计数、中间结果版本与清理策略。能算对只是第一步,能稳定重复地算对才算完成。