ARTICLE DETAIL

资讯详情

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

MapReduce原理与实战:从WordCount到二次排序、分区与数据倾斜调优

MapReduce原理与实战:从WordCount到二次排序、分区与数据倾斜调优 1. 初识 MapReduce它到底是什么又是为了解决什么问题而生的先用一句话把这件事说清楚MapReduce 是 Google 在 2004 年发表的一篇论文中提出的分布式计算编程模型后来 Hadoop 把它变成了大规模数据处理的事实标准。它的核心思想就八个字——“分而治之并行计算”。你不需要懂分布式系统的底层细节只要按照它规定的 Map映射和 Reduce归约两个阶段去写业务逻辑系统就会自动把一个大任务拆成无数个小任务丢到一群机器上同时跑最后再把结果汇总回来。我第一次接触 MapReduce 的时候其实是在做一份日志统计的需求。当时数据量大概有几百 GB分布在几十台服务器上如果用传统单机程序逐行读取估计要跑到第二天早上。后来我把处理逻辑改成 MapReduce 之后同一份数据跑完只用了二十多分钟。那个对比带来的震撼直到今天我还记得。所以这篇文章我不会光讲理论我会把你从“听说过 MapReduce”一路带到“能自己写 Map 和 Reduce、能理解 Shuffle 发生了什么、能排查常见的作业失败问题”的程度。这篇文章适合谁看两类人。一类是刚学大数据、正在做实训作业的学生尤其是那些卡在“头歌 MapReduce 排序”“自定义分组”“倒排序索引”这类题目上的同学另一类是在工作中需要处理海量数据但还没有系统性梳理过 MapReduce 内在逻辑的工程师。读完之后你不仅能把作业交上去还能真正理解每一行代码背后的运行机制。有人说 MapReduce 已经过时了现在大家都在用 Spark、Flink。这话对了一半。MapReduce 的实时性和迭代计算能力确实不如 Spark但它的模型思想是所有分布式计算框架的基石。你理解了一个 Key-Value 对如何从 Map 端流向 Reduce 端再去学 Spark 的 RDD、Flink 的 DataStream会发现它们的内核几乎一脉相承。所以我始终认为MapReduce 不是过时了而是沉淀成了你理解整个大数据生态的必修课。2. 核心原理拆解从数据流的角度看懂 MapReduce 的一生2.1 一个完整 Job 的五阶段旅程官方文档常常把 MapReduce 的流程画成一张很复杂的图有 InputFormat、Split、RecordReader、Mapper、Combiner、Partitioner、Shuffle、Sort、Reducer、OutputFormat 等一大堆名词。初次接触的人很容易被吓住。但如果我们把数据当成一条流水线上的工件整个流程其实就五个阶段。第一输入分片。HDFS 上的一个大文件会被切成若干个 Split每个 Split 对应一个 Map 任务。这里有个关键细节——Split 和 HDFS 的 Block 并不完全是一回事。默认情况下一个 Split 对应一个 Block通常是 128MB但你可以通过调整mapreduce.input.fileinputformat.split.minsize和maxsize来改变 Split 的大小。每个 Split 里并不是单独存放数据而是记录了这个分片对应 HDFS 上的哪一段数据以及起始位置和长度。第二Map 阶段。框架读取 Split 中的每一行数据交给 Mapper 的map()方法处理。map()接收一个 Key 和一个 Value输出若干个新的 Key-Value 对。这里要注意Map 阶段输出的数据是暂时存放在内存缓冲区里的不是直接写到磁盘。缓冲区的默认大小是 100MB当写入量达到 80%也就是 80MB 时后台线程就会开始把数据溢写到本地磁盘。这个阈值由mapreduce.task.io.sort.mb和mapreduce.map.sort.spill.percent控制实际调优时经常会动这两个参数。第三Shuffle 阶段。这是 MapReduce 里最复杂、也最容易被误解的环节。Map 输出的 Key-Value 对会根据 Key 的哈希值被分到不同的分区每个分区对应一个 Reduce 任务。分区内再按 Key 排序如果有 Combiner 的话还会在 Map 端先做一次局部合并减少网络传输的数据量。之后Map 端的输出就会被 Reduce 端拉取——注意Reduce 并不是等所有 Map 都跑完才开始拉数据只要某个 Map 完成了Reduce 就会去拷贝那个 Map 的输出。这个过程是异步的、并行的。第四Reduce 阶段。Reduce 拿到的数据已经是被分区、排序过的。同一个 Key 的所有 Value 会被组合成一个迭代器交给reduce()方法处理。你在reduce()里写的逻辑就是对一组相同 Key 的 Value 做聚合计算。框架保证传给同一个 Reduce 的 Key 是有序的但不保证不同 Reduce 之间的顺序。第五输出阶段。Reduce 的结果通过 OutputFormat 写入 HDFS。默认的 TextOutputFormat 会把每对 Key-Value 写成一行的文本Key 和 Value 之间用 Tab 分隔。从整体上看MapReduce 的数据流就是 Input - Map - Shuffle - Reduce - Output。理解这条主线之后你再看那些复杂的参数和源码就只是往主线上挂细节而已。2.2 深入 Shuffle 内部排序、分区、合并到底在干嘛Shuffle 是 MapReduce 里最容易出问题、也最能体现工程师水平的地方。网上很多文章把 Shuffle 说得玄乎我换个方式用洗衣房的流程来做类比Map 端就像每个人把脏衣服丢进洗衣机前先自己按颜色粗分一遍分区然后按照色号深浅排好序排序相同色系的衣服装进同一个袋子合并这些袋子再集中送到对应的分拣台Reduce 端。分拣台收到各个来源的袋子后再把同一色系的衣服归拢在一起按深浅排好归并排序最后交给专人处理。具体到代码层面Shuffle 中最重要的两个组件是 Partitioner 和 Comparator。Partitioner 决定了一个 Key 去哪个 Reduce。默认的实现HashPartitioner只是对 Key 的哈希值取模 Reduce 数量所以你会发现同一个 Key 一定会进同一个 Reduce。Comparator 则负责排序规则。默认情况下框架按照 Key 的自然顺序排序也就是字符串按字典序、数值按大小。但有时候我们需要自定义排序规则比如让数值大的排在前面或者让两个字段联合排序这时候就要写自定义 Comparator或者让 Key 实现WritableComparable接口。Combiner 是另一个容易被低估的优化点。Combiner 本质上是一个运行在 Map 端的 Mini-Reducer它的作用是在 Map 端先把部分重复 Key 的 Value 合并掉减少要写到磁盘、通过网络传输的数据量。最常见的例子是求和在 WordCount 中如果某个单词在一个 Map 任务中出现了几百次没有 Combiner 的话这几百条记录都会原样写到磁盘再由 Reduce 拉取加了 Combiner 后Map 端先把这个单词的计数加总成一条记录再输出网络 I/O 和磁盘 I/O 立刻降下来。使用 Combiner 的前提是它的输入输出类型必须与 Mapper 的输出类型一致并且逻辑必须满足交换律和结合律。否则你可能会得到错误的结果——比如算平均值就不能直接在 Map 端做局部均值合并。2.3 数据倾斜为什么某个 Reduce 总是跑得特别慢这是我在实际项目里踩过最深的坑。所谓数据倾斜就是大量相同 Key 的数据全部涌向同一个 Reduce导致其他 Reduce 早就跑完了唯独那一个 Reduce 还在忙碌。表现出来的现象是整个 Job 的运行时间被一个任务拖长集群利用率很低甚至会因为内存溢出而失败。倾斜的原因有很多。最常见的是 Key 的分布本身就不均匀比如用户日志中某个热门用户 ID 占据了 80% 的数据。这种情况下用默认的 HashPartitioner 必然导致那个 ID 对应的 Reduce 成为瓶颈。处理思路也分两条线。一条是从数据源头入手做加盐Salting处理给倾斜的 Key 加上随机前缀把它拆散成多个子 Key让它们分布到不同 Reduce处理后再把前缀去掉做二次聚合。另一条是调整 Combiner 和 Reduce 的合并逻辑让数据在 Map 端被压缩得更彻底减少倾斜 Reduce 端的数据量。还有一个比较隐蔽的倾斜源是自定义分组时使用了不合理的 Key 结构。比如我们要做“分组排序”业务把某个字段作为分组键另一个字段作为排序键。很多人会把这两个字段拼成一个复合 Key然后只重写 GroupingComparator却忽略了分区和排序也要和这个复合 Key 配套。结果可能造成本应分到同一个 Reduce 的数据被哈希到不同分区分组时根本不在一起。这个问题在实训题里特别常见下面我会专门展开讲。3. 从零实现一个 MapReduce 程序以 WordCount 为例的手把手教学3.1 项目准备与环境搭建写 MapReduce 程序之前先把环境准备好。我建议你用 Maven 管理依赖这样不会为 jar 包冲突操心。在pom.xml里加上 Hadoop Client 的依赖版本要和你运行的集群版本保持一致否则可能出现序列化协议对不上的问题。如果只是本地开发测试可以用 2.10.x 或 3.3.x 这类相对稳定的版本。dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version3.3.6/version /dependency然后准备一个输入文件。比如我建了一个input/word.txt里面放几行英文文本每行用空格分隔单词。这里的输入路径是指 HDFS 上的路径不是本地路径。你先把文件上传到 HDFShdfs dfs -mkdir -p /user/hadoop/input hdfs dfs -put word.txt /user/hadoop/input/ hdfs dfs -rm -r /user/hadoop/output # 输出目录不能存在需要先删掉旧的为什么输出目录不能存在因为 MapReduce 框架出于安全考虑如果检测到输出目录已经存在会直接报FileAlreadyExistsException。这是新手最容易踩的坑之一。3.2 Mapper、Reducer、Driver 三段式代码解析整个程序可以拆成三个类。第一个是 Mapper 类负责把一行文本拆成单词输出单词, 1。public class WordCountMapper extends MapperLongWritable, Text, Text, IntWritable { private Text word new Text(); private final static 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.length() 0) { continue; } word.set(w); context.write(word, one); } } }注意这里的输入 Key 是LongWritable它代表行在文件中的字节偏移量通常我们用不到但类型不能写错。Text是 Hadoop 自己封装的可序列化字符串类型对应 Java 的StringIntWritable对应Integer。这些 Writable 类型都是实现了 Hadoop 序列化接口的MapReduce 传输数据时靠它们完成二进制转换。第二个是 Reducer 类负责把同一个单词的所有计数累加。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()方法里传入的values是一个迭代器你不能把它整体保存下来等下轮再遍历因为框架为了节省内存迭代器底层复用的是同一个 Value 对象。如果你真想保存所有 Value必须先 deep copy否则得到的全是最后一个值。这个坑我见过很多次代码里“看起来”没问题运行结果却总是错的。第三个是 Driver 类也就是 main 方法所在的地方负责组装整个 Job。public class WordCountDriver { public static void main(String[] args) throws Exception { Configuration conf new Configuration(); Job job Job.getInstance(conf, word count); job.setJarByClass(WordCountDriver.class); job.setMapperClass(WordCountMapper.class); job.setCombinerClass(WordCountReducer.class); job.setReducerClass(WordCountReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); FileInputFormat.setInputPaths(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }值得注意的地方有四处。第一setJarByClass必须写而且要用包含 main 方法的类。这是为了让框架知道从哪里找到咱们的业务代码尤其是在集群模式下它会把类所在的 jar 分发到各个节点。第二我们这里把 Reducer 类同时设置为 Combiner 类因为单词计数这个聚合逻辑满足交换律和结合律Combiner 干了不会出错。第三setOutputKeyClass和setOutputValueClass设置的是 Mapper 和 Reducer 共同的输出类型。如果 Mapper 和 Reducer 的输出类型不一致就需要额外设置setMapOutputKeyClass和setMapOutputValueClass。第四waitForCompletion(true)里的参数表示是否打印进度信息生产环境建议设为 true方便观察任务状态。3.3 打包提交到集群并查看运行日志在 IDEA 或者命令行里用 Maven 打包mvn clean package -DskipTests打完包之后在 target 目录下会生成一个带依赖的 jar 或者普通 jar。如果用的是普通 jar并且依赖的 Hadoop 版本与集群一致直接提交即可hadoop jar target/wordcount-1.0.jar com.example.WordCountDriver /user/hadoop/input /user/hadoop/output提交之后控制台会滚动输出进度类似map 100% reduce 100%。跑完之后查看结果hdfs dfs -cat /user/hadoop/output/part-r-00000如果你看到的是hello 5、world 3这样的行说明你的第一个 MapReduce 程序已经成功跑通了。此时再回头看 2.1 节画出的那条数据流你应该会有一种“原来它们真的在按这个顺序执行”的感觉。4. 实战进阶排序、自定义分区与分组排序的实现套路4.1 全排序与二次排序先分清你要的是哪种网上搜“MapReduce 排序”会看到六七个相关热词什么“全排序”“二次排序”“分组排序”“自定义排序”。乍一看很乱我帮你理一理。全排序所有数据最终输出时严格按照某个 Key 的顺序排列。但需要牢记MapReduce 只在每个 Reduce 内部保证有序不保证 Reduce 之间的顺序。要做到全排序最简单的方案是只设置一个 Reduce代价是单点压力大、性能差更好的方案是自定义分区函数让数据分布到多个 Reduce 后还能保证前后 Reduce 的 Key 范围是连续的这需要你先对 Key 的分布有一个预估算划分出几个不重叠的范围。二次排序排序时不止看一个字段先按第一个字段排序第一个字段相同再按第二个字段排序。实现方式是把两个字段合并成一个复合 Key并编写分区器保证复合 Key 的第一个字段被分到同一个 Reduce编写排序比较器让复合 Key 先比较第一字段再比较第二字段编写分组比较器让相同第一字段的记录进入同一个 Reduce 方法。这套组合拳就是实训里最常见的“分组排序”题。自定义排序当默认的比较规则不满足需求时你让自定义类实现WritableComparable接口在compareTo()里写你自己的比较逻辑。比如求 Top N希望数值大的排在前面直接用 IntWritable 的话默认升序此时就需要自定义降序规则或者用LongWritable的负数技巧——当然不推荐负数这种 hack规范做法是实现接口。结合你在实训里可能遇到的题目第 1 关“MapReduce 排序—自定义排序”通常就是让你写一个自定义的 Bean 类实现WritableComparable按某列数据降序输出第 2 关“MapReduce 自定义分组”通常是在二次排序的基础上要求把第一字段相同的记录聚在一个 Reduce 里做聚合。下面我用一个完整的题目把这两个关卡串起来讲。4.2 经典实训题实战按订单 ID 分组按金额降序排序假设有一份订单数据每一行是订单ID 商品ID 金额例如A001 P001 200 A001 P002 100 A002 P001 300 A002 P002 150 A003 P001 50需求是按订单 ID 分组输出每个订单内部按金额从高到低排序并在每组后面输出该订单的总金额。这题你如果只用默认机制会发现同一个订单的几行数据不一定进同一个 Reduce——因为分区只依赖 Key 的哈希。所以我们需要构造复合 Key。第一步定义OrderBean类包含 orderId 和 amount 两个字段实现WritableComparableOrderBean。public class OrderBean implements WritableComparableOrderBean { private String orderId; private double amount; public OrderBean() {} public OrderBean(String orderId, double amount) { this.orderId orderId; this.amount amount; } Override public void write(DataOutput out) throws IOException { out.writeUTF(orderId); out.writeDouble(amount); } Override public void readFields(DataInput in) throws IOException { this.orderId in.readUTF(); this.amount in.readDouble(); } Override public int compareTo(OrderBean o) { // 先按订单ID升序再按金额降序 int cmp this.orderId.compareTo(o.orderId); if (cmp 0) { return Double.compare(o.amount, this.amount); } return cmp; } // getter、setter、toString 省略 }注意compareTo里金额是反着比较的Double.compare(o.amount, this.amount)表示当前对象的金额若小于对方则返回正数于是金额大的排在前面。这就是自定义降序排序的核心。第二步自定义分区器OrderPartitioner。要求 orderId 相同的记录进同一个分区可以把 orderId 的哈希值对 Reduce 数量取模。其实默认的 HashPartitioner 只要 Key 是 OrderBean它会拿整个 OrderBean 的哈希值去取模由于包含 amount 字段同一 orderId 的哈希值可能不同所以必须自己写。public class OrderPartitioner extends PartitionerOrderBean, Text { Override public int getPartition(OrderBean key, Text value, int numPartitions) { return (key.getOrderId().hashCode() Integer.MAX_VALUE) % numPartitions; } }第三步自定义分组比较器OrderGroupingComparator。分组比较器只比较 orderIdorderId 相等则认为这些记录属于同一个 reduce() 调用。注意分组比较器是WritableComparator的子类并且要重写compare(WritableComparable a, WritableComparable b)方法。public class OrderGroupingComparator extends WritableComparator { protected OrderGroupingComparator() { super(OrderBean.class, true); } Override public int compare(WritableComparable a, WritableComparable b) { OrderBean oa (OrderBean) a; OrderBean ob (OrderBean) b; return oa.getOrderId().compareTo(ob.getOrderId()); } }第四步在 Driver 里设置这个三个自定义组件job.setMapperClass(OrderMapper.class); job.setReducerClass(OrderReducer.class); job.setPartitionerClass(OrderPartitioner.class); job.setGroupingComparatorClass(OrderGroupingComparator.class); job.setMapOutputKeyClass(OrderBean.class); job.setMapOutputValueClass(Text.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(Text.class); job.setNumReduceTasks(2); // 至少 2 个 Reduce 才能看出分区效果Mapper 输出时把每行文本切成字符串数组orderId 放第一段金额转成 double然后context.write(new OrderBean(orderId, amount), new Text(整行))。Reducer 里遍历同一组的所有 Value由于排序器已经保证同一订单内金额降序所以直接把每行的 Text 原样输出就可以实现组内降序排列同时累加金额得到订单总额。这套解法覆盖了自定义排序、分区、分组三个知识点你做完之后再回头做“倒排序索引”那个题会有更清晰的感觉。倒排序索引的核心是把文档名作为 Key 的一部分把单词作为过滤条件本质上也离不开自定义 Key 的设计。4.3 倒排序索引的两种实现思路倒排序索引在实训里也高频出现题目一般给定若干文档要求输出单词, 文档名:词频的形式也就是一个单词对应的文档列表。思路有两个。第一种思路Map 阶段输出单词-文档名, 1在 Reduce 阶段先把同一单词-文档名的词频累加然后再次做一轮 MapReduce 把同一单词的多个文档合并。这种方式的优点是两个 Job 都非常浅显缺点是跑两轮耗时翻倍。第二种思路也是我推荐的做法只跑一个 Job在 Map 阶段输出单词, 文档名#词频的复合格式利用自定义分组按单词聚合Reducer 里遍历当前单词的所有文档并拼接结果。由于每个文档对同一个单词只贡献一个计数Map 端可以用一个小的 HashMap 先统计当前输入分片里每个单词在每个文档出现的次数再输出。这样可以显著减少 Map 端的输出量。这种思路的难点不在于代码本身而在于你能否理解Reduce 接收到的values迭代器遍历完一遍后如果再遍历一遍得到的是空列表因为迭代器是一次性的。我自己做这类题的经验是先把输入数据和期望输出写在草稿纸上模拟一遍数据会经过哪些 Map 任务、哪些分区、哪些 Reduce 任务然后把每个阶段输出的中间数据画出来。画完之后代码基本不会写错。5. 常见故障与性能调优我踩过的那些坑和排查思路5.1 内存溢出的三种典型场景与对策MapReduce 跑挂最常见的原因就是内存溢出。症状一般是Java heap space或者Container killed by ApplicationMaster。我归纳成三种典型场景。第一种是 Map 端缓冲区问题。默认每个 Map Task 的内存缓冲是 100MB如果单条数据特别大或者 Map 吐出的 Key-Value 对特别多溢写磁盘时会频繁排序合并导致内存和磁盘 I/O 双双飙升。这时适当调大mapreduce.task.io.sort.mb到 200MB 或 256MB同时把mapreduce.map.sort.spill.percent从默认的 0.8 调整到 0.7 或 0.9。调大缓冲可以降低溢写次数但不要超过单个 Container 可用内存的一半否则数据还没落地进程先被系统 OOM 杀掉。第二种是 Reduce 端拉取数据过多。默认情况下 Reduce 端从 Map 端拉取数据的内存缓冲由mapreduce.reduce.shuffle.input.buffer.percent控制默认值是 0.7。如果所有 Map 的输出都压向同一个 Reduce这个 Reduce 的缓冲会被灌满然后频繁落盘合并最终可能内存溢出。你得结合数据倾斜方案一起处理单纯调大内存治标不治本。第三种是mapreduce.map.memory.mb和mapreduce.reduce.memory.mb设置不当。许多集群的默认值只有 1GB如果你的 Mapper 内部要加载字典、做复杂的资源初始化1GB 根本不够。调大这些参数时同时要调大相应 Container 的虚拟内存比例比如yarn.nodemanager.vmem-pmem-ratio否则会出现明明物理内存够用但 yarn 判定超限而杀掉 Container 的假性 OOM 现象。5.2 任务卡在 99% 或长时间 Running 的检查顺序任务卡在 99% 是个极其常见的现象尤其是大作业联调的时候。我总结了四条排查步骤。第一步看日志。打开 YARN 的 ResourceManager Web UI找到失败或卡住的 Application进入 Logs优先查看syslog。MapReduce 的完整日志分为 stdout、stderr 和 syslog真正的异常栈通常在 stderr 中但 syslog 里会有框架级别的大量线索。第二步检查数据倾斜。看有没有某个 ReduceTask 的 Shuffle 字节数远大于其他任务。如果有按照 2.3 节的加盐思路处理。第三步检查资源分配。如果集群资源紧张多个作业同时提交Map 和 Reduce 任务会在队列里干等表现就是进度卡住不动。这时候看一眼 YARN 的队列调度情况必要时降低并行度或者错峰提交。第四步检查自定义代码中的死循环。这个坑很隐蔽比如你在reduce()里写了一个while (true)但退出条件永远不满足或者你的compareTo方法逻辑不对称——比如 A.compareTo(B) 返回正数B.compareTo(A) 也返回正数——这会让排序算法陷入混乱表现就是任务长时间不结束。这类问题没有捷径只能反复审查自己的比较逻辑。5.3 实战调优参数速查表我把日常调优中用到的高价值参数整理成一张表参数的具体含义在不同 Hadoop 版本里略有差异但大体一致。参数名默认值作用我的建议mapreduce.task.io.sort.mb100Map 端排序缓冲区大小数据量大且内存充足时调到 200-256mapreduce.map.sort.spill.percent0.8缓冲区溢写阈值内存吃紧时降到 0.7减少单次溢写数据量mapreduce.reduce.shuffle.input.buffer.percent0.7Reduce 端 shuffle 内存占堆比例若 Reduce 逻辑简单可调大到 0.8mapreduce.reduce.shuffle.parallelcopies5Reduce 并行拉取 Map 输出的线程数集群网络好时可调到 10-20mapreduce.job.reduces1Reduce 任务个数根据数据量和集群规模设为 2 到集群节点数之间mapreduce.map.memory.mb1024Map Container 内存复杂 Map 逻辑调到 2048-4096mapreduce.reduce.memory.mb1024Reduce Container 内存与 reduce 的 shuffle 缓冲联动调整需要特别强调mapreduce.job.reduces默认是 1很多实训题里如果不显式设置不管数据多大都只有一个 Reduce这会让你的任务耗时很长时间而且无法体现分布式效果。一般来说Reduce 数量可以设成接近集群可用核数的 70% 到 90%。不是越多越好因为每个 Reduce 都有启动和调度开销还会产生大量小文件。5.4 日志与调试技巧怎么快速定位坏数据MapReduce 里经常遇到一种情况程序逻辑完全正确但跑到中间就报NumberFormatException或者反序列化失败。这通常是输入数据里有脏数据。排查方式不是去改代码加一堆 try-catch而是先定位是哪一行数据出了问题。通行的做法是在 Mapper 里捕获异常并输出上下文信息。你可以临时写一个 SafeMapper把map()方法包起来异常时打印当前行号和行内容到 stderr。Override protected void map(LongWritable key, Text value, Context context) { try { String[] fields value.toString().split(\t); // 业务解析逻辑 } catch (Exception e) { System.err.println(Bad line at offset key.get() : value); // 选择跳过或抛出异常 } }但这里有一个取舍如果你直接throw new RuntimeException(e)作业会失败但你可以快速从日志里找到脏数据的位置如果你吞掉异常作业会继续跑完但结果可能缺失数据。我的习惯是第一次先让它失败看到脏数据后再决定要不要跳过。生产环境我会把脏数据单独写到一条特殊 Key 上攒到一个专门的输出路径里方便定期检查数据质量。日志查看模式也有技巧。我习惯在map()和reduce()的开头增加一个计数器context.getCounter(MyCounters, processed_lines).increment(1);这样在任务结束时可以从 Job 的 Counter 面板看到总共处理了多少行数据和输入文件的行数做对比立刻能确认有没有数据被遗漏。这个技巧比翻日志高效得多。6. 从实训到生产我对 MapReduce 学习路线的个人体会很多读者是从“头歌”的实训题点进来搜这篇文章的我特别能理解那种被作业蹂躏的感觉。当时我做“自定义分组”那一关时反反复复改了五六遍最后才发现自己只是在 Driver 里忘了设置setGroupingComparatorClass导致所有排序做了但分组逻辑压根没生效。以后遇到这类问题我建议你先列一个检查清单分区器装了没有排序比较器写了没有分组比较器配了没有Map 输出类型和 Reduce 输出类型是否混淆了这个清单能帮你解决 60% 的作业 Bug。再往深一层说做实训题的目的是让你通过代码理解框架而不是把代码背下来。我的学习路线是先拿 WordCount 跑通流程再拿排序题理解分区和分组然后拿倒排序索引理解多阶段数据流转最后自己设计一个小项目——比如把网约车订单数据做清洗和统计就正好对应热词里的“网约车大数据综合项目”。这类项目会让你面对真实场景数据有缺失、有重复、有格式不一致你的清洗逻辑必须扛得住脏数据这才算真正出师。招聘数据清洗这个实训也挺典型。清洗不是说把空行删掉就行而是要看字段是否符合业务规则。比如年龄字段不能在合理范围之外手机号要满足位数要求。用 MapReduce 做清洗时我建议你在 Map 阶段只做过滤和格式化在 Reduce 阶段做聚合统计不要把一个流程全都塞进 Map。这样如果你要扩展新的清洗规则只需要在 Map 里加规则Reduce 基本不用动。还有一个容易被忽视的重点理解 HDFS 和 MapReduce 的关系。MapReduce 的计算是流动的它会尽量把计算任务调度到数据所在的节点上这就是数据本地性。如果你提交作业时输入路径是 HDFS 上的文件而集群节点上恰好有这些文件的副本框架就会优先在那个节点上启动 Map 任务省去大量网络传输。我见过有人把几百 GB 数据从本地上传到 HDFS结果还是慢原因就是数据只存在于一个节点上所有的 Map 任务都要通过网络去那个节点拉数据根本没有本地性可言。解决方法是增加副本数或者用hdfs dfs -setrep -w 3把关键数据文件的副本数调上去。关于 MapReduce 的未来我的观点前面已经说了它并不会因为你学了 Spark 就变得没有价值。MapReduce 训练的是你“把任意业务拆成可并行化的键值处理”的思维方式。这个思维方式和具体的计算引擎无关。我在写 Flink 作业时遇到复杂事件处理仍然会先问自己如果我用 MapReduce 会怎么拆这一问往往能帮我把一个问题从混沌状态理清楚。最后留一个我常用的小技巧在本地调试 MapReduce 时不需要每次都启动集群。你可以把mapreduce.framework.name设为local这样整个 Job 会在本地 JVM 里以模拟模式运行方便打断点调试。hadoop jar wordcount.jar com.example.WordCountDriver -Dmapreduce.framework.namelocal input output要注意的是本地模式默认只跑一个 Map 一个 Reduce无法模拟分布式网络传输但用来验证业务逻辑是否正确完全够用。等你把逻辑调通了再去集群上跑出错的概率会低很多。这个习惯让我少走了非常多弯路希望也能帮到你。
返回列表