尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

从MapReduce原理到实战:核心机制、Shuffle与排序应用详解

从MapReduce原理到实战:核心机制、Shuffle与排序应用详解 MapReduce这个名词在Hadoop生态里存在感极强但真正能把它讲透的人不多。很多人学了WordCount就觉得掌握了MapReduce结果一到面试问Shuffle细节、自定义排序、数据倾斜就卡壳。这篇文章我想把多年实战里对MapReduce的理解系统地梳理一遍——从它解决的分布式计算本质问题到作业提交后的完整生命周期再到排序、倒排索引这类高频实战场景一次性把理论和代码串起来。适合正在学Hadoop、准备大数据面试、或者工作中需要写MapReduce但只停留在调包阶段的同学。我会尽量用大白话拆解原理给出可以直接复现的代码和命令让这篇文章真正能帮你少走弯路。1. 先把MapReduce读懂它到底在解决什么问题1.1 大数据处理的三个本质困境我经常跟初学者说MapReduce不是一门编程语言而是一种解决问题的思维框架。要理解它得先回答一个问题为什么单机处理不好大数据第一数据量超出单机承载能力。当一份数据有几十TB甚至PB级一台服务器的磁盘根本放不下内存更是完全装不下。就算勉强塞进去单机CPU的计算能力也让处理时间变成了天文数字。第二分布式计算引入了全新复杂性。既然一台机器搞不定那就多台机器一起上。但数据分散在多台机器上谁负责哪一块机器之间怎么通信某台机器突然挂了怎么办如何保证结果正确这些问题如果让每个程序员自己处理写出来的程序大概率是bug满天飞。第三传统编程模型对分布式不友好。你写一个单机程序从头到尾顺序执行逻辑很清晰。但放到分布式环境里数据是分片的、计算是并行的、节点是会宕机的如果还按照传统方式写代码会极度混乱。MapReduce的价值就在这它把这套分布式计算的复杂度封装成两个接口——map和reduce。你只需要告诉它怎么处理一条数据和怎么合并同类的数据,剩下的切分、调度、容错、通信全部由框架接管。这就是设计者当年提出分而治之思想的初衷。1.2 移动计算比移动数据更划算这是MapReduce最核心的设计哲学也是新手最容易忽略的点。假设你有1000台机器每台机器上存了1GB数据你想对这批数据做统计。最直觉的做法是把所有数据通过网络拉到一台中心服务器上处理。但数据总量是1TB网络传输就成了瓶颈——哪怕千兆网卡传输也需要几十分钟甚至更久而且中心服务器也成为单点瓶颈。MapReduce的做法完全反过来把计算逻辑jar包、代码分发到每一台存有数据的机器上让每台机器只计算自己本地的数据最后再把各自的结果聚合。这个思想叫做数据本地性Data Locality本质就是用计算移动代替数据移动。我举一个更生活化的例子你要统计一个小区里每栋楼有多少人你不会把所有人叫到物业办公室数一遍而是派一个统计员去每栋楼门前数人最后把各楼数字汇总就行。统计员就是map任务汇总就是reduce任务数据不动人计算在动。在Hadoop 1.x时代任务调度器会优先把map任务调度到数据所在节点上执行到YARN时代虽然调度粒度变成了容器Container但本地性优化依然是重要考量。理解这一点你就明白了为什么MapReduce适合批处理、不适合低延迟查询——因为懒加载和本地性优化本身就是为算大数设计的。1.3 Mapper、Reducer与Driver各自扮演什么角色任何MapReduce作业都跑不了三个角色我用最直白的话来定义它们Mapper映射器接收一条记录输出零到多条键值对。它做的是细粒度的筛选、清洗、转换。比如从一行日志中提取IP和时间或者把一行文本按空格切成单词。Reducer归约器接收某个key对应的所有value集合对这些value做聚合输出最终结果。它是粗粒度的汇总、合并、计算。比如统计所有IP出现的总次数。Driver驱动类配置作业参数并提交整个作业。它不参与计算更像是施工现场的总监工负责告诉Hadoop我要跑什么、跑在哪里、怎么跑。我在实际写项目的时候会把业务逻辑拆成一句话map阶段解决把数据变成什么reduce阶段解决把同类数据怎么合。比如统计流量日志中每个用户的访问量map把每行日志提取出用户ID并输出(userId, 1)reduce把同一个userId的所有1累加起来输出(userId, totalCount)。思路清晰了代码写起来就顺了。2. 从提交到落盘一个MapReduce作业的完整生命周期这一章节是面试重灾区也是排查问题必须掌握的知识。很多人在集群上跑作业报了错一脸懵地问我为什么结果发现是对作业执行流程没有整体概念。我按时间线从作业提交讲起。2.1 InputFormat与Split数据到底怎么切作业提交到集群后计算框架做的第一件事就是决定哪些数据被谁处理。这个决策由InputFormat接口负责默认实现是FileInputFormat。FileInputFormat会把输入目录下的所有文件逻辑上切成若干个输入分片InputSplit。分片不是物理上真的把文件切块而是一个逻辑概念它记录了要处理哪一段数据的起始偏移量和长度。默认情况下每个分片的大小接近HDFS block大小默认128MB这样分片才能大概率落在某个block所在节点上从而实现数据本地性。这里有个容易混淆的点HDFS block是存储层面的物理分块InputSplit是计算层面的逻辑分片。一个split可以由多个block组成如果文件被压缩且不可切分一个block也可以被多个split引用理论上。InputFormat的另一个职责是提供RecordReader它负责把一个split里的数据一行一行读出来转换成(key, value)对。对于文本文件默认的LineRecordReader以每行文本为一条记录key是行首字节偏移量LongWritablevalue是行内容Text。这就是为什么你写WordCount的map方法接收的参数永远是(LongWritable key, Text value)。2.2 Mapper与Shuffle框架里最复杂的环节map函数对每条记录处理后输出(key, value)。接下来这条输出并不会直接送到reduce端而是进入一个被称为Shuffle的中间环节。Shuffle是MapReduce中最难啃、也最值得深挖的部分理解它就是理解MapReduce性能调优的钥匙。Map端的Shuffle过程大致包含这几步分区Partition每条map输出都要经Partitioner决定去哪个reduce任务。默认是HashPartitioner它计算key的哈希值并模上reduce任务数保证相同key进入同一个reduce。这也是所有value能聚到一起的前提。环形缓冲区与Spillmap输出的数据不是直接存磁盘而是先写入一个环形内存缓冲区默认100MB。当缓冲区写满80%时后台线程开始把数据溢写到本地磁盘这个过程叫spill。溢写前会做两件事按key排序和合并相同key的value可选combiner。Sort与Merge一次map任务可能会产生多个spill文件最终map任务结束前这些spill文件会合并成一个更大的有序文件。合并时按key排序、按分区排列最终形成一个分区内有序、整体按分区顺序排放的文件。到Reduce端后每个reduce任务从各个map任务拉取属于自己分区的数据这个动作叫Fetch或Copy。数据先放到reduce端内存缓冲区再经历一次合并排序最终形成一个有序的、完整的(key, list[value])数据集喂给reduce函数。我把这个流程总结成一句话map输出——分区——排序——溢写——合并——reduce拉取——再排序——分组喂给reduce。顺便说一句如果你在map和reduce之间配置了Combiner它在spill时和merge时会对相同key的value做一次局部聚合。注意Combiner使用的函数必须是可重复执行的不能影响最终结果。最典型的例子是求和可以在combiner做但求平均值不能随便在combiner做。2.3 Reducer与OutputFormat结果如何落盘reduce函数接收(key, value列表)经过业务逻辑处理后输出最终结果。它输出的(key, value)对交由OutputFormat写入目标存储——最常见的是FileOutputFormat即写入HDFS指定目录下每个reduce任务生成一个part-r-xxxxx文件。这里有个重要细节最终生成的part文件个数等于reduce任务个数。所以当你发现输出目录里part-r-00000只有一个文件时说明实际只有1个reduce任务在跑。这往往不是好事可能是你显式设置了setNumReduceTasks(1)或者key类型导致单个reduce函数内数据量极大。在数据量大的场景下这会引起严重的数据倾斜和性能瓶颈。另外一个比较容易忽略的点是map-only作业没有reduce阶段的输出文件命名是part-m-xxxxx且文件个数等于map任务个数。这类作业常用于数据清洗、格式转换、ETL场景因为不需要聚合跑起来速度很快。我对初学者的建议是先把整个作业的生命周期画成一张输入分片→map→spill→merge→copy→merge→reduce→输出的行动路线图每写一个作业就对着这张图想一遍我的数据现在在哪一步。这样排查问题会有方向得多。3. 环境搭建与第一个WordCount程序3.1 学习阶段该选伪分布式还是集群很多同学一开始就照着网上的教程搭三节点甚至五节点集群结果环境没配好光折腾YARN的ResourceManager和NodeManager通信就花了两三天。我个人的经验是学习MapReduce阶段伪分布式Pseudo-Distributed就完全够了。伪分布式的意思是在一台机器上同时运行HDFS的NameNode和DataNode、YARN的ResourceManager和NodeManager所有守护进程都在但都在本机。它模拟了真实分布式环境的完整流程作业提交、任务调度、Shuffle全部都会真实发生只是都在一台机器上完成因此足够用来验证思路和学原理。如果一上来就搭集群你会面临几个麻烦机器资源不够导致任务跑不动节点间SSH配置错误导致通信失败防火墙问题导致端口不可达。而我实际做项目时的经验是先在伪分布式上把逻辑跑通再迁移到集群排错效率和自信心都会好很多。3.2 Hadoop配置要点与启动验证以Hadoop 3.x为例伪分布式需要重点修改以下四个配置文件core-site.xml配置默认文件系统为HDFS例如fs.defaultFS设为hdfs://localhost:9000。hdfs-site.xml设置副本数为1伪分布式只有一个DataNode默认3副本会卡在副本状态并可以设置NameNode和DataNode的目录位置。yarn-site.xml开启YARN的辅助服务配置ResourceManager地址单机环境还需要设置yarn.nodemanager.aux-services为mapreduce_shuffle。mapred-site.xml指定MapReduce运行在YARN之上设置mapreduce.framework.name为yarn。配置完成后格式化NameNodehdfs namenode -format然后执行start-dfs.sh和start-yarn.sh启动服务再用jps命令检查进程。如果看到NameNode、DataNode、ResourceManager、NodeManager几个进程都在环境就基本就绪了。我建议执行一次最简单的验证用hdfs dfs -put上传一个本地文件到HDFS再执行hdfs dfs -cat读回来确认HDFS读写正常。这一步能帮你区分HDFS有问题和MapReduce有问题长期来看很省事。3.3 WordCount代码逐行拆解与运行命令WordCount是MapReduce的Hello World但对于刚接触的人来说它的代码里藏着不少值得注意的细节。我直接贴一个标准实现附上注释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; public class WordCount { // Mapper把一行文本拆成单词输出 (单词, 1) public static class TokenizerMapper extends MapperLongWritable, Text, Text, IntWritable { private final static IntWritable one new IntWritable(1); private Text word new Text(); Override protected void map(LongWritable 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对每个单词的 value 列表求和 public static class IntSumReducer 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); } } 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); } }运行步骤很简单用Maven或直接javac编译打包成wordcount.jar然后在集群上执行hdfs dfs -mkdir -p /input hdfs dfs -put test.txt /input/ hadoop jar wordcount.jar WordCount /input/test.txt /output/wc hdfs dfs -cat /output/wc/part-r-00000几个细节我再啰嗦一下job.setCombinerClass这里直接复用了Reducer类因为求和这种操作满足交换率和结合律可以在map端先做一遍局部求和。这会显著减少map端到reduce端的网络传输量。代码中job.setJarByClass这句不能省。没有它分布式环境下YARN找不到jar包会报ClassNotFound错误。输出目录必须不存在。Hadoop出于安全考虑防止覆盖结果如果输出目录已存在会直接报错。我经常看到有人反复跑同一个作业报错就是这个原因。4. 排序实战从自定义排序到分组排序4.1 为什么默认排序不满足业务需求MapReduce框架对key的排序是天然行为——在Shuffle过程中map端和reduce端都会按key排序。但默认的排序只对WritableComparable类型生效并且只能按key的自然顺序排。比如IntWritable按数值升序Text按字典序升序。实际业务中这远远不够。举个例子你要统计每个用户的总消费金额并按照金额从高到低输出。默认排序是按用户ID字典序排的你拿到的reduce输出是ID有序而不是金额有序。另一个高频场景是二次排序按年份升序排列同一年份内部按温度降序排列并输出每年最高温度。这种复合排序必须自定义key。我总结一下凡是排序字段不是key本身、或者排序规则不是自然序、或者需要用多个字段联合排序的场景都需要自定义WritableComparable。4.2 自定义WritableComparable的完整实现实现自定义key的标准做法是继承WritableComparableT接口覆写write、readFields、compareTo三个方法。下面我以流量日志按上行流量排序为例。输入格式是手机号 上行流量 下行流量需求是按下行流量降序输出手机号及其流量。java import org.apache.hadoop.io.WritableComparable; import java.io.DataInput; import java.io.DataOutput; import java.io.IOException; // 自定义key以手机号和下行流量为组合维度 public class FlowBean implements WritableComparableFlowBean { private String phone; // 手机号 private long downFlow; // 下行流量 public FlowBean() { super(); } public FlowBean(String phone, long downFlow) { this.phone phone; this.downFlow downFlow; } Override public void write(DataOutput out) throws IOException { out.writeUTF(phone); out.writeLong(downFlow); } Override public void readFields(DataInput in) throws IOException { this.phone in.readUTF(); this.downFlow in.readLong(); } // 核心逻辑按下行流量降序相同则按手机号升序 Override public int compareTo(FlowBean o) { if (this.downFlow o.downFlow) { return -1; } else if (this.downFlow o.downFlow) { return 1; } else { return this.phone.compareTo(o.phone); } } Override public String toString() { return phone \t downFlow; } // 注意必须重写equals和hashCodePartitioner会用到 Override public boolean equals(Object obj) { if (this obj) return true; if (obj null || getClass() ! obj.getClass()) return false; FlowBean that (FlowBean) obj; return downFlow that.downFlow phone.equals(that.phone); } Override public int hashCode() { return phone.hashCode() (int) (downFlow ^ (downFlow 32)); } }这个key在map阶段作为输出key框架会自动按compareTo定义好的顺序排序。如果你的需求还要让同一个reduce处理同一组key的数据比如所有手机号都到同一个reducer做汇总那你需要控制分区——不过在这个例子里更重要的是理解compareTo决定了全局排序的顺序。我特别提醒一点write和readFields的字段顺序必须完全一致。序列化时先写phone再写downFlow反序列化时也必须先读phone再读downFlow。顺序错了读出来的数据就是错乱的而且报错往往不是你想象中的类型转换错误而是数据语义错乱排查起来很费劲。4.3 二次排序与GroupingComparator实战自定义key的另一种典型应用是二次排序。我现在把最经典的年份-温度案例完整跑一遍。需求有一批气象数据每行是年份 温度要求按年份升序同一年份内按温度降序排列并在reduce端输出每年最高温度。第一步定义组合keyDateTemperature包含year和temperature两个字段。public class DateTemperature implements WritableComparableDateTemperature { private int year; private int temperature; // 构造函数、write、readFields 略与上面类似 // 二次排序的核心年份升序年份相同则温度降序 Override public int compareTo(DateTemperature o) { int yearCompare Integer.compare(this.year, o.year); if (yearCompare ! 0) { return yearCompare; } return Integer.compare(o.temperature, this.temperature); } }第二步Mapper输出(DateTemperature, NullWritable)输入原样输出即可。第三步如果不做任何额外设置reduce端会收到一组按年份升序、年份内温度降序排好的数据对。此时每年第一个value就是该年最高温度。但有个问题每个年份有好多条数据你无法直接在reduce里只取每年第一个value然后调用一次reduce方法因为框架会把所有不同key的数据分别调用reduce。这时候就需要自定义分组比较器GroupingComparator。它的作用是指定哪些key算同一组应该汇聚到一个reduce方法里。默认情况下分组是完全按key的compareTo结果来的即每个唯一key一组。我们要让同年份的key归为一组就需要告诉框架只要年份相同就是同一组。import org.apache.hadoop.io.WritableComparable; import org.apache.hadoop.io.WritableComparator; // 分组比较器只看年份忽略温度 public class YearGroupingComparator extends WritableComparator { public YearGroupingComparator() { super(DateTemperature.class, true); } Override SuppressWarnings(rawtypes) public int compare(WritableComparable a, WritableComparable b) { DateTemperature dt1 (DateTemperature) a; DateTemperature dt2 (DateTemperature) b; return Integer.compare(dt1.getYear(), dt2.getYear()); } }第四步在Driver中注册分组比较器并在Reducer中取出组内第一个value。job.setGroupingComparatorClass(YearGroupingComparator.class);Reducer的reduce方法中因为每年数据已经按温度降序排列组内第一项就是最高温Override protected void reduce(DateTemperature key, IterableNullWritable values, Context context) { // values里全是NullWritable真正有用的数据在key中 // 而且因为组内第一个key就是温度最高的那一条 // 直接取第一次迭代的key即可 for (NullWritable val : values) { context.write(key, NullWritable.get()); break; // 只取第一条 } }这个案例能让你彻底分清三个排序层级key的compareTo决定reduce端收到的key顺序全局排序。自定义分区器决定哪些key进入哪个reduce任务跨reducer分配。分组比较器决定哪些key会被合并为同一个reduce方法的调用组内聚合粒度。我用一个生活类比帮助理解假设你要整理全班同学的考试成绩compareTo决定按分数从高到低排座位分区器决定哪些同学分在哪个考场分组器决定哪些同学属于同一个班级成绩单打印在一页纸上。5. 倒排序索引组合key的典型应用5.1 从搜索引擎场景理解倒排索引倒排序索引Inverted Index这个名词听起来高大上其实就是搜索引擎里最基础的数据结构。我给你个直观例子翻开一本书最后的索引页它列出的是关键词 → 页码让你能快速找到关键词出现在哪些页。这个从词到文档位置的映射就是倒排索引。在MapReduce中做倒排索引经典需求是给定一批文档输出每个单词出现在哪些文档中、各出现了几次输出格式一般是word - doc1:count1, doc2:count2。这个场景非常适合撑握组合key的用法因为它需要在一个key里同时包含单词和文档ID。5.2 实现思路与代码实战实现倒排索引一个最直接的方法是把文档ID和单词拼接成一个组合key比如word#docId然后reduce阶段对同一组合key的value求和最后输出。但这样有个问题Reducer的输出里单词和文档是混在一个key里的无法满足同一个单词的所有文档归并在一起的目标。所以我在实战中倾向于用两阶段MapReduce来实现第一阶段把(word, docId)作为组合key求词频第二阶段以word作为key把属于同一个word的所有docId:count拼接到一起。不过单阶段也可以实现做法是map阶段读入(文件名, 行内容)对每个单词输出(word#文件名, 1)在reduce阶段手动按#拆开先聚合到(word, 文件名) → 词频再做一次value拼接。因为reduce接收的key已经按字典序排好同一word的所有word#docId会自动相邻只要遍历一遍即可合并。我给出单阶段的核心代码片段public static class InvertedIndexMapper extends MapperLongWritable, Text, Text, Text { private Text outKey new Text(); private Text outValue new Text(); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { // 从InputSplit中获取文件名作为文档ID FileSplit fileSplit (FileSplit) context.getInputSplit(); String fileName fileSplit.getPath().getName(); String line value.toString(); StringTokenizer itr new StringTokenizer(line); while (itr.hasMoreTokens()) { String word itr.nextToken(); // 组合key: word#docId outKey.set(word # fileName); outValue.set(1); context.write(outKey, outValue); } } } public static class InvertedIndexReducer extends ReducerText, Text, Text, Text { private Text result new Text(); Override protected void reduce(Text key, IterableText values, Context context) throws IOException, InterruptedException { String[] parts key.toString().split(#); String word parts[0]; String docId parts[1]; int sum 0; for (Text val : values) { sum Integer.parseInt(val.toString()); } // 注意这里不能直接context.write因为我们要把同一word的结果拼接 // 但reduce是逐key调用的因此word#docId的不同组合会分别调用 // 所以拼接逻辑可以放在Driver的cleanup阶段配合TreeMap或者这里先输出 // 单阶段简化版先输出 (word, docId:count) result.set(docId : sum); context.write(new Text(word), result); } }不过上面的单阶段输出其实会产生(word, docId:count)这样的多行内容同一word的多行并不会自动合并。如果希望输出真正满足一个word对应一串文档统计更规范的做法是在Reducer内部用一个LinkedHashMap或TreeMap收集当前word的所有(docId, count)在cleanup阶段统一拼接输出。这要求你把同一个word的所有数据都交给同一个Reducer——由于组合key的word部分在字典序上是相邻的默认哈希分区很难保证这一点。为了保证同一个word全部落到同一个Reducer需要自定义Partitioner只对组合key的word部分做哈希。这也是我在面试时非常喜欢考察的一个点自定义Partitioner配合组合key才能实现逻辑分组的精准控制。如果你能把这个链路讲清楚面试官就知道你是真动过手的人而不只是背了概念。6. 面试高频问题与实战调优心得6.1 高频面试题背后的原理考察我发现MapReduce的面试题翻来覆去就是那几个但很多人只会背答案。下面我把这些题目背后真正要考的原理点指出来Shuffle阶段发生了什么这是必考题。考察点不是你能不能背出分区、排序、溢写、合并、拉取、再合并这几个词而是你是否理解为什么需要这些步骤。我在前文第二节已经详细拆解过这就是教科书级别的答案。MapReduce为什么不适合实时查询考察点在于对批处理本质的认识。作业启动有开销、Shuffle有大量磁盘I/O、调度有延迟这些都决定了它面向分钟级响应而不是毫秒级。数据倾斜怎么解决考察点在于是否能定位倾斜发生在map端还是reduce端。常见解法包括对热key加随机前缀打散、调整Partitioner、增加或减少reduce任务数、使用Combiner做局部聚合。如何实现自定义排序考察点就是第四章讲的自定义WritableComparable和GroupingComparator。建议把代码背下来因为面试官很可能会让你现场画类图和写关键方法。我个人的感受是单纯背题目容易真正理解需要把代码跑一遍观察不同的setNumReduceTasks、不同的Partitioner对输出文件数量和数据分布的影响。只有亲手做过实验面试被追问细节时才不会慌。6.2 实际集群中的常见故障与调优项目里跑久了你会发现真正折磨人的不是业务逻辑而是作业跑得慢、报错看不懂。我把高频出现的两类问题列一下第一类是小文件过多导致大量map任务。MapReduce的map任务数取决于输入分片数量如果上游产生了成千上万个几KB的小文件每个文件都会被当成一个split导致启动几百个map任务光任务调度开销就把作业拖垮了。解决办法通常是先做一次小文件合并或者在数据写入阶段就避免产生过多小文件。第二类是内存参数设置不当导致任务被杀。YARN中MapReduce的容器内存涉及mapreduce.map.memory.mb、mapreduce.reduce.memory.mb、mapreduce.map.java.opts等参数。曾经我在生产环境遇到reduce任务频繁被杀Container killed by the ApplicationMaster最后排查发现是堆内存设置超过了容器内存上限java.opts配的堆大小比.memory.mb还大。调整思路是让-Xmx保持在容器内存的75%—80%左右留出框架自身的内存余量。第三类是Reducer数量拍脑袋设。我见过有人把setNumReduceTasks(1000)当性能优化手段结果大量小reduce任务反而拖慢整体。合理的reducer数量需要综合考虑数据量和每个reducer期望处理的数据大小一般以每个reducer处理1GB左右数据为参照。更稳妥的做法是先按默认跑一版测试观察各reduce处理的数据分布情况再针对性调整。6.3 我踩过的坑和留下的习惯关于MapReduce我最后分享几个自己的习惯这些都是在踩坑后总结出来的第一凡是涉及自定义key必然重写equals和hashCode。这是血的教训。有一个线上作业自定义对象没重写hashCode导致相同的key被分到了不同的reduce任务单看每个reduce的结果都正确合起来却发现统计数翻了好几倍。排查了整整半天原因就是HashPartitioner用了key的hashCode来分区。第二不要在reduce里维护大量全局状态。reduce方法本身是逐key调用的但每个reduce任务实例的生命周期里所有key共享同一个Reducer对象。如果你在类里定义了一个HashMap用于暂存数据在数据量大时会内存暴涨。如果要跨key聚合务必评估清楚内存上限或者考虑改用两阶段作业。第三先小数据验证再大数据跑。哪怕逻辑已经很熟了我也习惯先用几十行的小数据跑一版验证输出格式和结果正确性再换全量数据。Hadoop作业一旦跑起来失败重试的时间成本太高小数据验证能帮你节省大量时间。这些经验不一定能直接写进教科书但确实是项目里最实用的东西。MapReduce这个框架虽然如今在不少场景被Spark、Flink等新引擎取代但理解它的核心思想、Shuffle细节、排序机制仍然是深入大数据领域的基石。以后你学Spark的Shuffle、Flink的状态管理时会发现很多设计思路是一脉相承的。把MapReduce吃透了那些看似过时的技术其实都在滋养你后续学习新框架的能力。
返回列表