很多人在面试里背 MapReduce 的时候,只会说一句话:“就是先把数据拆开,map 一下,然后再 reduce 合并一下。”这话没错,但就跟你把火锅解释成“把菜放水里煮”一样,听起来对,实际等于没说。真到了线上,一个作业跑了几个小时不动,或者某个 Reduce 任务卡住,GC 个不停,你根本不知道你对那个作业动了什么手脚、改哪个参数能救回来。我自己带项目的这些年,身边几乎所有追过 MapReduce 底层的同学,最后都有一个共识:搞懂 MapReduce 不是在背 API,而是要沿着一条数据从 HDFS 出来、经过 Map、溢写、排序、洗牌、拉取、归并、Reduce、再落回 HDFS 的完整轨迹,把每个环节的触发条件、默认参数、为什么这样设计的逻辑全部理清,才算真懂。这篇我就用最细的方式,把这条轨迹一条一条拆给你看,配合 WordCount 实例和 HDFS 综合实训的思路,尽量让新人少走弯路。
1. MapReduce 到底解决的是哪门子问题
1.1 没有 MapReduce 之前,我们怎么处理海量数据
先把时间退回到数据量还停留在“单机能扛住”的年代。那时候处理一批数据,无非是把文件读进内存,写个循环,遍历、统计、输出,完事。几千兆的数据没问题,一台性能好点的机器也能压得住。
但是当数据量变成 TB 甚至 PB 级别的时候,单机就彻底没戏了。你不光要把数据切成很多块,还得分别放到很多台机器上,然后让这些机器同时干活,再把结果拼在一起。听起来简单,做起来全是坑:哪台机器负责哪一块?A 机器算完了,B 机器崩溃了怎么办?网络传输又慢又容易断,中间结果往哪里放?各种结果合并的时候会不会互相冲突?
MapReduce 这个编程模型,本质上就是把上面这些“分布式系统工程师的脏活累活”都封装到框架层,让业务开发者只需要关心两个函数:map 和 reduce。你写一个 map 函数处理一条记录,写一个 reduce 函数合并一组相同 key 的数据,剩下的事——切分、调度、容错、网络传输、排序、合并——全部交给框架。
所以第一个要建立的认知是:MapReduce 的核心价值,不是多么高深的算法,而是“把分布式计算的复杂度从用户手里拿走”。那几年它能在工业界统治一批批大数据框架,靠的也不是性能碾压别人,而是“你不需要懂分布式也能算大数据”的极低门槛。
1.2 “分而治之”说的不是一句口号,是四层分工
很多人把“分而治之”理解成“拆数据 + 合并结果”,这个理解不错,但太粗了。真正落到框架实现上,至少可以分成四层:
- 作业层:一个完整的业务需求,对应一个 Job。
- 任务层:一个 Job 被拆成很多个 Map 任务和若干 Reduce 任务,这些任务并行跑在不同节点上。
- 切片层:输入数据被 InputFormat 计算成多个 InputSplit,每个 InputSplit 恰好对应一个 Map 任务。
- 记录层:每个切片内部又被 RecordReader 逐行读成
<key, value>记录,map 函数就是逐条处理这些记录。
这四层的关系,很像开一家大型连锁餐厅:老板(客户端)把订单下达给总店;总店大堂经理(ApplicationMaster)把订单拆成一道道菜,分配给各个档口;每个档口的厨师(Map 任务)拿到食材后按菜单做菜;最后传菜员(Reducer)把所有人做好的菜按桌号归位、拼盘上桌。
这里值得多说一句的是,MapReduce 最初在 Hadoop 里的角色分工其实经历过一次演进。早期叫 MRv1,有一个全局的 JobTracker 负责所有作业的调度和监控,节点上跑 TaskTracker 接受指令。这个架构最大的槽点是单点故障和性能瓶颈,一个 JobTracker 扛所有作业,集群一大必然撑不住。后来演进到 MRv2,也就是大家熟知的 YARN 架构,把“资源管理”和“作业管理”彻底拆开。ResourceManager 只负责管全局资源分配,每个作业单独拉起一个 ApplicationMaster 来管这个作业的任务拆分、调度、重试。这个改动很关键,它让一个集群可以跑多种计算框架,也让 MapReduce 的作业不再受“一个 Tracker 扛所有”的限制。
1.3 Map 和 Reduce 之间还有一大批“看不见”的角色
如果只盯着 Mapper 和 Reducer 这两个类,你永远理解不了 MapReduce 全貌。在一条数据从输入到输出的完整旅程里,真正干活的还包括:
- InputFormat:负责校验输入目录、计算分片、提供 RecordReader。
- Partitioner:决定 Map 输出的每条记录进入哪一个 Reduce 分区。默认是
(key.hashCode() & Integer.MAX_VALUE) % numReduceTasks。 - 环形缓冲区:Map 输出的暂存空间,默认 100MB,满了按比例溢写到磁盘。
- Combiner:跑在 Map 节点的“局部 Reducer”,能在数据传出去之前先做一次聚合,减少网络传输。
- Shuffle:Map 和 Reduce 之间的数据搬运过程,是整个框架最复杂、也最值得深挖的一段。
- OutputFormat:负责把计算结果写到目标存储,最常用的是 hdfs 上的
TextOutputFormat。
搞清楚这些角色各自的位置和职责,才能往下走。接下来我就按一条数据从客户端提交作业开始,一直到 HDFS 上看到输出文件,把这个过程整个串一遍。
2. 一条数据从进入集群到写出结果的完整路线图
2.1 作业提交:Client 先做一堆事情才轮到 ResourceManager
很多人以为作业提交就是敲一条hadoop jar xxx.jar命令,然后框架就自动跑起来了。实际上,客户端在提交之前会被迫做很多“重活”,这恰恰是很多人没注意到的。
第一步,客户端会先检查输入输出路径是否合法,读入作业的各种配置,比如 mapper 类、reducer 类、输出 key/value 类型、Reduce 数量等。然后重点来了:客户端会调用 InputFormat 的 getSplits 方法,把输入数据切成若干 InputSplit,并把分片元数据直接算好。注意,这一步是在客户端完成的,不是在集群上完成的。也就是说,你提交一个作业之前,你的客户端机器已经知道“这个作业要分成多少个 Map 任务”。
第二步,客户端把作业所需要的资源(jar 包、配置文件、计算出来的分片元数据)上传到一个 HDFS 目录里,这个目录通常长这样:/user/xxx/.staging/job_id/。上传到 HDFS 的原因也很朴素:后续 ApplicationMaster 可能在集群任意节点启动,它必须能拿到这些资源;同时这也是作业容错的一部分,万一 AM 挂掉重启,还能从 HDFS 再把资源拉回来。
第三步,客户端向 ResourceManager 发一个submitApplication请求。ResourceManager 收到后,会为这个作业分配一个 ApplicationMaster 容器,让 NodeManager 在某个节点上把 AM 进程拉起来。从这时候开始,客户端就不直接指挥任务了,它变成了“监工”,只通过 AM 的进度报告来判断作业是不是跑完了。
这里要插一个实时感受:为什么老会看到有人说“小文件多导致作业慢”?因为客户端切分时,每个小文件至少会生成一个 InputSplit,分片数量越多,AM 要管理的 Map 任务数就越多,调度开销、资源开销、启动开销全部成倍增长。几百个 100KB 的小文件,分片数和资源浪费能让你怀疑人生。
2.2 分片的计算决定了 Map 任务的“粒度”
分片是 MapReduce 里“工作量切分”的最小单位,一个 InputSplit 对应一个 Map 任务。分片不是把数据物理复制一份,它只是一个逻辑概念,保存的是元数据:数据在哪个文件、起始偏移量、长度。真正的数据还在 HDFS 的 block 里。
默认情况下,分片大小跟 HDFS 的 block size 一致,就是 128MB。这个规则不是拍脑袋定的,而是由公式决定的:
splitSize = max(minSize, min(maxSize, blockSize))其中,minSize对应参数mapreduce.input.fileinputformat.split.minsize,默认 1;maxSize对应mapreduce.input.fileinputformat.split.maxsize,默认 Long.MAX_VALUE。因为blockSize在两者之间,所以默认分片大小就是 128MB。如果你手动把maxSize调小,比如改成 64MB,那么分片会更小、Map 任务更多,并行度更高,但调度开销也会更大。
还需要注意一个特殊场景:压缩文件能不能切分,取决于压缩格式。像 Gzip 这种不支持随机读取的格式,文件再大,也只能生成一个分片,由一个 Map 任务读完整份文件。这会导致严重的单 Map 任务瓶颈。ZIP 和 LZO 相对好一点,但也要看是否建了索引。实际项目中选压缩格式时,这个“是否可切分”的属性往往比压缩率还要关键。
一个分片的数据读出来后,由 RecordReader 按行解析成<key, value>。默认的 TextInputFormat 会把每行开头的字节偏移量作为 key,这行文本本身作为 value。于是,Map 函数看到的世界就是:一堆“偏移量 + 一行业务日志”的记录。
2.3 Map 函数执行完成后,输出不会立刻落盘
Map 函数内部拿到一条记录后,会执行你写的业务逻辑,然后调用context.write(key, value)输出中间结果。这个中间结果不会像新手想的那样“直接写到磁盘再传给 Reduce”,而是先进入一个环形内存缓冲区。
这个缓冲区默认大小是 100MB,由参数mapreduce.task.io.sort.mb控制。数据写进去之后,会先做两件事:分区和排序。每条记录会根据 key 经过 Partitioner 算出要发往哪个 Reduce 分区;在分区内部,又按照 key 的字典序排好。这里用的是快排,只在内存中做,速度很快。
当缓冲区的写入量达到阈值——默认 80%,即 80MB,由mapreduce.map.sort.spill.percent控制——一个后台线程就开始把缓冲区里的数据溢写到磁盘,生成一个临时文件,叫 spill 文件。注意一个细节:溢写线程不会等到缓冲区全满才行动,因为全满时 Map 就会阻塞,等溢写完成后才能继续写,这样计算效率会断崖下跌。80% 这个阈值是在“充分利用内存”和“留出余量避免阻塞”之间找的平衡点。
一个 Map 任务处理完所有数据后,可能会生成多个 spill 文件。这些文件最终会被归并(merge)成一个大的输出文件,同时按照分区和 key 排好序。如果配置了 Combiner,会在溢写和归并的过程中执行 Combiner,提前合并相同 key 的局部结果。归并完成后,Map 任务还会告诉 ApplicationMaster:“我的输出文件在这里,你记一下位置。”但是,Map 的输出文件不会立刻删掉,它是放在本地磁盘而不是 HDFS 上,因为 Reduce 后面还要来拉取。
这就顺势带出一个调优直觉:Map 输出如果很大,本地磁盘 IO 就会很重。所以实际生产里,经常会对 Map 输出做压缩,既减少本地磁盘占用,也减少后面 Reduce 拉取时的网络开销。
2.4 Shuffle 与 Sort:Map 和 Reduce 之间的“隐藏高速公路”
Shuffle 是整个 MapReduce 里最精华、最容易把新手绕晕的一段。简单说,Shuffle 就是“把 Map 输出的数据搬运到 Reduce 节点的过程”,但“搬运”这两个字背后,藏着非常多的细节。
整个 Shuffle 可以拆成两段:Map 侧的 shuffle 和 Reduce 侧的 shuffle。
Map 侧,前面说到的分区、排序、溢写、归并,其实都属于 shuffle 的一部分。Map 任务跑完以后,它的输出是按分区排列好的文件。每个 Reduce 任务需要的数据,是其中某一个或某几个分区的数据。
Reduce 侧要复杂一点。Reduce 任务启动后,并不会等到所有 Map 任务跑完才开始工作,它会尽早启动一个或多个拉取线程(默认 5 个,mapreduce.reduce.shuffle.parallelcopies),循环向 ApplicationMaster 询问“有哪些 Map 输出已经 ready 了,位置在哪里”,然后拿着位置信息,通过 HTTP 协议把对应分区的数据拉到本地。这也是为什么你在 YARN 的日志里经常能看到 reduce 的状态一直停在copy阶段——它在等最后一个慢 Map 的碎片。
拉回来的数据先放在 Reduce 节点的内存缓冲区,缓冲区满了就溢写到磁盘,跟 Map 侧的思路几乎一样。等到所有 Map 输出都被拉完,Reduce 节点会对内存和磁盘上的所有数据做一次归并排序。这一步归并完成后,数据就变成“同一分区的数据已经合并在一起,并且分区内相同 key 的记录是相邻的”。
这里必须要提一个特别容易被误解的点:Reduce 端拿到的不是“按 key 排好序的一堆数据”,然后自己去遍历找相同 key。真正的实现是,Reduce 节点的输入数据经过归并排序后,框架会把相邻 key 相同的记录分组,每组调用一次 reduce 函数。你写的 reduce 方法签名里那个Iterable<IntWritable> values,其实就是一组相同 key 对应的所有 value 的迭代器。
为什么说排序是必需品?因为 Reduce 要对相同 key 的 value 做聚合,如果相同 key 的数据分散在几千万条记录里,靠哈希表去维护,内存会爆炸;而把它们排到一起,Reduce 只需顺序扫描一遍,就能自然地按 key 分组处理。这也是为什么你会听到“MapReduce 的排序是框架自带的、无关业务”的说法。
在特殊场景下,你还可以通过自定义 Comparator 控制“排序规则”和“分组规则”,实现二次排序。最典型的需求是:按 key 分组,但组内按 value 排好序。比如“每个用户的所有订单按时间升序排列”,就需要让 key 由“用户 + 时间”组成,但分组只按用户分,排序按用户和时间同时排。
2.5 Reduce 函数真正拿到的是“按 key 分组好的迭代器”
Reduce 阶段看起来比 Map 简单,但它隐藏着一个很反直觉的设计:reduce 函数拿到的 values,不是一次性加载进内存的集合,而是一个迭代器。也就是说,框架不会把所有相同 key 的 value 都堆到内存里再交给 reduce 函数,而是边遍历边喂给你。
为什么这么设计?因为一个 key 的 value 数量可能是百万甚至千万级别。如果全装进内存,内存分分钟被撑爆。Iterator 的本质是懒加载,让你需要多少处理多少,处理完就丢,内存开销跟单个 key 的数据总量无关,只跟“你同时在手里攥着多少数据”有关。
这给写代码的人提了个醒:如果你在 reduce 里做的是类似List<Integer> all = new ArrayList<>(); for (IntWritable val : values) { all.add(val.get()); }的操作,把迭代器里的值全存到一个自定义集合里,那“Iterator 防 OOM”的设计就形同虚设了。真碰到 key 特别多的情况,你这是主动踩雷。更好的做法是:边遍历边维护中间状态,例如累加器、最大值、TopN 堆等。
Reduce 处理完成之后,结果同样不会直接写到 HDFS,它先写到节点本地临时文件。只有当整个作业成功提交后,框架才会把这些临时文件移动到最终输出目录。这个“最后一步才算成功”的设计,是为了保证最终输出的一致性,避免半途看到残缺结果。如果作业中途失败,输出的临时文件会被清理掉,重新跑的时候不会污染数据。
2.6 收尾阶段:Commit、清理与用户感知
Reduce 全部完成后,ApplicationMaster 会向 ResourceManager 报告作业成功。随后,客户端通过轮询 AM 的作业状态,发现jobFinished后,会打印出一行“Map-Reduce Finished”的信息,并展示一组计数器统计。
同时框架会清理中间产物:Map 输出的本地文件会被删除,staging 目录里的临时文件也会被清理。最终 HDFS 上看到的就是输出目录里的part-r-00000、part-r-00001这类文件,有多少个part-r,就说明这次作业用了多少个 Reduce 任务。
有个细节值得提一下:Map 任务数量一般不由我们显式指定,而是由分片数量决定;而 Reduce 任务数量是可以通过job.setNumReduceTasks(n)或者mapreduce.job.reduces参数指定的。Reduce 数设置多少,直接影响数据分布和最终输出文件数量。设置的太大,每个 Reduce 拉取的数据量少,但启动和调度开销大;太小则并行度不够,集群资源大多闲置。实践中通常结合数据量和集群规模先估算一个范围,再通过测试对比调整。
3. 为什么设计者要这样做:那些关键参数背后的设计逻辑
3.1 为什么分片默认 128MB,而不是越小越好
分片大小的默认值之所以是 128MB,核心是平衡“并行度”和“开销”之间的矛盾。
如果分片太小,比如 1MB,一个 1GB 的文件会被切成 1024 个分片,也就是 1024 个 Map 任务。每个 Map 任务启动都需要申请容器、加载 jar、初始化 JVM,这个启动过程的开销可能比任务本身还大。这种场景下,你看到的作业运行时间大部分都浪费在了“创建任务”而不是“计算数据”。
如果分片太大,比如 1GB,一个文件只生成 1 个 Map 任务,那即使你有 1000 个计算节点在待命,也只有 1 个节点在干活。更麻烦的是,单个分片过大,处理时间太长,一旦任务失败,重跑的成本也很高。
128MB 这个值,刚好跟 HDFS 默认块大小对齐。这样带来的额外好处是:一个分片通常对应一个本地 block,Map 任务可以优先调度到该 block 所在的节点上执行,数据直接本地读,不用跨网络拉数据,这就是数据本地性(Data Locality)。网上很多性能调优文章让你把maxSplitSize调小以增加并行度,我建议你先想想自己的数据分片现状,别盲目乱调。
3.2 环形缓冲区:为什么是 100MB 配 80% 溢写
环形缓冲区这个名字听起来唬人,其实就是一个首尾相接的字节数组,Map 输出写进内存时,用一块区域放 key/value 数据,再用另一块区域放这些数据在内存中的索引信息。为什么要“环形”?因为内存缓冲区会被反复使用,写满一部分就溢写一部分,溢写完的区域腾出来继续写,像一个循环使用的蓄水池。
默认 100MB 的大小,对于大部分中小作业够用,但如果是 Map 输出很大的作业,100MB 很容易频繁溢写。溢写会触发磁盘 IO,一次溢写就是一次写盘,如果溢写次数很多,任务时间会显著拉长。所以对 Map 输出很大的作业,适当把mapreduce.task.io.sort.mb调到 200MB 或 256MB,往往能看到明显的提速。
80% 溢写阈值,需要理解成“缓冲区写着写着,到 80% 就触发后台溢写,但 Map 还能继续往剩下 20% 写”。如果阈值设成 100%,Map 就必须堵在缓冲区门口等溢写完,计算线程和 IO 线程互相等待,性能惨不忍睹;如果阈值设得太低,比如 20%,缓冲区利用率太低,溢写频繁,也没必要。80% 是个很务实的默认值,既兼顾内存使用率,也考虑了 IO 和计算并行。
3.3 一个让所有人都懵过的点:为什么 Map 输出也要排序
很多初学者卡在同一个问题上:我 Map 输出的 key 是随机无序的,Reduce 直接按 key 聚合不是也可以吗,为什么非要排序?
这里的难点在于,分布式场景下,Reduce 需要拉取的是分布在多台机器上的多个 Map 输出中的同分区数据。如果这些数据不排序,Reduce 把数据聚齐之后,还要自己建一个大哈希表,把所有 key 都塞进去再逐个聚合。数据量一上去,内存必定不够。
排序之后,所有相同 key 的数据在文件里都是连续的一段。Reduce 做归并时,只需要用类似“多路归并”的算法,把多份有序数据流合并成一份有序的大数据流,再顺序扫描,遇到 key 变化就切分组。整个过程的额外内存开销极小,时间复杂度也更漂亮。
所以排序的意义不是“让数据好看”,而是用一次全局有序的代价,换 Reduce 阶段几乎零内存压力的分组能力。
这个设计思路,我建议每个学 MapReduce 的人都记在心里。因为后面学 Spark 的时候,你还会看到 shuffle 里同样强调排序和分区,逻辑是一脉相承的。
3.4 Combiner 什么时候能加,什么时候绝对不能加
Combiner 是个“看起来很美好,用错就翻车”的功能。它在 Map 节点上先做一次局部聚合,减少要传送到 Reduce 的数据量。比如 WordCount 里,Map 输出 1000 条(word, 1),Combiner 在本地先聚合成几条(word, 100),明显减少网络 IO。
但 Combiner 不是随便什么逻辑都能当 Combiner 用的。核心约束是:Combiner 的输入和输出类型必须跟 Reduce 一致,而且 Combiner 的处理逻辑必须多次执行后结果不变。也就是说,它得满足“可交换”和“可结合”的数学性质。
加法满足,乘法满足,求最大值、最小值也满足。但求平均值就不行。举个例子:两个(key, 2)和(key, 4),如果直接聚合一次,平均值是 3;但如果在不同节点上先分别算出平均值 2 和 6,再拿去求平均,得到的是 4,显然不对。还有像去重类逻辑,前面用不对的 Combiner 会把全局结果直接改坏。
所以我的建议是:默认情况下,如果 Reducer 的聚合逻辑简单到“满足交换律和结合律”,比如 sum、max、min,可以放心复用 Reducer 类做 Combiner。但凡是涉及复杂统计、去重、组合计算的业务,宁愿先不加 Combiner,等确认逻辑没问题再优化。
4. 实战:从 WordCount 复现全流程,再到 HDFS 综合实训
4.1 一个能够对应到每个机制的 WordCount
如果只选一个 MapReduce 编程实例来理解整个原理,WordCount 当之无愧。下面是一个标准的实现,我会逐段把它跟前面讲过的机制对应起来,而不是简单贴代码。
public class WordCount { public static class TokenizerMapper extends Mapper<Object, Text, Text, IntWritable> { private final static IntWritable one = new IntWritable(1); private Text word = new Text(); public void map(Object key, Text value, Context context ) throws IOException, InterruptedException { StringTokenizer itr = new StringTokenizer(value.toString()); while (itr.hasMoreTokens()) { word.set(itr.nextToken()); context.write(word, one); } } } public static class IntSumReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private IntWritable result = new IntWritable(); public void reduce(Text key, Iterable<IntWritable> values, Context context ) throws IOException, InterruptedException { int sum = 0; for (IntWritable val : values) { sum += val.get(); } result.set(sum); context.write(key, result); } } public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "word count"); job.setJarByClass(WordCount.class); job.setMapperClass(TokenizerMapper.class); job.setCombinerClass(IntSumReducer.class); job.setReducerClass(IntSumReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }逐段对应来看:
Mapper<Object, Text, Text, IntWritable>四个泛型依次表示:输入 key、输入 value、输出 key、输出 value。输入 key 是行偏移量,所以在 map 函数里几乎不会被用到。context.write(word, one)这行代码执行时,数据进入了 Map 的环形缓冲区,并立即做了分区 + 排序。分区逻辑直接使用默认 HashPartitioner,所以相同 word 会进同一个 Reduce 分区,不同 word 则随机分散。setCombinerClass(IntSumReducer.class)这一行非常关键。因为 WordCount 的聚合是加法,满足结合律交换律,所以我们可以直接复用 Reducer 作为 Combiner。Map 节点上会先把局部单词计数求和,再传给 Reduce,网络传输的数据量会明显下降。setOutputKeyClass和setOutputValueClass同时决定了 Mapper 和 Reducer 的输出类型。如果 Mapper 和 Reducer 的输出类型不一致,还需要额外设置setMapOutputKeyClass和setMapOutputValueClass,这是很多新手容易漏掉的地方,漏掉的后果是作业跑起来才发现类型不匹配。
运行这个作业之前,需要先把输入文件放到 HDFS 上:
hdfs dfs -mkdir -p /wordcount/input hdfs dfs -put words.txt /wordcount/input/words.txt hadoop jar wordcount.jar WordCount /wordcount/input /wordcount/output作业跑完后,查看输出:
hdfs dfs -ls /wordcount/output你大概率会看到一个part-r-00000文件。文件名的r代表 Reduce 输出,数字是 Reduce 任务的编号。如果你没显式设置 Reduce 数量,默认是 1。这也很说明问题:整个作业只有一个 Reduce,所有 Map 输出最终都会拉到一个节点上处理,数据量大时必然慢。真实生产里,我会根据数据量手动把 Reduce 数量调成合适值,比如:
job.setNumReduceTasks(8);4.2 Counter 和日志:你在哪里能看到作业内部正在发生什么
代码写对了,作业跑完了,只是第一步。要判断一个作业“跑得好不好”,Counter 是你最好的显微镜。
作业运行结束后,控制台会打印一堆计数器,比如:
Map input records:Map 读取的总记录数。Map output records:Map 写出的记录数。Spilled Records:溢写到磁盘的记录数。这个值如果比 Map output records 大很多,说明溢写太频繁,内存紧张,需要考虑加大mapreduce.task.io.sort.mb。Combine output records:经过 Combiner 后的记录数。把它和Map output records对比,能算出本地聚合的压缩效果。Reduce shuffle bytes:Reduce 从 Map 拉取的数据量,单位是字节。如果这个值很大,网络传输会是性能瓶颈,考虑开启 Map 输出压缩。Reduce input records:Reducer 实际接收到的记录数。
Counter 的价值在于,它让你不用猜就知道作业内部发生了什么。有一次我帮同事排查一个作业,Map 任务很快就跑完了,但 Reduce 一直卡在 66% 附近不动。打开 Counter 一看,某个 Reduce 的Reduce shuffle bytes比其他 Reduce 大出几百倍,基本就断定是数据倾斜,后来在业务 key 上做了改造,问题立刻缓解。
如果想看更详细的日志,可以用 YARN 的命令行工具:
yarn application -status <application_id> yarn logs -applicationId <application_id> -log_files stdout在 AM 的日志里,你能看到更完整的调度细节:多少个 Map 任务成功、多少个失败重试、每个任务的启动时间、GC 时间、节点分布等。这些都是判断作业健康状况的直接证据。
4.3 HDFS + MapReduce 综合实训怎么做才有含金量
很多学校的综合实训课就是让跑一遍 WordCount,然后就没有然后了。说实话,这种程度离“懂”还差得很远。真正有含金量的实训,建议按下面的思路安排:
第一步,把数据装进 HDFS。用hdfs dfs -put上传一份真实感更强的数据,比如电商订单流水、某段时间的服务器访问日志,而不是一笔带过的小文本。数据量建议至少几个 GB 甚至更大,否则体会不到分布式计算的必要性。
第二步,设计一个有点复杂的业务需求。比如按用户 ID 统计每个用户每月的消费总额,或者按来源 IP 统计每天各时段的访问量。这些需求都能拆分出明确的(key, value)结构,但又比“数单词”更接近真实业务。
第三步,实现代码并设置对比实验。跑一遍默认配置,记录运行时间和 Counter;然后试着调整 Reduce 数量、调整 Map 内存、开启 Combiner,再跑一遍,对比时间差异。这种“同一个需求、不同配置、可量化的结果对比”才是实训里最有价值的部分。
第四步,观察 HDFS 上的输出结构,确认输出文件数与 Reduce 数一致,然后用hdfs dfs -cat抽几条结果验证正确性。有条件的话,再把结果用 Hive 或 Spark 读一次,感受不同框架对同一份 HDFS 数据的处理差异。
如果严格按照这个流程走一遍,你对 HDFS 和 MapReduce 的理解,会远超“会调 API”的水平。后面再学 Spark、Flink 时,很多概念(分区、洗牌、数据本地性、容错)都是相通的,学起来会轻松非常多。
5. 大数据作业的性能问题与排查实录
5.1 数据倾斜是最常遇到的“隐形杀手”
数据倾斜几乎是大数据场景最经典的头号问题,MapReduce 尤其容易踩中。表现上,就是某个 Reduce 任务运行时间超长,其他 Reduce 早早结束等着它,整个作业最终耗时被这一个任务拖死。
原因通常是业务 key 本身分布不均。比如热点商品、热门用户、某个地区的数据量特别大,HashPartitioner 又是按 key 哈希,哈希后热点 key 还是会落到同一个分区。于是,某台节点默默扛了几百倍于其他节点的数据量。
定位数据倾斜,我一般分两步:
- 先看 Counter。对比每个 Reduce 的
Reduce input records或者Reduce shuffle bytes,如果某个任务显著偏大,基本可以锁定倾斜。 - 再看日志和 Web UI 的任务列表。YARN 的 ResourceManager 界面里能看到每个任务读取的记录数和运行时间,一列出来,谁是“拖油瓶”一目了然。
解决思路也有常规套路。第一种是给热点 key 加随机前缀,把它拆成多个子 key,让它们分散到不同 Reduce,再做一次全局聚合。第二种是自定义 Partitioner,把业务上已经明确的少数热点 key 单独路由到指定分区,剩下的走默认哈希。第三种是两阶段聚合:先在 Map 端 Combiner 聚合一次,Reduce 端再聚一次,能把倾斜规模同时降下来。
这里有一个我踩过的坑:加随机前缀后,虽然均衡了负载,但因为同一个 key 被拆成了多个子 key,如果你在 Reduce 里做的是需要全局顺序的计算,比如排序、TopN,结果就会错乱。所以加随机前缀只适合求和、计数这类可再次聚合的运算。
5.2 调优方向:先看内存,再看并行度,最后看数据分布
很多人的调优顺序是反的,一上来就调并行度,结果没效果。我个人的调优顺序是:内存 → 并行度 → 数据分布。
- 内存方面,先看 Map 输出是否频繁溢写,
Spilled Records指标是不是异常高;再看 Reduce 端拉取数据时是否频繁落盘。如果内存吃紧,优先调大mapreduce.task.io.sort.mb,或者调整任务容器内存mapreduce.map.memory.mb和mapreduce.reduce.memory.mb。 - 并行度方面,看 Map 数量是不是过少。如果一个集群有几百个核,但 Map 任务只有几十个,资源没有充分利用。增加分片数或 Reduce 数能提升利用率,但别调到调度开销大于计算收益。
- 数据分布方面,确认没有明显的热点分区。冷热不均,再多的资源都会浪费在“等慢任务”上。
另外,Map 输出压缩也是个常被忽略的好手段。开启mapreduce.map.output.compress=true,选择 Snappy 或 LZO 编解码器,可以显著降低 Reduce 拉取的数据量,对网络紧张的场景帮助很大。代价是多了压缩解压的 CPU 开销,但大多数业务下收益都大于损耗。
5.3 一张表记住最实用的调优参数
业界流传的信息太多,我干脆把实践中最高频的参数整理成了一张速查表,方便你做作业前先过一遍:
| 参数 | 默认值 | 调优建议 |
|---|---|---|
mapreduce.task.io.sort.mb | 100 | Map 输出很大时调到 200~256,减少溢写次数 |
mapreduce.map.sort.spill.percent | 0.80 | 一般保持默认,不用刻意改 |
mapreduce.map.output.compress | false | 大作业开启,配合 Snappy 压缩减少网络 IO |
mapreduce.job.reduces | 1 | 按数据规模和集群核数调大,别让并行度浪费 |
mapreduce.reduce.shuffle.parallelcopies | 5 | Map 节点多时可提高到 10 左右,加快拉取 |
mapreduce.map.memory.mb | 1024 | 数据量大的 Map 任务适当加大防 OOM |
mapreduce.reduce.memory.mb | 1024 | Reduce 要聚合的数据多时考虑加大 |
mapreduce.map.maxattempts | 4 | 集群不稳时可调大,但要防止“无限重试”掩盖代码 Bug |
mapreduce.reduce.maxattempts | 4 | 同上 |
mapreduce.map.speculative | true | 长尾任务多时开启;代码有副作用则关掉 |
mapreduce.reduce.speculative | true | 长尾任务多时开启;效果不如 Map 侧明显 |
最后一个参数,推测执行,也值得单独说两句。它的逻辑是:同一个任务,如果检测到某个节点跑得明显比其他节点慢,框架会在另一台节点再启一个同任务的副本,谁先成功就用谁的结果。听起来很美好,但我在实际项目里见过它“帮倒忙”:任务本身有输出副作用,或者数据源不支持重复读取,跑出两遍结果就出问题。这种场景下,果断关掉mapreduce.map.speculative更稳妥。
最后再分享一个个人体会:MapReduce 这套框架比起后来的 Spark、Flink,确实重、确实慢,上手体验也不够“现代”。但我带过的所有新人里,凡是愿意花一个周末把这条完整流程亲手跑通、把 Counter 一个个点开看明白的人,后面学任何分布式框架都快得离谱。因为分片、洗牌、数据本地性、容错重试、任务推测这些核心概念,MapReduce 全都给你演了一遍。这篇能帮你把原理和实操串起来,我就觉得没白写。