ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

深入解析MapReduce:设计哲学、核心机制与工程实战

深入解析MapReduce:设计哲学、核心机制与工程实战 1. 从一次洗数据洗到怀疑人生说起如果你刚接触大数据或者已经在数据岗干过一阵子多半经历过这样一个场景老板丢过来一批日志说帮我把这几百个G的数据清洗一下提取出有效字段你心想几百个G也不算大开个Python脚本慢慢跑呗。结果脚本跑了两天还没结束内存报警磁盘IO飘红最后进程直接OOM死掉。这时候你才意识到单机处理真的扛不住大数据。MapReduce就是为了干掉这种窘境而生的。作为Hadoop体系里最经典的分布式计算模型它的核心思路只有四个字分而治之。把海量数据拆成小块扔到集群里多个节点并行处理再把局部结果汇总成最终答案。这套思想是Google在2004年公开的论文里提出来的后来Hadoop实现了开源版本成为大数据生态的起家本领。哪怕现在Spark、Flink这些实时计算框架风头正劲MapReduce依然是理解分布式计算的第一块基石——你如果能把MapReduce彻底搞懂再去学Spark会轻松非常多。这篇文章我想从一个实际使用者的角度把MapReduce掰开揉碎了讲一遍它背后的设计哲学、核心运行机制、怎么上手写代码、遇到数据倾斜怎么排查以及那些网上教程不爱写但实战中绕不开的坑。适合刚接触MapReduce的初学者也适合准备大数据面试或者需要在集群上跑离线作业的工程师参考。2. MapReduce的设计哲学为什么它能把大问题变简单2.1 单机瓶颈与移动计算胜过移动数据先说一个最根本的问题为什么几千个G的文件不能用一台好点的服务器处理说白了是两个瓶颈。一个是内存瓶颈数据量超过内存容量后系统开始疯狂做磁盘交换性能直接悬崖式下跌另一个是IO瓶颈单机硬盘的读写速度再快也有物理上限你把数据从磁盘搬到内存、再从内存写回磁盘整个过程全卡在IO上。分布式系统的思路是靠堆机器来解决问题但堆机器又带来了新问题数据分布在不同节点上怎么协同集群里有机器挂掉了怎么办每个节点的计算结果怎么汇总MapReduce最厉害的地方就是把这些分布式系统底层的脏活累活——任务调度、节点通信、故障恢复、数据合并——全部隐藏起来让使用者只需要关心两件事Map阶段怎么写Reduce阶段怎么写。这里面还有一个关键理念叫移动计算胜过移动数据。传统思路是你要分析某份数据就把数据拉到你的程序里来处理MapReduce反过来把处理逻辑也就是你的Map和Reduce代码分发到数据所在的那台机器上让程序走过去处理数据。数据不动代码动省下了海量数据跨网络传输的成本。这个思路在计算密集型任务里效果极其明显现在Spark、Flink也都延续了这个理念。2.2 一个MapReduce作业的完整生命周期理解MapReduce最好的方式不是死记概念而是跟着一份数据走完整个作业流程。假设你有一个100GB的文本文件分散存储在HDFS的多个Block上默认一个Block 128MB你要统计每个单词出现的次数。作业提交后框架会做这些事情第一步输入分片Input Split。Hadoop会根据HDFS上的Block信息把这份100GB的文件切成若干个分片每个分片默认对应一个Block的大小128MB也就对应一个Map任务的输入。这里要注意分片Split和块Block不是一回事块是物理存储单位分片是逻辑计算单位分片边界会考虑Block边界做划分保证每个Map任务处理的数据都在同一台机器上这就是数据本地性能大幅减少网络传输。第二步Map阶段。每个Map任务把分片里的数据一行一行读出来调用你写的map函数。你的map函数输入是一行文本输出是一系列key-value对比如输出(hello, 1)、(world, 1)这样的中间结果。第三步Shuffle与Sort阶段。这是MapReduce里最核心也最复杂的环节。Map输出的kv对不会立刻送到Reduce而是先按key做分区、排序、合并。所有key相同的记录会被分到同一个分区Partition分区数量与Reduce任务数量一致。这一步是在Map端先做一次本地排序等Reduce任务来拉取数据时收到的每个分区的数据就已经是排好序的了方便做归并。第四步Reduce阶段。每个Reduce任务负责一个分区把属于自己分区的数据聚合起来调用你写的reduce函数。reduce函数收到的是(key, 一个迭代器)迭代器里是所有该key对应的value比如(hello, [1,1,1,1,1])你在reduce里把这些value加起来输出(hello, 5)这样的最终结果。第五步输出。Reduce的输出写到HDFS上通常是part-r-00000这样带编号的文件。整个生命周期里作为一个使用者你只写了map和reduce两个函数剩下的分布式调度、并行计算、容错重试框架全包了。这正是MapReduce简单到极致的设计哲学——把复杂留给框架把简单留给用户。2.3 类比一下MapReduce就像一个流水线工厂把MapReduce比作一家工厂的流水线会非常直观。Map阶段就像流水线上的初加工车间原材料原始数据送到每个工位Map任务工人map函数把每份原材料粗加工成半成品中间kv对放在工位旁边的货架上本地磁盘。Reduce阶段像总装车间总装工Reduce任务会按订单key把不同初加工车间的半成品集中起来重新整理排序然后组装成最终成品输出结果。工厂里最麻烦的事情是什么不是加工本身而是物流调度——初加工车间生产出来的半成品怎么高效地送到对应总装车间而且路上不能丢乱了要重新排序。MapReduce框架干的就是这个物流调度工作做得还很优雅先在初加工车间原地排一次序Map端Sort再由专门的调度员Partitioner按key分好路线然后才发往总装车间。这一套流程下来总装车间收到的货基本就是有序的整个工厂的运转效率就高了。3. 核心机制拆解Map阶段、Shuffle、Reduce阶段3.1 输入分片与Map端实现要点Map阶段的输入处理在Hadoop 1.x和2.x及Hadoop 3.x中的默认实现是不同的。老的FileInputFormat按Block大小64MB或128MB切分分片Split新的FileInputFormat则更智能地处理了文件边界问题支持了更灵活的切分策略。不过核心原则一直没变分片大小决定了Map任务的粒度。分片太小比如只有几MB会导致Map任务数量过多一个集群几万个Map任务同时起调度开销、JVM启停的损耗都会拖垮效率分片太大比如超过一个Block又可能跨节点拉数据破坏数据本地性。所以生产环境里除非有明确的逻辑需求否则我基本不动默认的分片大小保持和HDFS Block一致是最稳妥的做法。在写map函数时有几点很关键map函数的输入key是行偏移量不是行号是这一行在文件中的字节偏移很多新手在这里会踩坑。输出kv对的类型要选对。Hadoop没有直接用Java的泛型而是用Text、LongWritable、IntWritable这些可序列化类型因为它们能被高效地通过网络传输。map函数是无状态的每一行输入都是独立的不要在map里维护全局状态比如计数器除非你清楚自己在做什么——因为Map任务可能被重启重启后状态就丢了。3.2 Shuffle分布式计算的物流神经系统Shuffle是MapReduce里最容易让初学者困惑、也最值得深入研究的一部分。它负责的是Map输出到Reduce输入之间的数据流转具体包括六个子环节1. 环形缓冲区每个Map任务都有一个环形内存缓冲区默认100MB由mapreduce.task.io.sort.mb控制。map函数输出的kv对会先写进这个缓冲区而不是直接落盘。2. 溢写Spill当缓冲区写满80%由mapreduce.map.sort.spill.percent控制后台线程会把缓冲区里的数据溢写到本地磁盘形成一个溢写文件。在溢写之前数据会先按key分区、再按key排序。这个先分区再排序的顺序很重要。3. 合并MergeMap任务执行过程中可能产生多个溢写文件等到map函数跑完这些溢写文件会被合并成一个大文件。合并过程中会进行归并排序还可以做Combiner文章后面会细说。4. 分区Partition上面提到过Shuffle的排序是按key做分区是按key的哈希值决定该去哪个Reduce任务。默认分区器HashPartitioner对key取哈希值后模上Reduce任务数公式是hash(key) % reduceNum保证相同key一定落在同一个Reduce里同时让不同key尽量均匀分布。5. 拉取Fetch每个Reduce任务会启动一对Fetcher线程从各个Map任务所在节点上拉取属于自己的那部分数据。这里有个隐患如果Map任务的输出很小Reduce拉取的开销占比就会很大这也是为什么Map端输出有压缩配置mapreduce.map.output.compress小集群上开启压缩收益非常明显。6. 归并排序Reduce端拿到多个Map任务的输出文件后会再归并排序一次保证输入给reduce函数的key是严格有序的。你可以在mapred-site.xml里配置这些参数也可以在每个MR作业的Configuration里单独配置。我的建议是先别急着调参把默认流程跑通再决定哪些参数值得动。3.3 Reduce阶段与输出Reduce函数的输入是一个(key, values)的迭代器这里有个官方文档写得少但极其重要的特性values迭代器是流式的不是一次性把所有value加载到内存。这意味着即使某个key有几百万个value只要你不手动把它们收集到List里内存压力依然可控。但现实是很多新手会在reduce里写这样的代码ListIntWritable allValues new ArrayList(); for (IntWritable val : values) { allValues.add(val); }这样写的话几百万个value全摊在内存里OOM迟早找上门。正确的姿势是在迭代器里边读边算比如求和就直接累加不要缓存。Reduce输出的默认方式是按key排序后的顺序写入part-r-00000文件。输出文件的数量和Reduce并行度保持一致这也是后续Hive查询、Spark读入时的重要假设。4. 手把手写一个MapReduce程序WordCount也是要优化过的4.1 环境准备与第一个作业在动手写代码之前你得先有一个能跑的Hadoop环境。如果是学习单机伪分布式模式就够了一个节点同时充当NameNode、DataNode、ResourceManager、NodeManager不用一上来就搭三台机器。伪分布式的配置过程比较多环境坑我第一次搭的时候折腾了一天在这里不重复了只提醒几个高频坑确保配置了JAVA_HOME环境变量并确认Hadoop版本对应的JDK版本。core-site.xml里的fs.defaultFS是hdfs://localhost:9000不是file:///写错会导致访问HDFS时报错。启动后多跑几遍jps命令看NameNode、DataNode、ResourceManager、NodeManager四个进程是否齐全缺一个都说明启动有问题。环境就绪后我用一个优化过的WordCount来演示核心代码。常规的教程版WordCount网上到处都是但那个版本在真实数据上跑起来性能和稳定性都很一般。我一般会加Combiner和自定义Partitioner实际效果差很多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; import java.util.StringTokenizer; public class WordCountMapper extends MapperLongWritable, Text, Text, IntWritable { private final IntWritable one new IntWritable(1); private final Text word new Text(); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line value.toString(); StringTokenizer tokenizer new StringTokenizer(line); while (tokenizer.hasMoreTokens()) { word.set(tokenizer.nextToken()); context.write(word, one); } } }这个map函数和教程版最大的区别是我做了复用word和one只在类里初始化一次避免每个单词都new一个对象的开销。别小看这个细节在几十亿行数据上对象创建和GC的开销能差出好几倍性能。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 ReducerText, IntWritable, Text, IntWritable { private IntWritable result new IntWritable(); Override protected void reduce(Text key, IterableIntWritable values, Context context) throws IOException, InterruptedException { int sum 0; for (IntWritable val : values) { sum val.get(); } result.set(sum); context.write(key, result); } }reduce函数没什么特别的注意一点我在循环里直接sum val.get()没有把values收集到List里这就是前面说的流式处理。Job配置部分是这个Demo的灵魂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 WordCountJob { public static void main(String[] args) throws Exception { Configuration conf new Configuration(); Job job Job.getInstance(conf, word count optimized); job.setJarByClass(WordCountJob.class); job.setMapperClass(WordCountMapper.class); job.setCombinerClass(WordCountReducer.class); job.setReducerClass(WordCountReducer.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); } }这里我特别加了一行job.setCombinerClass(WordCountReducer.class)。为什么要加Combiner在WordCount这个场景里combiner本质上就是在Map端先做一次本地合并同一个Map任务里比如(hello,1)出现了10次Combiner就是把这10个1先合并成10再输出这样Reduce需要拉取的数据量直接缩小到原来的几十分之一。Combiner用的是和Reducer同一个类因为WordCount的reduce满足合并操作可结合可交换的条件。你可以在mapred-site.xml里开启相关信息来看Combiner执行次数。4.2 让作业在集群上真正跑起来代码写完之后打包成Jar然后用hadoop命令提交作业hadoop fs -mkdir -p /input hadoop fs -put localfile.txt /input/ hadoop jar wordcount.jar WordCountJob /input /output注意/output路径不能已存在否则作业直接报错。运行过程中你会看到终端打印出Map和Reduce的进度百分比以及Counter统计信息。作业跑完后用下面的命令看结果hadoop fs -cat /output/part-r-00000我更习惯在这里看一眼Counter信息尤其是FILE_BYTES_WRITTEN和HDFS_BYTES_READ这几个指标它们能直接告诉你数据经历了多少磁盘IO。如果你发现FILE_BYTES_WRITTEN远大于输入数据量说明Shuffle阶段的数据量膨胀了这时候就可以考虑调压缩参数或者增加Combiner。4.3 用实操理解Design 模式如果把MapReduce当成一种编程范式很多经典数据处理模式是反复出现的。熟悉这些模式遇到实际问题时能快速套模板求和模式Summation上面演示的WordCount就是典型代表把一堆值加起来。对应场景统计访问量、PV/UV、销售额汇总等。去重模式DistinctMap阶段输出(value, null)Reduce阶段直接把key输出。对应场景清洗数据时找出不重复的用户ID等。排序模式SortMap阶段输出(value, null)框架按value排序后输出。注意如果要全局有序设置成单Reduce任务效率其次但数据全局有序。对应场景获取Top N。平均/统计模式StatisticsMap阶段输出(key, value)Reduce阶段可以是求和、平均数、方差、最大值、最小值等。常见的就是Hive里GROUP BY后聚合。反转模式Inverted IndexMap阶段输出(word, documentId)Reduce阶段输出(word, [doc1, doc2, ...])。对应场景搜索引擎的倒排索引。二分模式分治查找适合分阶段聚合的场景比如先按天聚合再按周聚合最后按年聚合。实际工作中遇到的问题80%都可以用这几种模式的组合来解决。初学者不要急着去看复杂框架先把你手头的问题抽象成这几种模式再决定Map和Reduce各自要做什么会顺手很多。5. 真实案例招聘数据清洗到底在洗什么5.1 场景分析网上有个热词是招聘数据清洗我就在实际项目里接到过类似需求几十个G的招聘平台原始数据格式是JSON或者CSV里面有大量脏数据——字段缺失、格式错误、重复记录、嵌套JSON里有些字段需要展开。清洗的流程如果用MapReduce来做天然就是分阶段的。先想清楚目标最终要生成一个干净的结构化表格每行是一条有效招聘信息包含公司名、职位、薪资范围、工作地点、发布时间、学历要求等字段。那Map和Reduce的职责是什么呢我的设计思路是Map阶段负责解析和过滤。每行输入是一条原始数据解析成若干字段做合法性校验比如必填字段不能为空、薪资必须是数字范围不合法或无效的数据直接丢弃。输出(position_id, 解析后的记录)这样一条kv对。Reduce阶段负责去重和最终输出。由于相同position_id可能会出现多次比如同一个岗位被多个渠道抓取Reduce端可以对同一个position_id做合并去重输出最终结果。这个设计的核心思想是尽量早地丢弃无效数据。在Map端就过滤掉的脏数据不会进入Shuffle节省了网络IO和磁盘IO。这个思路对任何大数据任务都是通用的过滤动作越早执行整体成本越低。5.2 关键代码思路如果你用Python写MapReduce现在很多学习场景会用因为不用编译打包可以直接用mrjob库pip install mrjob对应的清洗逻辑可以这么写import re from mrjob.job import MRJob class CleanApplyJob(MRJob): def mapper(self, _, line): # 解析CSV这里简单假设格式为 id,company,title,salary_min,salary_max,edu fields line.strip().split(,) if len(fields) ! 6: return # 格式不对直接丢弃 apply_id, company, title, salary_min, salary_max, edu fields if not apply_id or not company or not title: return # 关键字段缺失丢弃 try: salary_min int(salary_min) salary_max int(salary_max) except ValueError: return # 薪资不是数字丢弃 if salary_min 0 or salary_max salary_min: return # 逻辑错误丢弃 yield apply_id, (company, title, salary_min, salary_max, edu) def reducer(self, apply_id, values): # 去重同一个id只保留一条 seen set() for v in values: if v not in seen: seen.add(v) for v in seen: yield apply_id, v if __name__ __main__: CleanApplyJob.run()这段代码体现了一个重要的项目经验清洗逻辑必须能定位到为什么丢弃这一条。实际开发中我建议Map端除了输出有效记录还可以用context.getCounter()mrjob里是self.increment_counter来统计各种丢弃原因的数量。比如定义几个计数器INVALID_FORMAT、MISSING_FIELD、BAD_SALARY这样作业跑完你一眼就能看到各类脏数据的占比方便出清洗报告。6.3 一道经典面试题Combiner能不能乱用Combiner可以在Map端做局部聚合FinalReduce之前再做汇总这句话面试官很喜欢问。我实战中踩过一个大坑某次在Map端加了Combiner结果全局结果错了。原因是我那个reduce任务的聚合逻辑是计算唯一值数量SET的size而Combiner里如果执行同样的逻辑在Map端做的是局部去重Reduce端再做一次全局去重看上去好像没毛病。但实际上如果我让Combiner和Reducer完全复用同一个类而Reducer里有全局去重后还要基于全局再做一次排序或其他操作的逻辑最终结果就会偏离预期。Combiner能用的前提是reduce函数的运算满足可交换和可结合。比如求和Sum、最大值Max、最小值Min都满足但计算平均数不行——Map端的局部平均和Reduce端的全局平均不是一回事一旦用了Combiner结果必错。所以无论是面试还是实战记住一句话Combiner是锦上添花不是雪中送炭用之前先想清楚你的聚合函数是否满足交换结合律。6. 集群部署策略从伪分布到生产集群6.1 部署模式选择MapReduce作业要跑得灯稳底层集群部署策略很关键。不少学习者在单机伪分布模式下熟悉了用法一到生产环境就各种懵。这里我理一下三种常见模式单机模式Local Mode所有进程跑在同一台机器同一个JVM里用于代码调试、单元测试。好处是启动快坏处是根本没法模拟真正的分布式。伪分布式Pseudo-Distributed在一台机器上跑多个独立的Java进程模拟NameNode、DataNode等。学习和开发调试最推荐配置正确了基本不踩坑。完全分布式Fully-Distributed至少三台机器一个NameNode主节点多个DataNode从节点ResourceManager和NodeManager各自部署。生产环境的常态。规划集群时我习惯把NameNode和ResourceManager放在不同机器上避免Master节点负载过载DataNode和NodeManager放在同一批机器上这样Map任务才能利用数据本地性在本地处理HDFS的副本数据。6.2 资源规划与参数调优生产集群部署时我总结过一个经验公式基于Hadoop 3.x的默认容器模型假设你有10台服务器每台是64GB内存、16核CPU。那每台机器分配给YARN的资源大概可以这样算保留20%的内存给操作系统和HDFS的DataNode等进程剩余约50GB给YARN容器。每CPU核对应一个Map/Reduce任务容器但受限于内存每个容器默认最大内存可根据作业需求设置比如2~4GB。所以单台机器大约可以同时运行12~16个Map/Reduce容器受限于内存和CPU的较小值。这样10台机器集群整体可以支撑一百多个Map任务并发吞吐非常可观。注意给操作系统保留内存这个细节很多新人在小内存机器上把所有内存都分给YARN结果NameNode和DataNode反而因为内存不足频繁FullGC集群稳定性一塌糊涂。7. 常见问题与排查技巧实录7.1 数据倾斜跑得慢的隐形杀手数据倾斜是MapReduce实战中遇到最多的性能问题没有之一。它的典型表现是所有Map任务都跑完了但某个Reduce任务卡在99%几个小时不动其他Reduce任务早就完成了。为什么会倾斜因为Shuffle阶段是按key分区的如果某个key的数量远大于其他key比如电商日志里搜索请求这个key占了90%的数据那分配到这个Reduce任务上的数据就是别人的几十倍处理时间自然被无限拉长。排查方法第一看任务进度页面或YARN ResourceManager UI里每个Reduce的输入字节量。如果发现某个Reduce的输入量远超平均值的几倍甚至十几倍基本可以断定是数据倾斜。解决方案我按场景分几类加随机前缀Map阶段输出key时给热点key加一个随机后缀让它分散到多个Reduce任务Reduce端再做一次聚合。这个方法逻辑上等价于本地聚合先算一轮全局聚合再算一轮。自定义Partitioner如果热点key是已知的比如某个channel_id可以单独给这个key定义分区策略让它不去凑热闹——但这是带病运行不推荐长期依赖。多做一层预处理比如先用一次MR作业统计各key的数量分布就是频次统计找出热点key再针对性调整策略。7.2 小文件问题为什么Map任务多到爆炸小文件问题指的是HDFS上一堆几KB甚至几B的文件。MapReduce默认每个分片对应一个Map任务如果输入文件特别多特别小就会产生海量Map任务。举个例子一万个1KB的文件会生成一万个Map任务每个任务启动JVM的CPU/内存开销比数据本身还大效率极低。解决思路很直接总结三个动作合并输入文件用CombineFileInputFormat替代默认的FileInputFormat它会把多个小文件打包到一个分片里减少Map任务数。先合并再计算如果文件来源固定且持续产生小文件可以在ETL阶段用流批任务或用Hive的INSERT OVERWRITE合并小文件提前做合并。控制输出小文件Reduce输出时设置合适的Reduce数量如果业务上不需要分区尽量让Reduce写一个文件而不是几百个。7.3 Mapper阶段OOM、Reduce阶段OOM怎么区分OOM分两种处理方式完全不同。Map端的OOM一般发生在溢出写磁盘、排序或Combiner阶段优先看mapreduce.map.memory.mb是否设置得太小调大些其次确认Map函数是否缓存了过多数据比如一次把所有数据都load到内存。Reduce端的OOM多半是因为把大量value写进List再处理或者Reducer里用了全局HashMap等数据结构。解决办法就是拥抱流式迭代器能流式处理就别内存缓存。7.4 作业启动慢、卡在ACCEPTED状态你提交完作业发现日志一直显示Application has been accepted但就是不动。可能原因有两个一是集群资源不足所有NodeManager有大量Container在跑别的作业你的作业在队列里排队二是某个NodeManager上的可用内存比AMApplicationMaster要求的内存还小导致AM一直无法分配。排查方法在ResourceManager的Web UI上查看集群可用资源以及查看你的AM内存参数是否合理。尤其注意yarn.app.mapreduce.am.resource.mb这个参数它在Hadoop 3.x里默认值变大了如果集群内存不足这个参数反而成了启动瓶颈。8. 调优参数速查与学习路线建议8.1 常用调优参数清单我整理了一个工作中最常用的MapReduce调优参数表按核心程度排序参数名默认值作用我的经验推荐mapreduce.task.io.sort.mb100Map端环形缓冲区大小小文件多的场景下调到200~300MBmapreduce.map.java.opts默认200MB堆Map任务JVM堆大小比mapreduce.map.memory.mb小留出堆外空间否则会OOMmapreduce.map.memory.mb1024Map任务Container内存大内存机器上调到2048/3072mapreduce.reduce.memory.mb1024Reduce任务Container内存视Reduce函数复杂度调整mapreduce.map.output.compressfalseMap端输出是否压缩必开IO开销减少一半以上mapreduce.job.reduces1Reduce任务并行度计算公式min(集群CPU核数 * 1左, 数据总量 / 每Reduce期望处理量)mapreduce.task.io.sort.spill.percent0.8环形缓冲区内触发溢写的阈值不用动0.8是个好数字mapreduce.reduce.shuffle.parallelcopies5Reduce并行拉取Map输出的线程数集群机器多时可调大到108.2 学习路线建议如果你想系统掌握大数据核心技术完全可以从MapReduce出发形成一条递进的路线先学HDFS理解分布式存储的机制、Block、副本策略这样MapReduce的数据本地性才有意义。再学MapReduce把本文里的设计思想和代码实操过一遍重点吃透Shuffle。然后上手Hive你把MapReduce的调优搞懂了Hive其实只是SQL翻译成MapReduce作业的壳遇到Hive跑得慢你能很快定位到是底层哪个环节的瓶颈。最后学Spark现在很多公司做好离线数仓都在用Spark SQL但Spark的核心RDD操作依然处处能看到MapReduce的影子。理解了MapReduce的IO模型和延迟痛点你就明白为什么Spark要引入内存计算、DAG调度、按Stage划分任务——它就是为了弥补MapReduce在某些场景下每一步都要落盘的低效设计。这个路线走下来你对大数据离线计算的主流技术栈就有了一套知其所以然的体系而不是东一个工具西一个框架地学。9. 写在后面的一点体会跑了这么多年的MapReduce作业我有一个很深的感受真正难的不是写map和reduce函数而是理解数据在分布式系统里是怎么流动的以及当数据量级上去后那些平时看不出来的问题会怎么集中爆发。数据倾斜、小文件、JVM堆配置、分区策略……这些都是教科书上写得少、但生产环境里天天要面对的东西。如果你现在正准备用MapReduce做第一个真正的项目我建议你先拿一份真实的数据集在伪分布式集群上亲手把WordCount改造成一个像样的清洗工具。第一次跑通的感觉和看一百篇教程的感觉是完全不一样的。等你亲手排查过一次任务卡死、OOM、结果对不上你就真正把这门技术焊死在脑子里了。
返回列表