1. 从“一脸懵”到“真看懂”:MapReduce到底是什么
我到现在还记得第一次翻开Hadoop源码、看到MapReduce两个单词时的感受——名字高端、文档晦涩、示例代码看完一遍脑子还是空的。后来在真实项目里把一个几千万行的日志清洗任务跑通、调优、稳定上线之后,才慢慢发现,MapReduce这个老家伙能在大数据领域屹立这么多年,是真的有它的一套逻辑。
如果你现在也处于“听说过MapReduce、背过面试题、但没真正理解它到底怎么工作”的阶段,这篇文章就是写给你的。我会用偏实战的视角,把它拆开揉碎讲清楚:核心设计思想是什么、完整工作流程长什么样、代码怎么写、和HDFS怎么配合、踩过哪些坑。不求让你看完就成为专家,但至少下次别人聊起MapReduce,你能接住话,并且能动手跑通一个自己的作业。
先说一个很多人容易混淆的点:MapReduce不等于Hadoop,它是Hadoop体系里负责“计算”的部分,对应Google那篇著名的MapReduce论文,而HDFS负责的是“存储”。这两个东西是配套的,经常一起出现,所以后面我也会花篇幅讲它们的集成原理。理解了这一点,你再看任何MapReduce教程,框架感会清晰很多。
2. 分而治之:MapReduce的核心设计哲学
2.1 为什么大数据的计算非得分“Map”和“Reduce”两步
要理解MapReduce,先忘掉代码,想一个生活化的场景。
假设你的任务是统计一座城市所有图书馆里“技术类”书籍的总数,图书馆有100家,分布在不同城区。最笨的办法是一个人挨家挨户跑,把100家都数一遍,最后加总——这就是单机处理,数据量一大就崩了。换个思路,你找来100个人,每人负责一个图书馆,先各自数自己馆里的技术书数量,再派一个人把100个数字汇总起来——这就是Map。每个人数自己负责的那一堆书的过程,就是Map阶段;汇总的环节,就是Reduce阶段。
MapReduce的名字就是这么来的:Map负责“映射、处理、归类”,Reduce负责“汇总、聚合、计算”。它的设计哲学就四个字:分而治之。一个大任务被拆成成千上万个独立的小任务,分布式地跑在很多台机器上,再把结果合并起来。
这个设计放在2004年前后确实是极具颠覆性的:在那之前,处理海量数据要么靠超级计算机(贵得要命),要么靠写复杂的分布式并发代码(要命)。MapReduce把并行计算的复杂度封装了起来,让普通工程师只需要写两个函数,就能利用一个机器集群的算力。
2.2 为什么MapReduce合适解决大规模离线计算问题
判断一个技术方案好不好,要看它解决了什么问题。MapReduce能火,核心是它精准击中了那个年代最普遍的一个痛点:海量数据离线批量处理。
注意“离线”和“批量”这两个词。MapReduce的设计目标不是毫秒级响应,不是处理一条数据流,而是在一段时间内(几分钟到几小时)把存储在磁盘上的大批量数据完整地跑完一遍,产出结果。典型的场景包括:
- 全量日志的清洗和格式化(比如把几TB的Nginx日志解析成结构化表)
- 全量数据的离线统计(比如“昨天全网每个商品的PV/UV”)
- 大规模索引构建(早期的搜索引擎就是这么干的)
这类任务有几个共同特征:数据量极大、计算逻辑不复杂(往往是过滤、映射、聚合)、对实时性没有要求。MapReduce的串行化、容错、自动并行机制,让这类任务的开发门槛降到了极低——你不用关心哪台机器在跑哪个任务,也不用操心某台机器挂了怎么办,框架全包了。
也正因如此,后来Spark出来的时候,大家一边欢呼它把中间结果放到内存、迭代计算快了上百倍,一边依然承认:在海量数据的一次性扫描处理上,MapReduce的稳定性和工程成熟度依然是标杆。很多公司线上跑了几年的定时离线任务,底层还是MapReduce,不是因为没人迁移,而是因为“能稳定跑完不犯错”本身就是巨大的价值。
3. 工作流程拆解:一个作业从提交到落盘的完整旅程
3.1 Job与Task:理解这两层概念是看懂任务调度的关键
在细看流程之前,先理顺两个基础概念:Job和Task。
一个MapReduce程序提交到集群后,会被框架看成是一个Job(作业)。这个Job是逻辑层面的完整任务,比如“统计2024年1月所有订单的金额总和”。而Job会被拆分成无数个Task(任务):
- MapTask:负责处理输入数据中一个分片的所有数据,实现Map函数
- ReduceTask:负责对一组相同Key的数据执行Reduce函数,实现汇总
默认情况下,一个MapTask处理一个输入分片(split),分片的数量往往由数据量和HDFS块大小决定;ReduceTask的数量由代码里job.setNumReduceTasks()设置,默认是1个。
举个直观的例子:如果一份数据在HDFS上有10个块(默认128MB一个块),那么MapTask通常就是10个;如果代码里设置了3个ReduceTask,那么每个ReduceTask会接收到所有MapTask输出中属于它负责的那部分Key的数据。
理解这个区分很重要,因为MapReduce框架的所有调度、并行、容错其实都是针对Task这个粒度的。Task个数越多,并行度越高,但也意味着调度开销越大;Task太少,那么大的集群就闲着,那就是在浪费资源。
3.2 从InputFormat到OutputFormat:手绘一张MapReduce全流程图
MapReduce流程网上有很多图,但我建议你自己画一遍,把每个环节串起来。一个标准流程包含以下阶段:
HDFS上的输入数据 ↓ InputFormat:把输入数据切分成多个split,并为每个split生成RecordReader ↓ Map阶段:对每个KeyValue对调用一次map函数,输出临时KeyValue对 ↓ Shuffle过程:分区(Partition)→ 排序(Sort)→ 溢写(Spill)→ 合并(Merge) ↓ Reduce阶段:对每个Key调用一次reduce函数,处理该Key对应的所有Value列表 ↓ OutputFormat:把结果写回HDFS每一环都有具体的类在干活:
- InputFormat:负责验证输入格式、把文件切分成split、提供RecordReader。常见的有
TextInputFormat(默认,一行一行读)、KeyValueTextInputFormat、SequenceFileInputFormat。 - RecordReader:把split中的原始数据解析成
<key, value>对。默认情况下,key是行首字节偏移量(LongWritable),value是这一行的内容(Text)。 - Mapper:用户实现的核心逻辑,输入是
<key, value>,输出是新的<key, value>。 - Partitioner:决定Map输出的每条数据去哪个ReduceTask。默认的是
HashPartitioner,即对key做哈希取模,均匀分布到不同的ReduceTask上。 - Combiner:可选的“迷你Reducer”,在Map端先做一次局部聚合,减少网络传输的数据量。
- Shuffle与Sort:框架自动完成的“洗牌”过程,保证每个ReduceTask拿到的数据都是按Key排好序的。
- Reducer:用户实现的核心逻辑,输入是
<key, Iterable<value>>,输出最终结果。 - OutputFormat:决定结果输出到哪、怎么写。默认是
TextOutputFormat,每个键值对输出一行,用制表符分隔。
3.3 Shuffle阶段凭什么被称为MapReduce的“灵魂”
我必须专门拉一节出来讲Shuffle,因为MapReduce面试的压轴题基本都在这里,而且后续调优几乎都是在调Shuffle。
Shuffle字面意思就是“洗牌”,它描述的是从Map输出到Reduce输入的这段过程。它由两部分组成:
Map端的Shuffle:
- Map函数跑完后,输出结果不是直接写磁盘,而是先进入一个环形内存缓冲区(默认大小100MB,通过
mapreduce.task.io.sort.mb配置)。 - 当缓冲区达到阈值比例(默认80%)时,后台线程开始“溢写”,把缓冲区数据写到磁盘的临时文件。
- 在溢写之前,会首先做一次分区(按Partitioner逻辑划分数据),然后对每个分区的数据按Key做排序。
- 如果配置了Combiner,会在溢写时对局部数据进行合并(注意:Combiner可能被调用多次,所以它的逻辑必须满足幂等性)。
- 每次溢写产生一个小文件,当MapTask结束时,会对这些小文件做归并排序,合并成一个大文件,同时告诉AppMaster这个文件的位置。
Reduce端的Shuffle:
- 每个ReduceTask启动后,先去AppMaster拉取元数据,知道哪些MapTask完成了,然后启动专门的线程到各个MapTask节点上拷贝属于自己分区的数据。
- 拉过来的数据先放入内存缓冲区,如果不够就落盘。
- 数据全部拉取完毕后,对多份数据进行归并排序,最终形成“按Key分组、组内有序”的输入序列。
- 这时候才真正调用reduce函数,对每个Key对应的Value列表做计算。
这里有个最值得注意的设计点:排序不是我们手动调了sort,而是框架自带的机制。Map输出会排序,Reduce输入也会排序,这意味着reduce函数天然接收到有序的数据流,这让很多聚合类的计算变得极为简单——顺序读一遍就能完成去重、排序、TopN等操作。
理解了Shuffle,你就能明白几个经典调优方向的底层逻辑:为什么Combiner能大幅提升性能?因为它把网络传输量减了;为什么要压缩Map输出?因为Shuffle过程中最贵的是I/O;为什么有时候ReduceTask特别慢?大概率是某个Key的数据量特别大(数据倾斜),导致某个Reduce分区的数据远多于其他分区。
4. 编程实例:WordCount从零到跑通
4.1 环境准备和代码骨架:Maven工程建起来
讲完原理,咱们来写代码。不管是学校的HDFS和MapReduce综合实训,还是工作中第一次写MR任务,“WordCount”(词频统计)都是绕不开的Hello World。我强烈建议你亲手敲一遍,别直接复制,敲的过程中你会对每个类的继承关系印象深刻。
假设你已经搭好了Hadoop集群(单机伪分布式也可以),并且Maven能用了。第一步是建一个标准Maven工程,引入依赖。
<dependency> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-client</artifactId> <version>3.3.6</version> </dependency>由于Hadoop相关的包特别多,直接用hadoop-client一个依赖就够了,它会把你开发MapReduce程序所需的全部API都带进来。版本号建议和你的集群版本保持一致,避免RPC协议不兼容。
工程的Java包名我一般用com.example.mr,然后写三个类:WordCountMapper、WordCountReducer、WordCountDriver。
4.2 Mapper和Reducer的编写细节
先看Mapper。它的任务就是把每一行文本拆成单词,输出“单词,1”这样的键值对。
import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; import java.io.IOException; public class WordCountMapper extends Mapper<LongWritable, Text, Text, IntWritable> { private final Text word = new Text(); private final IntWritable one = new IntWritable(1); @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line = value.toString(); String[] words = line.split("\\s+"); for (String w : words) { if (w.isEmpty()) { continue; } word.set(w); context.write(word, one); } } }注意,我特意把word和one设计成类的成员变量,而不是在循环里new。这是MR编程里一个很经典的性能优化——MapTask会为每条输入记录调用一次map()方法,如果每次调用都新建对象,几百万行数据就意味着几百万次对象创建和GC,非常浪费。复用对象可以显著降低内存压力。
再看Reducer。它的任务是把同一个单词的所有“1”加起来。
import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; import java.io.IOException; public class WordCountReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private final IntWritable result = new IntWritable(); @Override protected 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); } }这里需要注意:values是一个可复用的迭代器对象,也就是说在for循环里,val指向的对象每次可能都是同一个,只是内容被框架刷新了。如果你想把values保存下来留到循环之后再处理,一定要重新new一个对象并把值拷进去,否则会得到一堆重复的“最后一值”。
4.3 Driver类的标准写法
Driver类是这个作业的入口,负责组装整个作业的各种配置。它通常长这样:
import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; public class WordCountDriver { public static void main(String[] args) throws Exception { if (args.length < 2) { System.err.println("Usage: WordCountDriver <input path> <output path>"); System.exit(-1); } Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "Word Count"); job.setJarByClass(WordCountDriver.class); job.setMapperClass(WordCountMapper.class); job.setReducerClass(WordCountReducer.class); job.setMapOutputKeyClass(Text.class); job.setMapOutputValueClass(IntWritable.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); boolean success = job.waitForCompletion(true); System.exit(success ? 0 : 1); } }几个关键点给新手强调一下:
setJarByClass是必须的。它告诉Hadoop去哪找实现代码的Jar包。如果不设置,你在本机跑IDE里的作业可能会报ClassNotFoundException。setMapOutputKeyClass和setMapOutputValueClass设置的是Map阶段输出的类型;setOutputKeyClass和setOutputValueClass设置的是Reduce阶段输出的最终类型。如果你的Reducer的输出类型和Mapper不同(比如你在Reduce里做了类型转换),这两个地方的设置必须分开指定,否则序列化会出错。- 输出路径绝对不能已经存在。这是Hadoop的一个安全机制,如果输出目录已存在,它会直接报错,防止你把上一次的结果覆盖掉。所以每次跑之前记得删掉旧的输出目录,或在代码里加一个判断来删除。
4.4 打包提交:在集群上跑起来
在集群上执行有两种常见姿势:一种是把代码上传到服务器,用hadoop jar命令跑;另一种是在IDE里直接跑,但需要你本机能连通集群NameNode和NodeManager。
先说第一种。在Maven里打包后,用下面命令提交:
hadoop jar wordcount-1.0.jar com.example.mr.WordCountDriver /input/data.txt /output/wordcount这里的/input/data.txt和/output/wordcount都是HDFS路径。提交之前别忘了先往HDFS放数据:
hdfs dfs -mkdir -p /input hdfs dfs -put data.txt /input/跑起来之后,你会看到一堆日志,重点关心两个数字:Map input records和Reduce output records。如果Map阶段输出了正确的记录数且没有ERROR日志,基本就成功了。然后验证结果:
hdfs dfs -cat /output/wordcount/part-r-00000 | head -20会看到类似hello 5这样的输出,每一行一个单词加一个统计数字。这里有个细节值得注意:输出的文件叫part-r-00000,其中r代表它是Reduce的输出;如果你设置了setNumReduceTasks(0),Map的输出会直接写盘,生成的文件名是part-m-00000。
5. HDFS与MapReduce综合实训:从原理到一整套数据清洗
5.1 为什么HDFS是MapReduce的“最佳拍档”
很多人学习的时候会把HDFS和MapReduce分开看,但真实项目里它们永远一起出现。原因其实很本质:MapReduce处理的数据存在哪?——HDFS;结果写哪?——HDFS。这套搭配不是偶然的,而是设计上的必然。
HDFS提供的是分布式的存储:数据被切成块,每个块默认128MB,在集群的多台机器上做冗余备份。MapReduce提供的是分布式的计算:它可以把一个任务拆成很多个小任务,并行地在多台机器上跑。但是计算要想高效,有一个前提条件,那就是计算最好就在数据所在的机器上跑,而不是把数据从A机器搬到B机器再算。否则网络传输就会成为整个系统的瓶颈。
HDFS和MapReduce的“天作之合”就体现在这里:数据被分成块、冗余存储在多台机器上,MapReduce在调度MapTask时,会优先把任务调度到“保存有对应数据块副本”的那台机器上。这个机制在Hadoop里叫“数据本地性”(Data Locality)。它会尝试让计算任务和它的输入数据在同一个节点上,这样读取数据就是本地磁盘I/O,而不是网络I/O,速度快一个数量级。
如果数据块的副本恰好都在繁忙的机器上,任务调度会退而求其次,选择“机架感知”的级别——尽量调度到同一机架的机器上。如果你的机架网络带宽是千兆甚至更高,这个级别的本地性损失是可以接受的。数据本地性是衡量MapReduce集群性能的核心指标之一,可以在ResourceManager的Web UI上看到对应百分比。
5.2 案例:清洗一份Nginx日志
5.2 案例:清洗一份Nginx日志
理解了HDFS和MapReduce的配合方式,咱们来做一道综合实训题:清洗一份Nginx访问日志。
日志样本长这样(字段间用空格分隔):
192.168.1.10 - - [10/Jan/2025:08:30:25 +0800] "GET /api/v1/users HTTP/1.1" 200 532 "https://example.com" "Mozilla/5.0"需求是:从这些原始日志中提取出IP、访问时间、请求方法、请求路径、状态码、响应体大小,过滤掉状态码非200的记录,输出为制表符分隔的格式化文件,供后续数据分析使用。
5.3 解析逻辑实现
Mapper的核心逻辑是把一行原始日志拆解成结构化的字段。由于日志格式相对规整,直接用正则表达式或字符串切割都能做。
public class LogCleanMapper extends Mapper<LongWritable, Text, Text, Text> { private static final Pattern LOG_PATTERN = Pattern.compile( "^(\\S+) (\\S+) (\\S+) \\[([^]]+)] \"(\\S+) (\\S+) (\\S+)\" (\\d{3}) (\\S+)" ); private final Text outKey = new Text(); private final Text outValue = new Text(); @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line = value.toString(); Matcher matcher = LOG_PATTERN.matcher(line); if (!matcher.matches()) { context.getCounter("Log Cleaner", "Malformed Lines").increment(1); return; } String ip = matcher.group(1); String time = matcher.group(4); String method = matcher.group(5); String path = matcher.group(6); String protocol = matcher.group(7); String status = matcher.group(8); String bodyBytes = matcher.group(9); if (!"200".equals(status)) { context.getCounter("Log Cleaner", "Non-200 Lines").increment(1); return; } outKey.set(ip); String cleaned = String.join("\t", time, method, path, protocol, status, bodyBytes); outValue.set(cleaned); context.write(outKey, outValue); } }这里我引入了一个很实用的技巧:使用context.getCounter()来统计数据质量。通过Counter,你可以在作业跑完之后看到“总共有多少行格式不对、多少行被过滤掉”,这对清洗类任务极其重要。你不需要写额外的逻辑去统计这些信息,MapReduce框架会帮你汇总并在日志中打印出来。跑完作业后在控制台或HDFS的作业历史里就能看到Counter的值。
Reducer在这类清洗任务中一般不需要做太复杂的事。大部分场景下,我们只是把Map阶段解析出的数据直接原样输出。这种情况下可以把分区策略改成自定义策略,或者直接不设置Reducer,只用MapOnly模式:
job.setNumReduceTasks(0);这样Map的输出直接落盘,省掉了Shuffle的排序和网络传输开销,性能非常可观。清洗类任务本质上就是“解析+过滤”,不需要聚合,所以完全可以只跑Map。
5.4 如何避免生产环境下的小文件问题
清洗任务跑完,你很可能遇到一个麻烦:输出了一堆小文件。每个MapTask生成一个文件,如果输入文件有1000个split,那输出目录就有1000个小文件,每个只有几KB。 这种结果很不友好:后续如果再用这堆文件作为输入去跑MapReduce,又会启动1000个MapTask,无形中放大了任务调度的开销;若是Hive表的数据目录,小文件过多会影响查询性能。
解决办法有几种:
- 如果结果必须按某个Key分区聚合,就设置合理的ReduceTask数量,让多个Map输出合并到同一个Reduce输出文件中。但这会引入Shuffle开销,不适合纯清洗场景。
- 用
MultipleOutputs,在Map端把一个作业的结果写到多个目录,控制每个文件的写入数量。 - 最后用一条独立的合并命令,例如
hdfs dfs -getmerge /output/* /local_result.txt,把远端的小文件合并后拉到本地处理。 - 使用Hadoop的
CombineFileInputFormat,在输入阶段就把多个小文件合并成一个split,这样MapTask的数量不再等于文件数量,而是等于合并后的split数量。这个方法对处理“海量小文件作为输入”的场景非常关键。
我在实际项目中一般这样判断:如果清洗任务产出的是“中间结果”,后面还要继续做统计分析,那就用CombineFileInputFormat配合一个合适的splitSize,把输出控制在“个位数到几十个文件”;如果产出的是“最终结果”,并且需要同步到其他系统,就直接getmerge合并成一个文件。
6. 故障自救:MapReduce实战中那些经典大坑
6.1 内存配置不当导致的Container反复被杀
我第一次跑一个稍微大点的作业,就遇到了MapTask和ReduceTask反复失败、应用直接被Kill的情况。日志里写着:
Container killed on request. Exit code is 143一开始完全不知道怎么回事,后来才明白,这是Container请求的内存超过了YARN给它的上限,被NodeManager强制杀掉了。每个MapTask和ReduceTask跑在YARN的Container里,这个Container能用的物理内存受两个参数控制:
mapreduce.map.memory.mb:MapTask申请的Container内存上限,默认1024MBmapreduce.reduce.memory.mb:ReduceTask申请的Container内存上限,默认1024MB
如果你的mapred-site.xml里没有专门调大这两个值,而你的MapTask处理的数据量特别大(比如一行文本有几MB),Mapper里又做了很多字符串处理,就很容易撑爆默认的1GB。
解决的办法不是简单粗暴地把值调到8GB,而是要根据数据特点评估。一个更稳健的做法是同时调整Java堆大小和Container总内存,因为Container总内存包含Java堆以外的开销(如元空间、堆外缓冲等):
<property> <name>mapreduce.map.memory.mb</name> <value>2048</value> </property> <property> <name>mapreduce.map.java.opts</name> <value>-Xmx1536m</value> </property>一个合理的经验法则是:mapreduce.map.memory.mb比mapreduce.map.java.opts的Xmx大20%-30%,给堆外留足余量。否则即使容器内存加到了2GB,JVM堆上限只设1GB,GC一频繁,整体性能反而更差。
6.2 遇到数据倾斜:最头疼的“尾延迟”问题
数据倾斜是每个写MapReduce的人迟早会撞上的问题,也是最考验排查经验的问题。
表现很典型:整个作业有一大堆MapTask秒完,剩下的ReduceTask有的也秒完,但有一个或几个ReduceTask跑了几个小时都不动弹。原因通常是一个或少数几个Key的数据量远超其他Key——比如电商日志里“某个头部商家的访问量”占了总数的一半,或者用户关联数据里“NULL”值特别多。
排查方法第一步是看Counter或Spark Star类似的统计,确认ReduceTask接收的records数是否严重不均衡。确认之后有几种破解思路:
- 给倾斜的Key加随机前缀,把一个大Key拆成多个小Key去跑,最后再去掉前缀聚合。这个方式能治本,但需要业务上接受“多跑一轮”。
- 用Combiner做局部聚合,把“同一个MapTask内的重复Key先合并”,减少Reduce端的压力。这个对“计数求和类”问题非常有效,但前提是聚合逻辑满足交换律和结合律。
- 设置
mapreduce.job.reduces调整ReduceTask数量,让数据分布到更多Reduce上。这个方案治标不治本,如果倾斜很严重,调整几个Reduce也没用。 - 多阶段作业:先跑一个轻量Job统计Key分布,再根据分布定制Partitioner,让每个ReduceTask分到的数据量接近均衡。
我见过不少工程师在倾向上死磕,一个方案不行又换另一个,折腾半天。这里给句诚实的建议:如果业务允许,先确认这个Key是不是非法数据(比如空值),允许的话直接在Map阶段过滤掉,省下的时间是实打实的。
6.3 Map输出压缩:不调这个参数你迟早后悔
另一个实操里性价比极高的优化是开启Map端输出压缩。Shuffle阶段最耗时的环节就是把Map输出从MapTask所在节点拷贝到ReduceTask所在节点。如果中间结果数据量很大,网络传输的时间可能比计算本身还长。
开启压缩只需要两行配置:
<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在压缩率和解压速度之间取得了很好的平衡,特别适合Shuffle这种“边写边传边读”的场景。实测中,开启后Shuffle阶段的网络传输量能下降60%-80%,而代价是Map和Reduce两端各增加了一点和压缩解压相关的CPU开销。对大多数IO密集型的作业来说,收益远大于损失。
6.4 日志里那些“诚实的谎言”和排查招法
最后分享一个写MR作业时几乎必见的诡异情况:作业明明失败了,但日志里只有一个FailedTask,却看不到任何Exception。这时候别慌,用几招按顺序排查。
第一步,去YARN的ResourceManager Web UI找对应Application ID,点击进入日志聚合页面,查看那个失败Task的完整堆栈,尤其是System.err。第二步,查看HDFS上的用户日志目录(在mapreduce.jobhistory.intermediate-done-dir配置的路径下),那里的日志比控制台更全。第三步,重点看java.lang.ClassNotFoundException——这种情况九成是你打的Jar包没包含第三方依赖,试试Fat Jar(带依赖的包)提交。
7. 我认为的MapReduce正确“打开方式”
写到这里,MapReduce的核心概念、工作原理、编程方式、和HDFS的配合、实战调优和常见坑都过了一遍。最后聊几句我个人的体会。
很多新手学MapReduce的时候会被“老”“慢”“过时”这些印象影响,觉得这玩意没什么可学的。但我的观点恰恰相反。MapReduce是整个大数据生态里最适合入门“分布式计算”的教材。它把“分布式”这个概念拆得极其清晰:数据是怎么分的,任务是怎么调的,机器挂了怎么办,数据倾斜怎么处理。这些东西你学透了,再去看Spark、Flink的源码和调优文档,会发现很多概念都能一一对上——只不过人家换了个更快的引擎,但底层那些关于分区、排序、数据本地性、容错的智慧,本质上是一脉相承的。
所以,如果你正在学这个方向,建议别急着跳去Spark,先用MapReduce亲自写完十个作业,踩过那几个经典的坑。等你亲手调过内存参数、处理过一次数据倾斜、看懂过一条Shuffle日志之后,你就真正跨过分布式计算的“门槛”了。到那时,你再看任何新的计算框架,都会有一种“万变不离其宗”的踏实感。