
1. 项目概述MapReduce 是怎么一步步成为大数据计算基石的做大数据这块绕不开 Hadoop而 Hadoop 里最让人又爱又恨的就是 MapReduce。很多人一开始接触它觉得代码写起来比 Spark 啰嗦跑起来也比 Spark 慢但它确实是大数据处理历史上的一座里程碑也是理解分布式计算思维的必修课。如果你后面想搞懂 Hive、Spark 甚至 Flink 的底层调度逻辑MapReduce 这套“分而治之”的思想是无论如何都要吃透的。这篇文章我想从一个实际落地者的角度把 MapReduce 的核心机制、代码实战、常见报错、调优思路完整过一遍。我不打算只贴几个 API 就完事而是会把 MapReduce 为什么要这么设计、里面的 shuffle 到底经历了什么、写作业时哪些参数最容易踩坑都展开聊聊。适合人群很明确刚学完 Hadoop 基础、准备做课程设计或者面试前突击的初学者以及写了好几个 MR 作业但一直没搞明白内部机制的“半桶水”同学。文章里的案例基于 Hadoop 2.x/3.x 版本我在本地伪分布式和真实集群上都跑过。2. 分治思想与整体设计拆解2.1 MapReduce 解决的核心问题是什么首先我们需要回到一个问题在没有 MapReduce 之前处理大规模数据为什么难假设你有一个 1TB 的日志文件存在一台机器上单机 CPU 去扫描 统计可能要跑十几个小时甚至几天。更麻烦的是如果这份数据是持续增长的单机的磁盘和内存很快也会见底。这时候大家自然想到把数据拆成很多块分给多台机器同时算最后再汇总结果。想法很简单但实际落地会碰上一堆麻烦数据怎么切分切成多少份合适哪些机器去处理哪些分片机器宕机了怎么办每台机器算完的中间结果怎么合并网络传输、磁盘读写怎么保证效率MapReduce 的核心贡献就是把这堆问题抽象成了两个阶段Map映射和Reduce归约然后把中间的调度、容错、数据分发全部交给框架处理程序员只需要专注于“怎么把一条记录变成 KV 对”和“怎么把相同 Key 的 Value 列表聚合起来”。我举个生活化的类比。假设你要统计全校 100 个班、每班 50 人的姓氏分布。常规做法是找 100 个老师每个老师负责一个班把自己班里学生的姓氏写在纸条上这就是 Map 阶段每个班独立统计互不干扰然后你安排 10 个组长把写着同姓的纸条归到同一个组里汇总数量这就是 Reduce 阶段相同 Key 的 Value 汇聚。MapReduce 做的事情就是帮你自动完成“纸条传递”和“分组汇总”的过程哪怕某个老师中途请假它也能换人重算。2.2 为什么选择“移动计算而非移动数据”这个策略MapReduce 有一个非常核心的设计原则数据不动计算移动。这个原则新手很容易忽视但它直接决定了集群的调度策略。我们知道 HDFS 上数据是分块存储的默认一个块 128MB并且每个块有 3 个副本分布在不同的机器上Hadoop 3.x 默认副本数为 3。如果一个 Map 任务的输入数据在节点 A 上而 YARN 把这个任务调度到了节点 B 上那就意味着节点 B 要从 A 通过网络拉取 128MB 数据。如果整个作业有 10000 个 Map 任务网络传输量就是 TB 级别的这显然不可接受。所以 Hadoop 的做法是任务调度器会优先把 Map 任务分配给“输入数据所在节点”这个策略叫“数据本地性Data Locality”。如果数据块的副本恰好都在繁忙节点上调度器才会退而求其次选择同机架的其他节点Rack Locality最后才是跨机架Off-Switch。这也是为什么 MapReduce 适合“数据密集型”任务——它假设数据已经躺在 HDFS 上了只是把计算逻辑分发过去执行。注意这里的“移动计算”指的是把 JAR 包和任务配置分发到数据节点而不是把整个程序安装到每个节点。YARN 的 NodeManager 会在本地节点上启动一个 Container然后在 Container 里运行你的 Mapper 类。2.3 一个 MapReduce 作业的完整生命周期如果用一个流程图来描述MapReduce 作业从提交到结束会经历这些环节客户端向 ResourceManager 提交作业包括 JAR 包、输入路径、输出路径、各种配置参数。ResourceManager 分配一个 Container 用于启动 ApplicationMasterApplicationMaster 负责整个作业的调度与监控。ApplicationMaster 根据输入数据的分片信息向 ResourceManager 申请一组 Container 来运行 Map 任务。每个 Map 任务读取一个 InputSplit 对应的数据调用用户自定义的 map() 函数逐条处理输出中间结果KV 对。Map 任务输出的中间结果先写入环形缓冲区在溢写Spill过程中进行分区、排序、合并Combiner 可选最终写到一个或多个输出文件里。Reduce 任务从各个 Map 任务的输出文件里拉取属于自己分区的数据Shuffle 阶段然后对数据进行再次排序和分组。Reduce 阶段调用用户自定义的 reduce() 函数对相同 Key 的所有 Value 做聚合计算最终输出结果到 HDFS。这里我想特别强调第 4 步里的 InputSplit。InputSplit 并不是数据本身而是一个“逻辑描述”它告诉 Map 任务你要处理的数据在哪个文件的哪个字节区间。默认情况下Hadoop 会根据 HDFS 的块大小把文件切成若干个 Split但注意逻辑分片可以大于一个 HDFS 块这是很多性能问题的根源之一后面调优部分我会细说。3. 核心机制深度解析从 Map 到 Reduce中间发生了什么3.1 Map 阶段不是“处理完就完事”它内部藏着三道工序很多人以为 Map 阶段就是循环调用 map() 函数处理数据处理完了输出 KV 对就结束了。真实情况要复杂得多。一个 MapTask 内部从 map() 函数产生输出到数据真正落盘要经过分区 → 环形缓冲区 → 溢写排序三道工序。第一道工序分区Partition每一条 map() 输出的 KV 对框架首先要确定它应该交给哪个 Reduce 去处理。默认的分区器是 HashPartitioner它的逻辑非常简单public class HashPartitionerK, V extends PartitionerK, V { public int getPartition(K key, V value, int numReduceTasks) { return (key.hashCode() Integer.MAX_VALUE) % numReduceTasks; } }也就是说对于同一个 Key无论它在哪台机器的哪个 Map 任务里输出计算出来的分区号一定相同这样才能保证“同一个 Key 的所有记录最终都进入同一个 Reduce”。如果 Reduce 数量设为 1所有数据都会进同一个分区如果 Reduce 数量大于 1数据会被打散到多个分区中。这里有个新手容易犯的错误如果 Reduce 数量设置得很大但数据本身的 Key 分布不均匀就会出现数据倾斜。比如某个热门 Key 占了总数据量的 99%那么无论你有多少个 Reduce这 99% 的数据都会涌向同一个分区造成某个 Reduce 任务要处理的数据量远超其他任务。这种情况在后面问题排查部分我会给出对策。第二道工序环形缓冲区map() 函数每输出一条 KV 对并不是直接写到磁盘而是先写入一个内存中的环形缓冲区。这个缓冲区默认大小是 100MB可以通过mapreduce.task.io.sort.mb参数调整。缓冲区使用率超过阈值默认 80%时一个后台线程会把缓冲区中的数据溢写到磁盘。为什么是“环形”因为这块内存被设计成可以循环复用的数据不停地写入溢写线程不停地把数据刷到磁盘腾出空间继续写入。如果缓冲区写满而溢写线程又来不及刷盘map() 函数就会阻塞等待——这是 Map 阶段最常见的性能瓶颈之一。第三道工序溢写Spill与排序每次溢写缓冲区中的数据并不是原样写到磁盘的。在溢写之前后台线程会先对这些 KV 对进行一次分区内排序排序规则默认按 Key 的字典序可以通过mapreduce.job.output.key.comparator.class自定义比较器。如果有设置 Combiner还会在溢写过程中先做一次局部合并Map 端的合并减少网络传输量。多次溢写会产生多个 Spill 文件最终 Map 任务结束前这些 Spill 文件会被合并成一个有序的大文件。这个合并过程同样会做二次排序同时把相同 Key 的 Value 合并成一个列表方便后面 Reduce 去拉取。3.2 Shuffle 阶段Map 和 Reduce 之间的“快递干线”Shuffle 是 MapReduce 里最核心也最容易被忽略的环节。它描述的是Map 任务的输出如何被 Reduce 任务拉取的过程。整个过程横跨 Map 端和 Reduce 端我分开说。Map 端的 shuffleMap 任务每次溢写落盘后其实已经按分区把数据分好了。真正要告知 Reduce 的是我这个分区在哪个文件、什么位置。这个信息由 MapTask 的 MapOutputBuffer 通过 YARN 的 ShuffleHandler 服务来提供。Reduce 任务会向各个 Map 任务的节点发起 HTTP 请求拉取属于自己分区的数据。这里有一个关键点Map 任务输出文件的生命周期。只有当所有 Map 任务都运行完成且对应的 Reduce 任务成功拉取完数据后这些中间文件才会被清理。如果某个 Reduce 任务失败并对同一个 Map 任务发起重试拉取而 Map 任务已经结束并清理了中间文件就会出现“Fetch failure”。所以 Hadoop 对 Map 端的中间文件有一个保留策略具体由mapreduce.task.files.preserve.failedtasks和mapreduce.task.files.preserve.filepattern等参数控制。Reduce 端的 shuffleReduce 任务并不是等所有 Map 任务跑完才开始拉数据的而是会启动一个专门的 copier 线程持续从已完成 Map 任务的节点上拉取数据。拉取到的数据先放入 Reduce 端的内存缓冲区默认大小可通过mapreduce.reduce.shuffle.input.buffer.percent等参数调整当缓冲区达到一定比例时数据会溢写到 Reduce 节点的本地磁盘。Reduce 会有多个“合并Merge”阶段先把从不同 Map 节点拉来的数据合并成有序的批次最后在真正执行 reduce() 函数之前再把这些批次合并成一组“有序的、按 Key 分组的数据流”。也就是说reduce() 函数看到的输入是按 Key 有序排列的、且相同 Key 的所有 Value 连续出现的迭代器。这一点非常关键如果你的业务逻辑需要对同 Key 的多个 Value 做全局排序就必须依赖这个有序性。3.3 YARN 在 MapReduce 里扮演的角色MapReduce 在 Hadoop 1.x 时代有独立的 JobTracker 和 TaskTracker 管理资源与任务到了 Hadoop 2.x 之后资源管理和任务调度被拆分成了 YARN。YARN 的出现让 Hadoop 不再只服务于 MapReduceSpark、Flink 也能跑在同一个集群上。从这个角度来说理解 MapReduce 和 YARN 的配合关系对后面理解整个大数据生态至关重要。YARN 的核心组件有四个ResourceManagerRM、NodeManagerNM、ApplicationMasterAM、Container。ResourceManager 负责全局资源管理和调度NodeManager 负责单个节点上的资源管理和 Container 生命周期ApplicationMaster 则是每个作业的“大脑”负责向 RM 申请资源、监控任务进度、处理任务失败重试。Container 是资源分配的最小单位它封装了 CPU、内存、磁盘等资源。一个 MapReduce 作业在 YARN 上的运行流程前面已经提过但有一个细节值得展开AM 和 Map/Reduce 任务的关系是最典型的“双重调度”。第一重调度是 RM 给 AM 分配资源第二重调度是 AM 内部再把拿到的 Container 分配给自己管理的各个 Map/Reduce 任务。这种设计的好处是作业级别的调度策略可以由 AM 自定义每个 MapReduce 作业都可以有自己的调度逻辑而 RM 只负责最底层的资源切片。注意AM 本身也是一个 Container只不过它运行的是 MapReduceAppMaster 而不是用户代码。如果 AM 所在的节点宕机RM 会在另一个节点上重新拉起 AM并恢复之前作业的状态。所以我们在集群运维时经常强调“ResourceManager 最好单独部署不要和 DataNode 混布”就是为了避免 RM 宕机导致整个集群的资源调度瘫痪。4. 实操从零写一个 WordCount 并跑通伪分布式环境4.1 环境准备与输入数据规划在写代码之前先把环境准备好。我用的版本是 Hadoop 3.3.4JDK 1.8Ubuntu 20.04。如果你还没搭好环境强烈建议先用伪分布式模式把整个流程跑通再去折腾集群。伪分布式意味着所有守护进程NameNode、DataNode、ResourceManager、NodeManager都在本机运行非常适合学习和测试。检查 Hadoop 是否正常启动jps # 期望看到NameNode、DataNode、ResourceManager、NodeManager、SecondaryNameNode准备输入数据这里我造一个简单文件模拟日志文本hdfs dfs -mkdir -p /user/example/input echo -e hello hadoop\nhello mapreduce\nhadoop is powerful\nmapreduce is powerful\nhello world /tmp/input.txt hdfs dfs -put /tmp/input.txt /user/example/input/4.2 Maven 工程搭建与核心代码我推荐用 Maven 管理 Hadoop 项目的依赖。这里有个小坑Hadoop 3.x 的hadoop-client依赖如果直接引入会把一大堆传递依赖打进来导致 JAR 包非常大。建议用provided作用域打包时排除 Hadoop 自身的依赖。dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version3.3.4/version scopeprovided/scope /dependency然后编写 WordCount 的代码。下面这个版本我在类上加了注释方便你对照理解import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.LongWritable; 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; import java.io.IOException; public class WordCount { /** * Mapper 类 * 输入 Key 是行偏移量LongWritable输入 Value 是一行文本Text * 输出 Key 是单词Text输出 Value 是次数 1IntWritable */ public static class TokenizerMapper extends MapperLongWritable, Text, Text, IntWritable { private final static 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(); String[] words line.split(\\s); for (String w : words) { if (w.isEmpty()) { continue; } word.set(w); context.write(word, one); } } } /** * Reducer 类 * 输入 Key 是单词输入 Value 是迭代器包含该单词出现过的所有 1 * 输出 Key 是单词输出 Value 是该单词的总次数 */ public static class IntSumReducer extends ReducerText, IntWritable, Text, IntWritable { private final 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); } } 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); // 这里直接复用 Reducer 作为 Combiner 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); } }有两点我要特别提醒Combiner 的设置。我直接复用了 Reducer 类作为 Combiner因为 WordCount 的聚合操作满足交换律和结合律。但如果你遇到像“求平均值”这样的操作直接复用 Reducer 作为 Combiner 就会出错。比如你有一组 (3, 5, 7)Map 端如果先求出平均值 5再把 (5, 1) 传给 Reduce多个 Map 的局部均值汇总后再求平均和全局求均值结果完全不一样。所以 Combiner 能不能用必须看业务逻辑是否满足“局部汇总后再汇总”等价于“全局汇总”。Map 输出 Key/Value 类型的设置。如果 Mapper 的输出类型和 Reducer 的输出类型不一致就必须显式调用job.setMapOutputKeyClass()和job.setMapOutputValueClass()。因为 Map 端和 Reduce 端的序列化/反序列化配置不同很多初学者在这里踩坑报错信息通常是java.io.IOException: Type mismatch in key from map。4.3 打包并提交作业到 YARN在项目根目录执行 Maven 打包mvn clean package -DskipTests注意Hadoop 的 Driver 类在提交作业时需要把依赖的第三方库一并打包或者通过-libjars参数指定。我们的 WordCount 除了 Hadoop 本身的类没有其他依赖所以直接打包即可。如果你的项目引用了 Jackson、Guava 等库最稳妥的做法是用 Maven Shade 插件打一个 fat jar或者把依赖放到 HDFS 的某个公共路径再用mapreduce.job.classpath.files指定。打包完成后提交作业hadoop jar target/wordcount-1.0-SNAPSHOT.jar com.example.WordCount /user/example/input /user/example/output如果一切正常控制台会输出每个阶段的进度信息。跑完后检查结果hdfs dfs -cat /user/example/output/part-r-00000输出应该是按 Key 字典序排列的单词统计结果和预想的一致。4.4 跑一个实战案例日志流量统计只看 WordCount 肯定不够过瘾我再说一个真实业务中经常出现的场景统计某天访问日志中每个 URL 的请求次数和总流量。日志行格式类似192.168.1.1 - - [09/Dec/2024:10:23:45 0800] GET /index.html HTTP/1.1 200 532 192.168.1.2 - - [09/Dec/2024:10:23:46 0800] GET /api/users HTTP/1.1 200 1024我们关心的是 URL 和响应大小。这里直接解析字符串public static class AccessLogMapper extends MapperLongWritable, Text, Text, AccessLogWritable { private final Text url new Text(); private final AccessLogWritable out new AccessLogWritable(); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line value.toString(); // 简单解析找到引号内的 URL以及最后一个字段的状态大小 int firstQuote line.indexOf(\); int secondQuote line.indexOf(\, firstQuote 1); if (firstQuote 0 || secondQuote 0) return; String requestLine line.substring(firstQuote 1, secondQuote); String[] parts requestLine.split( ); if (parts.length 2) return; url.set(parts[1]); long bytes; try { bytes Long.parseLong(line.substring(line.lastIndexOf( ) 1)); } catch (NumberFormatException e) { bytes 0; } out.set(1, bytes); context.write(url, out); } }这里我自定义了一个AccessLogWritable作为 Map 输出的 Value它同时携带请求次数和流量字节数两个字段。这就要求我们实现 Writable 接口重写write()和readFields()方法。自定义 Writable 在 MapReduce 中非常常见因为默认的 IntWritable 只能携带一个值。这块知识如果你还不熟建议认真学一下它是很多复杂业务场景的基础。Reduce 端只需要把次数累加、流量累加然后输出即可。跑完的结果就是每个 URL 的两列汇总数据可以非常直观地看到哪些接口是热点、哪些接口消耗的带宽大。5. 常见问题排查与避坑实录5.1 经典的“jar does not exist or is not a normal file”错误打开搜索热词你会发现有一条非常具体的报错信息jar does not exist or is not a normal file: /usr/local/hadoop/share/hadoop/m...这条报错几乎每个新手都遇到过通常是在执行hadoop jar xxx.jar或yarn jar xxx.jar命令时出现而且报错路径指向 Hadoop 安装目录下 share/hadoop/mapreduce 目录。我先说一下根因这个报错的原因是 YARN 的 ApplicationMaster 在分发 JAR 时会尝试解析 JAR 包中的META-INF/MANIFEST.MF里的Class-Path、Main-Class等属性。如果你的 JAR 包是在 Windows 环境下用 IDE 或某些打包工具生成的MANIFEST 里可能残留了Class-Path: /usr/local/hadoop/share/hadoop/mapreduce/hadoop-mapreduce-client-core-3.3.4.jar这样的绝对路径。当 YARN 尝试把这个 JAR 包分发到各个 NodeManager 时发现这个路径在目标节点上不存在或者权限不对就会报这个错误。解决办法有以下几种重新打包去除 MANIFEST 中的绝对 Class-Path。用 Maven 时显式配置MAIN CLASS和依赖 scope避免把 Hadoop 自身的 JAR 打包进 fat jar 或写入 Class-Path。检查 JAR 包类型。Hadoop 对包含依赖的 fat jar 识别确实有些挑剔如果确认打包没问题还可以尝试直接执行java -cp方式运行 Driver 类的 main 方法前提是本地类路径里有 Hadoop 依赖。换一种提交方式。比如用yarn jar替代hadoop jar或者把 JAR 先放到 HDFS 上再通过-files参数指定。这几种方式对错误定位都有帮助。5.2 Map 阶段内存溢出Java heap spaceMap 阶段报OutOfMemoryError: Java heap space通常不是 JVM 默认堆内存太小而是环形缓冲区调整不当导致 map() 函数写入大量 KV 时内存不够。我们知道环形缓冲区默认 100MB配合溢写阈值 80%意味着 map() 函数最多可以写入约 80MB 的中间数据而不触发溢写。如果每条 KV 都很大比如一条记录的 Value 有几 MB那么内存中能容纳的 KV 数量就非常少溢写频繁甚至还没触发溢写内存就爆了。应对策略有几个方向调大mapreduce.task.io.sort.mb比如从 100 调整到 200 或 256但不要超过 NodeManager 给这个 Container 分配的内存。减小区块溢写阈值mapreduce.map.sort.spill.percent让溢写更早发生避免内存峰值过高。这个方法治标不治本但能缓解偶发性的 OOM。最根本的办法是优化 map() 函数减少无意义的大对象创建同时确认业务上是否真的需要缓存大量数据。另外还要检查mapreduce.map.memory.mb配置。这个参数是 YARN 为 Map 任务分配的物理内存上限默认 1024MB。如果你的 map 任务需要更多内存可以在提交作业时通过-Dmapreduce.map.memory.mb2048指定。注意JVM 堆内存mapreduce.map.java.opts和 Container 物理内存是两个概念JVM 堆默认是物理内存的 75% 左右也就是 768MB。如果你调大了物理内存但没调 JVM 堆Java heap space 还是不会变。5.3 Reduce 阶段一直卡在 33.33%很多新手跑 MR 作业发现 Reduce 进度停在 33.33% 不动了以为程序死锁或者集群卡死。其实这是正常现象原因是 Reduce 被划分成三个阶段Shuffle33.33%Reduce 从各个 Map 任务节点拉取数据。Sort66.67%Reduce 把拉取到的数据在内存和磁盘上进行合并排序。Reduce100%执行 reduce() 函数。当 Reduce 卡在 33.33% 时意味着它还在等 Map 阶段完成。如果 Map 任务有部分还没有跑完比如有 Map 任务失败正在重试Reduce 就会一直等待。所以你看到 33.33% 不要慌先去看 Map 任务的进度大概率是 Map 拖了后腿。还有一种情况Reduce 的 copier 线程一直在拉取数据但由于mapreduce.reduce.shuffle.parallelcopies默认只有 5 个并发线程如果有成百上千个 Map 任务拉取速度会很慢。适当调大并发线程数能显著缩短 Shuffle 时间但也不能无限调大否则会耗尽节点的网络带宽。5.4 数据倾斜运行时间从 3 分钟变 3 小时数据倾斜是分布式计算里最经典也最头疼的坑。我举一个真实例子有一次我统计一个电商订单表里每个用户的消费金额大部分用户只消费了不到 1 万元但有一个“测试用户”下单了几十万笔结果这个 Key 的数据量比其他用户多出几个数量级。跑起来之后其他 Reduce 任务几分钟就完成了唯独那个处理“测试用户”的 Reduce 跑了将近 2 小时而且因为数据量大频繁溢写磁盘最后甚至 OOM。解决数据倾斜的思路通常有两类过滤异常 Key。如果这些大 Key 本身不是我们关心的数据比如测试账号、爬虫用户直接在 Map 阶段把它们过滤掉问题消失。添加随机前缀Salting。如果这些大 Key 是必须统计的可以给 Map 输出的 Key 加上一个 0~N 的随机前缀把它们分散到 N 个 Reduce 任务里每个 Reduce 只处理原来 1/N 的数据量然后在 Reduce 阶段后再写一个作业去掉前缀做二次聚合。代价是至少要跑两轮作业但收益是整体时间大幅下降。此外在跑 Hive 任务时也可以利用 Hive 自带的skewjoin或者手动添加前缀的方式处理思路是通用的。5.5 常见问题速查表现象可能原因排查/解决方式提交作业时报jar does not existJAR 包 MANIFEST 中 Class-Path 含绝对路径重新打包去掉 Class-Path 或改为相对路径Map 阶段 OOM环形缓冲区过小或 JVM 堆内存不足调大mapreduce.task.io.sort.mb、mapreduce.map.memory.mbReduce 卡在 33.33%等待未完成的 Map 任务查看 Map 任务进度排查 Map 失败原因某个 Reduce 任务运行时间异常长数据倾斜过滤异常 Key 或添加随机前缀输出目录已存在导致失败FileOutputFormat 不允许覆盖同名目录删除输出目录或换一个新的路径本地运行正常但集群运行失败本地文件路径/权限问题把输入输出路径换成 HDFS 路径检查 HDFS 权限两个 DB 来源的日志时间格式不一致业务数据不规整在 Map 阶段做标准化解析正则匹配多种格式6. 调优经验让 MapReduce 作业跑得更快6.1 输入分片大小别让 Map 任务数变成“玄学”MapReduce 的 Map 任务数量并不是“越多越好”也不完全是“越少越好”而是和输入数据的大小与集群资源有关。默认情况下一个 InputSplit 对应一个 HDFS 块128MB也就是一个 Map 任务最多处理 128MB 数据。如果你的集群有 100 个可用的 Map 槽位而输入数据只有 50 个块那就会有 50 个槽位空闲资源利用率不足。反过来如果输入是大量小文件比如每个文件只有几十 KBFileInputFormat 会为每个文件生成一个 Split即使这个文件远小于 128MB。这样会产生成千上万个 MapTask而每个 MapTask 的启动、初始化、调度开销远大于它的计算时间整体效率极其低下。针对小文件场景常用的优化手段是在数据入口做合并把大量小文件合成大文件再入库。使用CombineFileInputFormat它可以把多个小文件打包成一个逻辑分片避免每个小文件单独成一个 Map 任务。如果数据是从 Kafka 之类的消息系统实时落地的可以在落盘时按时间窗口聚合。而针对大块数据如果单条记录很大比如每行是几 MB 的 JSON可以手动调整mapreduce.input.fileinputformat.split.maxsize和split.minsize让每个 Split 对应半个块或四分之一块从而增加 Map 并行度。这个参数的取舍逻辑是确保单个 Split 不会因为内存限制而 OOM同时尽量让 Map 任务数量和集群 CPU/内存资源匹配。6.2 合理使用 Combiner 和压缩Combiner 的作用我在前面已经解释过这里重点说压缩。MapReduce 中间数据量越大Shuffle 阶段的数据传输就越慢。对 Map 输出进行压缩可以大幅减少网络 I/O。Hadoop 默认支持 LZO、Snappy、Bzip2 等压缩格式生产环境我最推荐 Snappy压缩率高且解压速度快。启用 Map 输出压缩和 Reduce 输出压缩的配置如下property namemapreduce.map.output.compress/name valuetrue/value /property property namemapreduce.map.output.compress.codec/name valueorg.apache.hadoop.io.compress.SnappyCodec/value /property这个配置可以在mapred-site.xml里全局设置也可以在提交作业时通过-D参数动态指定。我实际测试过启用 Snappy 后一次连接日志分析作业的 Shuffle 时间大约缩短了 40%中间数据写盘量也明显减少。代价是 CPU 会多一些额外开销但对于 CPU 相对充裕、网络是瓶颈的集群来说收益非常可观。LZO 和 Snappy 怎么选如果你需要的是可切分Splittable的压缩格式需要注意Snappy 块压缩在 Hadoop 3.x 中默认支持但如果使用 LZO需要额外安装原生库。可切分意味着一个压缩文件可以分成多个块由多个 Map 任务并行处理这对大文件很关键。如果使用不可切分的压缩格式比如 Gzip一个文件只能由一个 Map 任务处理数据再多也白搭。所以当数据文件是压缩格式时我建议优先考虑可切分的大文件方案。6.3 Reduce 数量与分区策略的调整经验Reduce 数量怎么设置没有一个固定的公式但有几个经验参考Reduce 数量应尽量接近集群可用 CPU 核心数的整数倍。比如你有 20 个节点每节点 16 核地设置 320 个 Reduce 任务20×16就比较合理。如果 Reduce 数量过少每个 Reduce 需要处理的数据量会很大浪费并行能力如果 Reduce 数量过多大量小文件的写盘开销会让 NameNode 压力变大而且最终输出文件数量过多也不利于下游使用。对于不需要 Reduce 的作业比如简单的数据过滤、格式化输出可以显式把 Reduce 数量设为 0作业只执行 Map 阶段效率会高很多。设置方式-Dmapreduce.job.reduces0关于分区策略默认的 HashPartitioner 在 Key 分布不均匀时可能引发倾斜。你可以自定义 Partitioner根据业务特点把数据打散。比如统计用户消费时可以按照用户 ID 的区间范围分区让数据量更均衡。但要注意分区的逻辑改掉所有用户自定义的分区结果必须和后续 Reduce 的处理逻辑兼容否则会出错。6.4 Container 资源参数与集群稳定性的平衡调优有一条铁律尽量让单个 Container 的资源大小不超过一台机器的 1/4 到 1/2。原因很简单如果一台机器只有 16GB 内存你给一个 Container 分配 12GB那这台机器的其余任务都会被挤爆。YARN 的调度器虽然可以分配资源但它不会主动把任务从超载的机器上迁走很容易导致“连环 OOM”。我常用的参数调整思路是先确定每台机器的物理内存和 CPU 核心数。再设定单个 Container 的资源上限比如单容器给 4GB 内存、2 核 CPU。最后根据作业类型调整mapreduce.map.memory.mb和mapreduce.reduce.memory.mb。在伪分布式环境中资源本来就不多如果同时跑 Map 和 Reduce默认配置下可能出现“任务卡住”的现象。比如你的机器只有 4GB 内存而mapreduce.map.memory.mb默认 1GB、mapreduce.reduce.memory.mb默认 1GB加上 NameNode、DataNode、ResourceManager、NodeManager 自己吃掉的内存启动完所有 JVM 进程后系统内存可能就剩不下多少了。这时候跑作业往往会失败或变慢。正确做法是调小 Container 内存比如-Dmapreduce.map.memory.mb512 -Dmapreduce.reduce.memory.mb512 -Dmapreduce.map.java.opts-Xmx512m -Dmapreduce.reduce.java.opts-Xmx512m调成 512MB 后伪分布式基本就能平稳跑起来。7. MapReduce 相关生态工具联动7.1 Hive 与 MapReduce 的关系以及 Hive on MR 的执行流程很多人学完 MapReduce 会有个疑惑日常工作里我明明写的是 Hive SQL怎么它和 MapReduce 有关系这个关系在 Hive 默认执行引擎为 MapReduce 时非常直接每一条 Hive SQL 都会被翻译成一个或多个 MapReduce 作业翻译工作由 Hive 的编译器完成然后提交到 YARN 上执行。我举一个例子SELECT user_id, COUNT(*) FROM orders GROUP BY user_id;这条 SQL 翻译成 MapReduce 时大致流程是Map 阶段从订单表读取每行数据以 user_id 为 Key、1 为 Value 输出。Combine/Shuffle 阶段按 user_id 分组把相同 user_id 的 1 汇聚在一起。Reduce 阶段对每个 user_id 的 Value 列表求和输出聚合结果。你会发现这套流程本质上和手写 WordCount 的 Map 和 Reduce 没有任何区别。所以如果你对 MapReduce 底层机制理解得深那看 Hive SQL 的执行计划就很容易“脑补”底层发生了什么调优思路也自然清晰。7.2 用 Python 操作 HadoopHadoop Streaming 与 PyHive最热词里有“python与hadoop数据sql脚本的使用”说明很多同学希望用 Python 而不是 Java 来写 MR。Hadoop 官方提供了一种方式叫Hadoop Streaming它允许你用任意可执行程序比如 Python 脚本充当 Mapper 和 Reducer只要它从标准输入读数据、向标准输出写数据即可。下面是一个非常简单的 Python Streaming WordCount#!/usr/bin/env python3 import sys for line in sys.stdin: words line.strip().split() for word in words: print(f{word}\t1)这是 Mapper 脚本。Reducer 脚本可以从 stdin 读取 Map 输出的 KV 对按 Key 聚合#!/usr/bin/env python3 import sys current_word None current_count 0 for line in sys.stdin: word, count line.strip().split(\t, 1) count int(count) if current_word word: current_count count else: if current_word: print(f{current_word}\t{current_count}) current_word word current_count count if current_word: print(f{current_word}\t{current_count})提交命令hadoop jar $HADOOP_HOME/share/hadoop/tools/lib/hadoop-streaming-3.3.4.jar \ -input /user/example/input \ -output /user/example/output_py \ -mapper mapper.py \ -reducer reducer.py \ -file mapper.py \ -file reducer.py这里有个细节-file参数会把脚本从本地分发到集群的各个节点上。如果你的脚本需要依赖其他模块或配置文件也要一并-file上传。Hadoop Streaming 的优点是灵活但性能会比 Java 版本差一些主要原因是 Python 脚本每次处理一行都要经过进程内解析且标准输入输出的开销比二进制序列化大。对于离线报表这类非实时任务这个性能损失完全可以接受。如果你追求更高效的 Python 方案可以考虑 MRJob、PySpark 或者直接通过 PyHive 写 HQL但这些属于后话。7.3 Spark 与 MapReduce 的互补关系写到这里肯定会有人问既然 Spark 这么快我是不是可以直接跳过 MapReduce我的观点恰恰相反MapReduce 的价值不在于“算得快”而在于“算得明白”。Spark 的 RDD 流程、Stage 划分、Shuffle 机制很多核心概念都继承自 MapReduce。你只有先理解 MapReduce 的 Map → Shuffle → Reduce 模型再去学 Spark 的宽依赖、窄依赖、Shuffle 调优时才会有“原来如此”的感觉。MapReduce 和 Spark 的取舍也很直接MapReduce适合超大规模数据的离线批处理集群规模大、稳定性要求极高且资源允许场景。它对内存的要求低可以把中间数据放磁盘所以不太容易因为数据量暴涨而 OOM。Spark更擅长迭代计算、交互式查询和实时性要求稍高的场景因为中间结果尽量驻留内存避免反复落盘。生产环境中很多团队是 Hadoop Spark 混用日常 ETL 用 Hive SQL 跑在 MapReduce 引擎上或者直接切到 Spark 引擎数据挖掘类的迭代算法用 Spark MLlib小规模即席查询用 Hive on Tez 或 Presto。你懂 MapReduce再学这些新引擎会省力很多。8. 实战心得几个值得长期记住的经验判断最后分享几条我在实战中反复验证过的经验希望能帮你少走弯路。第一先确认“到底该不该用 MapReduce”再动手。如果你的数据量只有几个 GB单机数据库或者 pandas 处理可能更合适如果要做实时流处理MapReduce 也不是好的选择。MapReduce 的定位就是“海量离线数据的批量处理”判断标准很简单数据量是否远超单机处理能力对延迟是否不太敏感第二调优的顺序应该是“先压任务再调参数”。很多新手一上来就改mapreduce.map.memory.mb、mapreduce.reduce.memory.mb但忽略了代码本身的效率。我见过有人 map() 里对每个数据都创建了一个java.util.regex.Pattern对象结果导致 GC 非常频繁。先把代码层面的问题解决掉再谈参数调优才是正确顺序。第三学会读日志和看监控是基本功。遇到问题不要急着猜先看 YARN 的 Application 日志再到对应节点查看 Container 日志通常能很快定位。YARN 的 ResourceManager Web UI 会展示每个任务的状态和日志链接这是排查问题的第一入口。yarn logs -applicationId app_id这个命令可以捞回所有 Container 的日志在排查分布式任务问题时非常好用。第四MapReduce 学完一定要自己写一个完整的作业亲手跑一遍。无论是课程设计、面试准备还是工作中首次接触没有比“跑通一个真实场景案例”更快的成长方式。你可以找一份访问日志统计 PV/UV、Top URL、流量分布或者找一份订单表做用户维度分析。做完这些你才会真正理解为什么 Hadoop 生态能够成为大数据时代的基石。第五一个小技巧。在开发阶段如果要反复运行同一个作业可以把输出目录改成 UUID 或时间戳命名防止“输出目录已存在”的报错耽误时间。也可以把输入范围缩小到几条数据先验证逻辑正确再放开到全量数据这样迭代速度会快很多。