ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

深入理解MapReduce:原理、流程、实战与调优

深入理解MapReduce:原理、流程、实战与调优 如果你正在学大数据MapReduce这四个字一定会反复出现在你眼前。无论是课程作业、毕业设计还是面试题里让你手写WordCount都绕不开这个Hadoop生态里最经典的分布式计算框架。坦白讲很多人初学时容易陷入一个误区把MapReduce当作一堆API背下来能跑通一个统计程序就觉得“会了”结果一旦换一个计算场景或者集群出了问题立刻抓瞎。我从最早照着教程在虚拟机里跑Hadoop伪分布式到后来在真实集群上处理上百GB的日志数据中间踩过不少坑。这篇学习笔记不是简单罗列概念我会把MapReduce的核心设计思路、一次完整作业的运行流程、常见错误排查思路以及我个人的调优经验尽量用大白话讲清楚。这篇内容适合正准备接触大数据的初学者也适合要应付面试、需要系统梳理MapReduce知识体系的朋友。1. MapReduce的核心设计思路与整体架构1.1 为什么单机算不动非要用分布式框架先想一个问题一份日志文件有100GB你在一台8核16GB内存的机器上用单线程程序逐行统计关键词最快也得跑十几分钟甚至几个小时。如果数据量到了TB级别单机基本就是死路。这时候最朴素的想法就是多找几台机器一起算算完再把结果汇总。想法很简单落地却很难。你得考虑数据怎么切分、节点之间怎么通信、某个节点挂了怎么处理、负载不均衡怎么办。MapReduce这个框架存在的意义就是把“分布式计算的复杂度”从程序员手里接过去——你只需要关心“业务逻辑本身怎么写”至于数据切分、任务调度、节点通信、故障恢复全部交给框架处理。所以说MapReduce本质上是一个“编程模型 运行时环境”的组合。编程模型负责定义你该怎么写代码Map和Reduce两个阶段运行时环境负责在实际集群中把你的代码跑起来。1.2 Map与Reduce分治思想的两步走整个MapReduce设计核心就两个阶段Map阶段和Reduce阶段对应的是“分而治之”的思路。Map阶段可以理解为“分”把一个大任务拆成很多个小任务每个小任务处理一份数据切片。比如100GB的日志分成100份每份1GB同时交给集群里不同的节点去并行处理。每个Map任务只负责自己那一份输出一个中间结果。Reduce阶段可以理解为“合”把Map阶段产生的中间结果按照某种规则归并到一起做最后的汇总计算。比如把各节点统计出的单词频次再合并成一份全量结果。打个比方你要统计一个学校所有学生喜欢的社团类型Map阶段就是每个班自己统计本班的数据Reduce阶段就是收集所有班级的数据后合并出全校的统计结果。这个类比虽然简单但对应到数据流动、任务调度的各个环节你就能理解MapReduce的整体脉络了。1.3 移动计算而非移动数据有一句经典原则是“数据本地性Data Locality”MapReduce会尽量把计算程序调度到数据所在的节点上运行而不是把数据拉回到计算节点。为什么要这么做因为网络传输相比磁盘IO和内存计算是分布式环境下最昂贵的资源消耗。一份数据可能几百GB程序包可能只有几百KB把程序发过去显然比把数据拖回来划算得多。这个设计直接影响HDFS和MapReduce的关系。HDFS负责把数据切块分布式存储MapReduce在调度Map任务时会优先选择“数据块所在节点”来运行任务——这就是HDFS与MapReduce能够“亲密合作”的根本原因也是你搭Hadoop集群时DataNode和NodeManager通常部署在同一批机器上的原因。1.4 容错机制集群环境下的生存之道集群中节点宕机是常态而不是异常所以容错是框架的必修课。MapReduce提供了两个层面的容错任务失败重试。任何一个Map或Reduce任务如果执行失败比如进程崩溃、节点宕机ApplicationMaster会把任务重新调度到另一台健康的节点上重新执行。默认重试次数是4次如果4次都失败整个作业才会标记为失败。推测执行。集群里偶尔会出现一台性能较低的节点比如磁盘即将满、CPU资源被占导致某个任务迟迟无法完成拖慢整个作业。推测执行机制会检测到这种情况在同一份数据上另起一个“备份任务”哪个先完成就采用哪个结果另一个直接杀掉。这个机制在实际生产中经常能救你一次但也要注意如果集群负载已经很高推测执行反而会加剧资源竞争这时候可以手动关闭它。2. MapReduce核心机制与原理拆解2.1 从输入分片到RecordReader数据是怎么被切开的一个MapReduce作业从输入数据开始。输入数据的来源可以是HDFS上的文件也可以是其他存储系统。框架会先把输入数据逻辑上切成若干个“分片InputSplit”每个分片最终会交给一个Map任务处理。分片的大小有个默认规则尽量接近HDFS数据块的大小默认是128MB。为什么定到这个量级因为如果分片小于数据块会产生大量Map任务任务调度开销会超过计算本身如果分片远大于数据块又可能导致某个Map任务要跨节点远程读数据违背“数据本地性”原则。有了分片还需要一个“RecordReader”来按行或按其他格式读取分片内容把它们转换成“键值对”key-value交给Map函数。对文本文件来说默认的TextInputFormat会把每一行的行首偏移量作为key这一行的内容作为value。你可以通过自定义InputFormat和RecordReader支持JSON、CSV、多行合并等特殊输入场景。2.2 Mapper端处理业务逻辑的第一次加工Mapper端的业务逻辑由你自己编写继承Mapper类、重写map方法即可。map方法的输入是“键值对”输出也是“键值对”。这个阶段基本不做复杂的聚合操作只做“数据清洗、提取、转换”这类相对独立的工作。举个例子做用户访问日志分析时你可以把一行日志切割出时间、用户ID、访问URL输出key为“用户ID”value为“1”表示该用户产生了一次访问。每个Map任务处理完自己的分片后会先把结果写到本地磁盘——注意此时输出是写本地磁盘而不是HDFS因为中间结果不需要多副本冗余如果写HDFS反而会造成巨大的网络和存储开销数据还得走一遍复制流程没有必要。这里有一个很容易被忽略的性能隐患如果Map输出的中间结果非常大磁盘IO会成为瓶颈。所以后面会讲到Combiner和压缩就是在这一环节做优化的。2.3 Shuffle与SortMapReduce的灵魂枢纽MapReduce里最容易让人糊涂的就是Shuffle阶段。它发生在Map输出之后、Reduce输入之前负责把Map产生的“无序、分散”的中间数据整理成“分区有序、按key归并”的形式交给Reduce端处理。整个过程可以拆成这几步分区Partition。每个Map输出会先经过Partitioner决定每条数据应该去哪个Reduce任务。默认的分区算法是对key做哈希后对Reduce任务数取模hash(key) % R。这样相同key的数据会被分到同一个分区也就保证了“同一个key的最终结果会落到同一个Reducer”。排序Sort。每个分区内部框架会按照key进行排序。排序是MapReduce的默认行为哪怕你的业务不要求有序框架也一定会排。Combiner本地聚合。如果设置了Combiner会在Map端先做一次“预聚合”。比如WordCount场景中一个Map任务统计出(hello, 1)出现过1000次如果直接传给Reduce就是1000条记录有了Combiner它在Map端先合并成(hello, 1000)再传出去。这对减少网络传输量效果显著。归并Merge。Map端输出被Reduce端拉取后Reducer会按照key把所有来自不同Map任务的记录归并到一起。完成后每个key对应一个value迭代器供reduce方法遍历。如果你去看YARN的UI界面会发现Shuffle阶段往往占整个作业运行时间的一半以上“调优不调Shuffle、等于白调”这是有道理的。2.4 Reducer端处理与最终输出Reduce端的业务逻辑同样由你继承Reducer类、重写reduce方法来实现。reduce方法的输入是“key 一个value集合”输出是“计算后的key-value”最终写到HDFS上。这里有个面试常考的点reduce方法的values迭代器只能遍历一次。如果你需要对这个集合做多次操作比如先求总和再求平均值必须先把数据复制到本地集合里——很多人在这上面翻过车因为“迭代器只能消费一次”这个特性不太直观。输出阶段会通过OutputFormat来决定输出文件的格式和落盘位置。默认TextOutputFormat会把每个key-value写成一行key和value之间用制表符隔开。你可以自定义OutputFormat把结果写成分区文件、列式文件比如Parquet或者写入外部数据库。2.5 从WordCount全流程看数据是怎么流动的用最经典的WordCount把整个流程串一遍你会看得更清楚。假设有一个文本文件内容共三行hello world hello hadoop hadoop mapreduce整体数据流动如下输入分片阶段文件被切成一个或几个分片RecordReader按行读取输出(0, hello world)、(11, hello hadoop)、(24, hadoop mapreduce)这几个键值对给Map。Map阶段map方法按空格分词输出(hello, 1)、(world, 1)、(hello, 1)、(hadoop, 1)、(hadoop, 1)、(mapreduce, 1)。Shuffle阶段Partitioner根据key分配到对应分区相同key聚到一起内部排序后到Reduce端变成(hello, [1, 1])、(world, [1])、(hadoop, [1, 1])、(mapreduce, [1])。Reduce阶段reduce方法遍历value集合累加输出(hello, 2)、(world, 1)、(hadoop, 2)、(mapreduce, 1)。输出阶段结果按key \t value格式写入HDFS。整个过程看上去简单但你只要能在脑子里还原“每条数据在每个阶段的形态变化”MapReduce基本就吃透了。3. 实操环境搭建与MapReduce编程实例3.1 Hadoop环境准备从伪分布式开始最稳学习阶段不建议一上来就搞多台机器的集群。先在虚拟机或自己的电脑上用单机部署Hadoop的伪分布式模式所有守护进程跑在同一台机器上足够跑通MapReduce作业的开发链路也能帮助你理解各个组件进程之间的关系。伪分布式搭建需要准备几件事JDK版本。Hadoop 3.x版本建议用JDK 8或JDK 11太新的JDK版本可能会遇到一些兼容性问题。装好后记得设置JAVA_HOME环境变量。SSH免密登录。虽然本机也需要但Hadoop启停时需要SSH连接本机来拉起进程建议配置一下。执行ssh-keygen -t rsa生成密钥然后把公钥追加到authorized_keys里。Hadoop解压与配置。下载Hadoop 3.x的二进制包解压后需要修改几个核心配置文件core-site.xml设置fs.defaultFS为hdfs://localhost:9000指定NameNode地址。hdfs-site.xml设置dfs.replication为1单机只有一份副本以及NameNode和DataNode存储目录。mapred-site.xml设置mapreduce.framework.name为yarn让MapReduce作业跑在YARN上。yarn-site.xml设置YARN的ResourceManager地址和NodeManager的附属服务需要启用mapreduce_shuffle辅助服务否则作业在Shuffle阶段会报错。配置好之后先格式化NameNodehdfs namenode -format注意每次重新格式化前要清空NameNode和DataNode的数据目录否则会出现集群ID不一致的问题这是一个很经典的坑。然后执行start-dfs.sh和start-yarn.sh用jps命令检查NameNode、DataNode、ResourceManager、NodeManager这几个进程是否都在。3.2 创建Maven工程并配置依赖我建议用Maven来管理项目依赖这样不用手动下载一堆jar包。创建一个普通的Maven工程然后在pom.xml里加上Hadoop客户端依赖。properties hadoop.version3.3.6/hadoop.version /properties dependencies dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version${hadoop.version}/version /dependency /dependencies build plugins plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-compiler-plugin/artifactId version3.8.1/version configuration source1.8/source target1.8/target /configuration /plugin plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-shade-plugin/artifactId version3.2.4/version executions execution phasepackage/phase goals goalshade/goal /goals /execution /executions /plugin /plugins /buildshade插件会打一个“胖jar包”把所有依赖都打进去提交作业时不容易出现ClassNotFound问题。3.3 编写WordCount代码并逐段理解以下是可运行的完整WordCount程序。注意这段代码不是“背下来就行了”我建议你手敲一遍每个类名、每个参数都去查一下是什么意思。import java.io.IOException; import java.util.StringTokenizer; 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.Mapper; import org.apache.hadoop.mapreduce.Reducer; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; public class WordCount { // Mapper类继承Mapper四个泛型依次是输入key类型、输入value类型、输出key类型、输出value类型 public static class TokenizerMapper extends MapperObject, Text, Text, IntWritable { private final static IntWritable one new IntWritable(1); private Text word new Text(); Override 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); } } } // Reducer类四个泛型依次是输入的key类型、输入的value类型、输出的key类型、输出的value类型 public static class IntSumReducer extends ReducerText, IntWritable, Text, IntWritable { private IntWritable result new IntWritable(); Override public 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); } } 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); } }有几个细节值得多说一句代码里job.setCombinerClass(IntSumReducer.class)直接复用了Reducer类作为Combiner。这样做在“求和”这种场景下没问题因为局部求和再全局求和的结果等价于直接全局求和。但不是所有场景都能直接复用后面我会专门讲这个问题。Hadoop序列化机制用的不是Java原生的int和String而是IntWritable和Text。这是为了减少序列化后的体积同时让对象可复用、减少频繁创建对象的开销。job.waitForCompletion(true)传入true表示作业运行期间持续打印进度信息方便你观察Map和Reduce的百分比。3.4 打包、提交作业与查看运行日志代码写完后在项目根目录执行mvn clean package -DskipTests在target目录下会生成一个带依赖的jar包比如hadoop-demo-1.0-SNAPSHOT.jar。接下来准备测试数据。在HDFS上创建输入目录上传一个本地文本文件hdfs dfs -mkdir -p /input hdfs dfs -put /tmp/wordcount.txt /input/注意输出目录不能事先存在。如果/output已存在作业会直接报错这是MapReduce为了防止覆盖数据而做的保护。如果你确认可以覆盖需要先手动删除hdfs dfs -rm -r /output然后提交作业hadoop jar target/hadoop-demo-1.0-SNAPSHOT.jar WordCount /input /output作业提交后终端会滚动输出Map和Reduce的进度。你可以看到类似map 100% reduce 50%的信息。执行完成后用下面命令查看结果hdfs dfs -cat /output/part-r-00000如果一切顺利你能看到每个单词和它的出现次数。注意输出文件名是part-r-00000r表示来自Reducer如果你跑的是纯Map任务会出现part-m-00000。3.5 在YARN界面观察作业运行过程作业运行期间打开浏览器访问http://localhost:8088进入YARN的资源管理界面。你能看到当前作业的状态、使用的资源量、运行时间。点进作业详情页能看到每一个Map任务和Reduce任务的状态以及任务的日志入口。这步一定要亲自做一次。我见过太多人作业跑完就完事了从不看界面。实际上YARN的日志是你排查问题最宝贵的工具——某个任务失败的原因、GC耗时、Shuffle字节数全在里面。养成“跑完作业看日志”的习惯比多背十个概念都管用。4. 常见问题与排查技巧实录4.1 最常见的报错现象速查表下面这些报错是我在学习和实际使用中遇到频率最高的整理成表格方便你对照排查。报错现象可能原因处理方式Input path does not exist输入路径写错或HDFS上根本没有这个目录用hdfs dfs -ls /确认路径是否存在Output directory already exists输出目录已存在框架防止覆盖删除旧目录或换一个新输出路径Container exited with a non-zero exit code 1日志里通常藏着真正的异常任务代码抛错查看Container日志找Exception关键字定位原因java.lang.ClassNotFoundException提交作业时没有打依赖包或类名写错确认使用shade插件打胖包检查setJarByClassjava.lang.NoClassDefFoundError: org/apache/hadoop/crypto/...运行时缺少相关Hadoop依赖检查classpath确认hadoop-common依赖被正确打包Virtual memory exceeded容器使用的虚拟内存超过YARN限制调大yarn.nodemanager.vmem-pmem-ratio或检查是否有内存泄漏connection refused相关服务没启动或端口配错用jps确认进程用netstat检查端口监听状态4.2 作业一直卡在Running状态不动了这是初学阶段很崩溃的一个问题。明明sumbit成功了但进度一直卡在某个百分比或者Map跑到100%、Reduce迟迟不启动。排查思路一般按以下顺序先看YARN界面。如果资源剩余不足集群只有1个NodeManager又被其他作业占满新作业只能排队等待。学习环境最常见的问题是上一个作业虽然结束了但容器资源没有完全释放隔几秒再提交就好了。再看日志。Map阶段100%但Reduce未启动很可能是Shuffle拉取数据时出了问题。一个高频原因是我之前说过的yarn-site.xml里没有配置mapreduce_shuffle辅助服务。这时候Reduce一直处于SHUFFLE状态让你误以为卡死。最后看数据量。如果Map输出量特别大Shuffle阶段确实需要较长时间这是正常现象可以耐心多看几分钟。4.3 本地跑IDEA直接报错别慌很多初学者喜欢在IDEA里直接右键运行main方法然后发现各种问题——明明集群能跑本地却报错。这是因为本地运行时没有Hadoop的native库和支持包Windows下尤其容易出问题。我的建议是学习阶段统一走“打包上传到服务器用命令行提交”的流程。这样既符合生产习惯也能避开本地环境各种坑。等你对Hadoop足够熟悉了再配置IDEA的远程调试或者本地模拟环境。如果你一定要在本地调试至少要做两件事配置HADOOP_HOME环境变量并在IDEA里将bin目录对应的native库路径加入java.library.path同时把core-site.xml等文件放到classpath下否则本地模式会使用默认配置连接不到集群。4.4 多个小文件拖垮性能怎么破HDFS对大量小文件非常不友好每个小文件都会产生一个数据块进而可能产生一个Map任务。如果有上万个小文件就会启动上万个Map任务调度开销远超计算开销作业会被活活拖慢。处理思路有三个用HDFS的hdfs dfs -getmerge把多个小文件合并成一个文件再上传适合一次性处理。在输入前用自定义InputFormat把小文件合并成一个大分片再交给Map处理可以把多个小文件的内容封装成一个逻辑文件。如果数据源可控写入时尽量使用SequenceFile或ORC等格式从源头上规避小文件问题。5. 进阶调优与实践心得5.1 Combiner不是“想用就能用”Combiner作为Map端的预聚合工具优化效果立竿见影尤其是统计、求和、最大值这类场景。但要注意Combiner的输入和输出格式必须与Mapper的输出一致因为它本质上是在Map端执行了一次“局部Reduce”。更重要的是使用Combiner前你必须保证“重复计算不会改变结果”。比如求平均值就不能直接复用Reducer逻辑——Map端求了平均Reduce端再求平均最终结果是错的因为每个Map任务的记录数权重不同。正确做法是求平均值时Map端输出(key, (sum, count))这样组合的value对象Combiner做“局部求和与计数”Reduce端再做最终除法。包括用Combiner前一定要写个场景测试验证业务等价性别盲目套用。5.2 数据倾斜有效解决“一半节点忙死一半节点闲死”数据倾斜是分布式计算里最常遇到的性能杀手。表现是大部分Reduce任务很快跑完个别Reduce任务跑了很久甚至整个作业卡在那几个任务上。原因通常是某个key的数据量远超其他key。比如统计“全国人口”时某个超级大城市的key数据量特别大这个Reduce任务就必然比其他任务慢。常见缓解手段自定义Partitioner把数据量大的key单独拆成多个分区分散到多个Reducer。给倾斜的key加随机前缀打散后在不同Map端聚合再二次聚合消除前缀影响。增大Reduce任务数量让单个Reducer处理的数据量降下来。如果倾斜发生在Map端某些Map处理的数据块特别大需要先解决HDFS上数据分布不均的问题。5.3 常用参数调优参考下面是我在实践中比较常用的MapReduce调优参数给出了适合入门和中等规模集群的参考值。请根据自己的集群硬件条件调整不要照搬。!-- 每个Reducer处理的数据量建议值约256MB~1GB -- property namemapreduce.job.reduces/name value3/value /property !-- Map端内存缓冲区大小默认100MB可适当调大 -- property namemapreduce.task.io.sort.mb/name value256/value /property !-- Map端排序/Spill时触发溢写比例调大可减少溢写次数 -- property namemapreduce.map.sort.spill.percent/name value0.8/value /property !-- 开启Map端输出压缩 -- property namemapreduce.map.output.compress/name valuetrue/value /property property namemapreduce.map.output.compress.codec/name valueorg.apache.hadoop.io.compress.SnappyCodec/value /property !-- 调整Map和Reduce容器内存根据机器内存修改 -- property namemapreduce.map.memory.mb/name value1024/value /property property namemapreduce.reduce.memory.mb/name value2048/value /property尤其推荐开启Snappy压缩。Map端输出往往有大量重复字段压缩后Shuffle阶段网络传输量能减少60%~80%代价仅是少量的CPU开销。对于IO密集型的作业收益非常明显。5.4 面试常问的几个MapReduce深层问题把这几个问题想明白你对MapReduce的理解会上升一个台阶。MapReduce的二次排序Secondary Sort是怎么实现的核心思想是把组合key自然key 排序字段一起作为key自定义分区只按自然key分区、自定义排序按组合key排序再通过GroupingComparator只按自然key分组。为什么Map端输出不直接写HDFS因为中间结果只服务于当前作业不需要多副本容错写本地磁盘可以减少落盘和网络传输开销。如果Reduce端接收到的values无序MapReduce如何保证实际上Shuffle传输后Reducer端接收到的values是按照key排序后的迭代器但values之间的相对顺序并不保证稳定所以不要依赖values的顺序性。一个Map任务最多能处理多少数据取决于分片InputSplit大小HDFS块大小128MB意味着一个Map默认处理不超过128MB的数据具体取决于分片切分逻辑。5.5 学习路径建议别在MapReduce上恋战最后说一点个人体会。MapReduce虽然经典但由于中间结果频繁落盘、Shuffle开销大它在实时性和迭代计算场景下越来越吃力。这也是为什么后来会出现Spark、Flink这些更高效的计算框架。但我不建议你因此跳过MapReduce。它是理解“分布式计算思维”的最佳教材——分治、数据本地性、Shuffle、容错这些概念都源自MapReduce。你把MapReduce搞透了再去学Spark的RDD、Shuffle机制会发现很多概念是相通的学习成本会大幅降低。如果你正要准备面试建议把本文提到的几个核心机制、数据流动流程、WordCount实现烂熟于心。我在面试中问候选人的经验是能画出MapReduce全流程数据流图、能说出数据倾斜三个解决思路的人基本都算真学过。如果只能背概念一追问就露馅所以动手跑一次、打开日志看一次比死记硬背有效得多。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进