多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

MapReduce核心原理与实训实战:从排序到数据清洗

MapReduce核心原理与实训实战:从排序到数据清洗 如果让我用一句话概括MapReduce我会说它是大数据时代“先分后合”思想最经典的落地产物。你只需要写好一份Map函数和一份Reduce函数剩下的切分、调度、容错、排序框架全替你扛住。过去十年里几乎每一门大数据课程、每一个实训平台、每一张面试题里都有它的身影——从最基础的WordCount到排序、分组排序、倒排索引再到招聘数据清洗这类综合案例MapReduce既是入门Hadoop的第一道门槛也是检验你分布式思维成色的试金石。这篇文章我打算换个讲法不堆理论就围绕实训平台上大家反复刷的那些关卡展开为什么MapReduce能排序、自定义排序到底在改什么、分组排序和普通排序差在哪、倒排索引为什么要用复合键以及招聘数据清洗这类综合任务如何一步步拆成Map和Reduce。无论你是刚写完第一个WordCount的新手还是卡在某个排序关卡里的“头歌难民”这篇文章都值得你静下心读完。1. MapReduce核心设计思路与执行流程1.1 “先分后合”的计算哲学MapReduce的思想本质上是六个字分而治之。别小看这六个字它是整个框架的立身之本。一个TB级文件单机处理要跑十几个小时但把它切成几百个128MB的分片扔到集群的不同节点上并行处理可能几分钟就出结果了。MapReduce要解决的核心问题不是“怎么算”而是“怎么把一个大任务拆成无数小任务再把无数小结果合并成大结果”。我习惯用一个厨房流水线的类比来解释。你是一家餐厅的老板要准备一千份盒饭。一种做法是你一个人从头忙到尾洗菜、切菜、炒菜、装盒累死也做不完另一种做法是招一批临时工有人专门洗菜有人专门切菜有人专门炒菜最后统一装盒。MapReduce就是第二种做法Map阶段相当于“切菜工”每个人领一小堆原料切成标准大小的片Reduce阶段相当于“装盒工”把切好的菜按照订单归好类一份一份装进盒饭里。中间原料怎么分配、哪个切菜工负责哪一堆这些不用你操心框架自己调度。这里的关键在于“分片”的粒度。Hadoop默认的块大小是128MBMapReduce会把输入数据按照InputFormat逻辑切分成一个个split每个split对应一个Map任务。split可以横跨多个文件也可以比块小但理想情况下一个split里的数据就落在本地节点的块上——这就是“移动计算比移动数据更划算”这句名言的由来。网络带宽永远是稀缺资源把计算逻辑下发给数据所在的节点比把几百GB数据拉到一台机器上再算要高效得多。1.2 作业从提交到完成的全过程很多教程会把MapReduce执行流程画成一张复杂的图然后讲一堆“JobTracker”、“TaskTracker”之类的过时概念。如今Hadoop 2.x之后的YARN架构已经把这些名词扫进历史了但核心流程的逻辑骨架没变。我用当前主流版本的流程来讲。一个MR作业从提交到跑完大致经历这样几个阶段。首先是客户端提交作业把JAR包、配置信息和输入路径打包提交给ResourceManagerResourceManager会分配一个Container在某个NodeManager上启动一个ApplicationMaster这个AM负责整个作业的生命周期管理。接着AM根据输入数据的split数量向RM申请Map任务的资源每个Map任务在NodeManager的Container里运行。Map任务的核心逻辑是这样的InputFormat先把输入数据切成split每个split交给一个RecordReaderRecordReader按行读取数据解析成key, value对喂给Mapper。Mapper的输出不会直接写磁盘而是先写进一个环形内存缓冲区默认大小100MB。缓冲区写到一定阈值默认80%时后台线程开始把数据溢写到本地磁盘溢写过程中会做分区和排序——这一步是MapReduce能“自动排序”的秘密所在。溢写产生的多个小文件最终会被合并成一个大的输出文件这个文件里的数据已经按照分区号排好序每个分区内部又按照key排好序。Reduce任务启动后先从Map端拉取属于自己分区的数据这个过程叫Shuffle。拉过来的数据还要再做一次合并排序然后传给Reducer。Reducer的输入天然就是按key排好序、且相同key聚在一起的流你在reduce()方法里拿到的那个IterableTextvalues就是这一组相同key对应的全部分组值。最后Reducer的输出写到OutputFormat指定的路径上默认是key\tvalue一行一条记录。1.3 数据本地性与移动计算为什么要这么大费周章地分片、溢写、排序、拉取一句话数据本地性。数据放在哪个节点最好就在哪个节点上计算避免跨网络搬运。HDFS会把块复制三份散落在不同机架Map任务调度时会优先选择存有对应数据副本身的节点。如果本地节点资源不够再退而求其次选同机架的节点最后才考虑跨机架。这个细节决定了你的作业跑得快不快。我见过太多人一上来就开几十个G的测试数据放在集群上跑结果Map任务大量落到了“数据不本地”的节点上光拉数据就拉了半小时。合理的做法是先用小数据把逻辑调通再逐步放大数据量。数据本地性是分布式计算里最底层也最容易被忽略的性能因素理解了它你才真正明白为什么MapReduce要叫“MapReduce”而不是简单的“多线程并发计算”。2. 核心编程模型与关键组件细节2.1 Mapper、Reducer与Driver的标准骨架无论实训平台上的关卡怎么换皮MapReduce编程永远逃不出三件套Mapper类、Reducer类、Driver主类。我先把最标准的Java骨架贴出来后面的实操讲解都基于这个模板展开。public class WordCount { public static class TokenizerMapper extends MapperLongWritable, Text, Text, IntWritable { private Text word new Text(); private 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) { word.set(w); context.write(word, one); } } } 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); } }这套骨架有几个地方值得抠细节。map()方法的输入LongWritable, Text是Current版本的默认输入格式TextInputFormat产出的key是字节偏移量value是一行文本。你要改输入格式、改key类型、改value类型都从这四个泛型参数入手。还有一个我踩过坑的地方setOutputKeyClass和setOutputValueClass只管Reduce端输出Map端输出的类型如果和Reduce端不一致必须单独用job.setMapOutputKeyClass()和job.setMapOutputValueClass()声明。很多新手在实训里报ClassCastException八成就是在这个地方栽了——Map输出的value是TextReduce输入却按IntWritable解析框架在Shuffle阶段的序列化器直接抛异常。2.2 Writable序列化体系为什么不能直接用Java对象MapReduce的Shuffle过程要把数据写出磁盘、拉过网络这就要求Map输出的key和value必须能被序列化。Java原生的Serializable接口确实能用但它序列化出来的字节体积大、性能差Hadoop干脆自己搞了一套轻量级序列化体系所有key/value类型都要实现Writable接口。所以你在写Mapper的时候会发现Text对应Java的StringIntWritable对应intLongWritable对应longDoubleWritable对应doubleNullWritable对应null引用。这些Writable类怎么用直接决定你的代码能不能跑通。我自己踩过的典型坑是在继承WritableComparableT写自定义key时write()方法里用out.writeUTF()写入字符串但readFields()里却用readLong()去读字段顺序错位导致反序列化出来的对象全是乱的。自定义Writable类的write和readFields必须严格保持字段顺序一致而且最好在类里加一个无参构造器——Hadoop反射创建对象时默认调用它你不写就会在Reducer初始化时报NoSuchMethodException。2.3 Partitioner、Combiner与Comparator四个容易混淆的组件实训平台上的排序关卡之所以难是因为它把四个概念揉在一起考Partitioner分区器、Combiner合并器、SortComparator排序比较器、GroupingComparator分组比较器。这四个东西各管一段弄混了整道题就废了。Partitioner决定Map输出的每一条记录去哪个Reduce任务。默认的分区器是HashPartitioner它对key的hash值取模Reduce数量所以相同key必然进同一个Reduce。反过来你要自己写分区器比如让某个特定key单独进一个Reduce就继承PartitionerK,V重写getPartition()。Combiner本质是“Mini-Reducer”在Map端本地先做一次局部聚合减少Shuffle网络传输的数据量。但Combiner有个大坑它的调用次数不固定可能被调用0次、1次或多次所以Combiner必须满足交换律和结合律。你要是用Combiner算平均值结果就是错的——1和2的平均是1.51、2、3的平均却可能是(1.53)/22.25而不是2。求均值、求中位数这类操作千万别用Combiner。SortComparator管的是“Map输出和Reduce输入排序阶段用哪个比较器”。默认情况下用key的compareTo()方法。自定义排序要改排序逻辑最干净的方式是让key类实现WritableComparable并重写compareTo()如果不想改key类也可以单独写一个RawComparator子类通过job.setSortComparatorClass()挂载。注意RawComparator直接在字节层面比较序列化后的数据速度快但写起来晦涩实训题里一般用前者就够了。GroupingComparator则是Reduce端对“同一组”的判定标准。默认情况下排序时认为key相等compareTo返回0的记录会进入同一个reduce()调用。一旦你改了compareTo()排序时是“订单号相同且金额相同才算相等”但你业务上想把“同一订单号”的所有记录归到一组这时就必须用GroupingComparator把分组标准和排序标准解耦。2.4 InputFormat与OutputFormat的选择很多实训题不会明说要用什么InputFormat但读入数据的格式不同代码写法完全不同。文本文件默认走TextInputFormat一行一条记录CSV文件本质也是文本走的还是它但如果数据是自定义二进制格式就得自己写InputFormat。还有KeyValueTextInputFormat它把每一行按第一个分隔符默认是tab切分成key和value适合“键值对”格式的输入。SequenceFileInputFormat则对应Hadoop自有的二进制格式SequenceFile适合中间结果的存储。输出格式方面最常见的TextOutputFormat会把每条key, value写成key\tvalue换行结尾。实训平台比对结果时默认就按这个格式一行行比较。输出多了空格、多了tab、多了换行都会被判错。我建议所有输出格式都以“key和value之间用一个tab分隔、结尾只有换行符”为准这也是后面讲招聘数据清洗时最容易被忽略的扣分点。3. 排序与分组排序实训关卡的高频考点3.1 默认排序与“字典序”陷阱MapReduce天然会排序这句话让很多新手产生了误解以为框架会“按照我想要的规则”排序。实际上框架默认只做一件事按key的二进制字典序排序。数字按字典序排是什么概念10排在2前面因为字符1的ASCII码比2小。你要是把IntWritable字段直接转成Text去当key就等着排序结果惊掉下巴吧。实训平台的“MapReduce排序”第一关通常就是让你对一批数字做自定义排序让数字真正按数值大小排。这里就要用到WritableComparable了。一个比较常见的套路是定义一个IntPairWritable之类的复合key把原始数字包进去重写compareTo()比较逻辑是按数值比较而不是字符串比较。public class IntPair implements WritableComparableIntPair { private int first; private int second; public IntPair() { } public IntPair(int first, int second) { this.first first; this.second second; } public int getFirst() { return first; } public int getSecond() { return second; } Override public int compareTo(IntPair o) { int cmp Integer.compare(this.first, o.first); if (cmp ! 0) { return cmp; } return Integer.compare(this.second, o.second); } Override public void write(DataOutput out) throws IOException { out.writeInt(first); out.writeInt(second); } Override public void readFields(DataInput in) throws IOException { this.first in.readInt(); this.second in.readInt(); } Override public int hashCode() { return first * 31 second; } Override public boolean equals(Object obj) { if (this obj) { return true; } if (!(obj instanceof IntPair)) { return false; } IntPair o (IntPair) obj; return this.first o.first this.second o.second; } }这段代码是一个典型的复合key模板。first可以是主排序字段second是次排序字段。你可以自由调整compareTo()里的逻辑——比如让first降序、second升序只需把比较结果取反就行。这里要特别提醒hashCode()和equals()一定要实现好因为Partitioner默认用的就是key的hashCode()Reducer端合并时也会大量调用equals()来判断key是否相同。hashCode()返回值不靠谱会让相同key的记录分散到不同Reduce结果错得一塌糊涂。3.2 自定义排序的三个实现路径在实训里实现自定义排序一般有三条路可走。第一条最直观让key类实现WritableComparableT接口重写compareTo()。这个方案适用于key是自定义类型的场景代码量适中可读性好。第二条是给已有的Writable类型挂一个自定义RawComparator通过job.setSortComparatorClass()配置不改key类本身。它适合你想让Text按“忽略大小写”或者“倒序”排列的场景但写起来不直观出问题也难调试。第三条是干脆绕开排序Map端输出时把所有数据都设计成某个key的一个valueReduce端收到之后自己用Collections.sort()排一遍。这种做法只适合数据量小的场景否则一个Reduce会拖垮全作业。三条路怎么选我一般遵循一个原则排序需求是面向业务的就封装成自定义key类排序需求是临时调试的用RawComparator数据量小到能在内存里排用第三条偷懒。实训平台上的“自定义排序”关卡本质就是考第一条路。3.3 分组排序每个分组的Top N怎么取分组排序是MapReduce实训里让最多人卡壳的关卡它的典型场景是“求每个部门工资最高的员工”。假如数据是部门,姓名,工资三列你希望输出每个部门下面工资最高的那个人该怎么设计第一反应可能是Reduce端按部门分组循环比较每个员工的工资取最大值。这个思路没错但如果你想让“每个部门工资最高的人”作为整组第一条就出现输出上会更优雅。做法是设计一个复合key包含部门和工资两个字段compareTo()排序规则是“先按部门升序部门相同按工资降序”。这样在Reduce端同一个部门的记录会聚在一起并且工资最高的那条恰好排在组内第一条。然后你写一个GroupingComparator重写分组逻辑为“只按部门字段比较相等”那么Reducer收到的每一组第一条就是该部门工资最高的人。public class EmpGroupComparator extends WritableComparator { protected EmpGroupComparator() { super(EmpKey.class, true); } Override public int compare(WritableComparable a, WritableComparable b) { EmpKey ea (EmpKey) a; EmpKey eb (EmpKey) b; return ea.getDept().compareTo(eb.getDept()); } }然后在Driver里记得挂上job.setGroupingComparatorClass(EmpGroupComparator.class);这里有一个初学者必踩的坑光写了GroupingComparator没在Driver里调用setGroupingComparatorClass()或者写完之后重新打包的JAR是旧的导致分组逻辑完全不生效。我当年在实训里被这个问题卡了一整晚以为是排序逻辑写错了后来才发现是打包缓存的问题。记住改了Java代码一定要重新编译、重新打JAR、重新提交三个步骤一个都不能省。3.4 倒排索引原理解析与MapReduce实现实训平台上那个“倒排序索引”关卡名字看着像是在考“倒序排序”实际上是倒排索引Inverted Index搜索引擎的核心数据结构。它做的事是把“文档里有哪些单词”这个正向关系反转为“某个单词出现在哪些文档里、出现了几次”的反向关系。用MapReduce实现倒排索引核心难点在于Map阶段如何同时拿到“单词”和“文档名”两个维度的信息。HDFS上的每个文件都对应一个FileSplitMapper里可以通过context.getInputSplit()拿到当前处理的是哪个文件再取文件名public static class InvertedIndexMapper extends MapperLongWritable, Text, Text, Text { private Text wordDoc new Text(); private Text one new Text(1); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { FileSplit split (FileSplit) context.getInputSplit(); String fileName split.getPath().getName(); String line value.toString(); String[] words line.split(\\W); for (String word : words) { if (word.length() 0) { wordDoc.set(word \u0001 fileName); context.write(wordDoc, one); } } } }这里把word和fileName用\u0001这个不可见字符拼接成复合key输出value为1。Reduce端收到形如hello \u0001 doc1.txt的key和一组1求和后得到“hello在doc1.txt出现了几次”。但倒排索引的最终格式通常是word - doc1:次数, doc2:次数所以Reduce端不能直接输出复合key需要把key拆开重新组合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(\u0001); String word parts[0]; String fileName parts[1]; int sum 0; for (Text val : values) { sum Integer.parseInt(val.toString()); } result.set(fileName : sum); context.write(new Text(word), result); } }注意这样直接写会有个问题同一个单词如果出现在多个文档里它会被拆成多条独立记录而不是合并成一行。要真正形成“word - doc1:3, doc2:5”的格式还得加一个二次MapReduce或者在第一个Job的Reduce输出里用自定义OutputFormat做拼接。实训关卡的评测如果只比对“单词加文档名加次数”的中间结果那第一个版本就够了如果要求合并成一行就需要在Map端做一层二级聚合把同一单词的多条文档信息拼成一个字符串再输出。倒排索引这个例子能一次把“自定义key、自定义分组、Shuffle排序、复合键设计”全部串起来这也是实训平台把它放在排序章节里的原因。它的精髓不在“倒序”而在“用key的排序与分组特性把分散在不同文档里的词频信息汇聚到同一个Reducer里”。4. 招聘数据清洗综合案例4.1 清洗需求的业务拆解“招聘数据清洗”是MapReduce综合应用里最贴近真实业务的关卡。给你一份爬虫抓下来的招聘信息CSV里面什么脏数据都有公司名带空格、薪资写成“10K-20K·14薪”、城市写“北京”又写“北京市”、学历字段缺失、岗位ID重复。你要做的不只是简单过滤而是把数据规整到能直接做统计分析的程度。实训里通常会给你几个明确的清洗目标去重、剔除无效记录、统一字段格式、按城市统计岗位数、按学历统计平均薪资。拿到任务先别急着写代码先把数据结构和清洗规则列清楚。招聘数据一般有这些字段公司名称、岗位名称、工作地点、薪资范围、学历要求、经验要求、发布时间。清洗规则则包括去掉重复记录按公司岗位工作地点组合去重、过滤薪资为“面议”的记录、把“10K-20K”拆成最低薪资和最高薪资两个数值字段、把“北京”和“北京市”统一为“北京”。这一步的产出是一张“清洗规格表”明确每个字段的输入样例、清洗规则、输出格式。写MR作业之前把这张表定下来后面编码会顺畅很多。4.2 Map阶段的解析与清洗逻辑Map阶段在招聘清洗任务里不只做“分词”它承担的是“解析清洗规整”三件事。解析是指把CSV的一行文本拆成字段数组清洗是指过滤掉非法数据规整是指把字段转换成标准格式。public static class CleanMapper extends MapperLongWritable, Text, Text, Text { Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line value.toString(); String[] fields line.split(,, -1); if (fields.length 7) { return; } String company fields[0].trim().replace(\, ); String job fields[1].trim().replace(\, ); String city normalizeCity(fields[2]); String education fields[4].trim(); String salaryStr fields[3].trim(); if (company.isEmpty() || job.isEmpty() || city.isEmpty()) { return; } if (!education.contains(大专) !education.contains(本科) !education.contains(硕士) !education.contains(博士)) { return; } if (salaryStr.contains(面议)) { return; } int[] salary parseSalary(salaryStr); if (salary null) { return; } context.write(new Text(city), new Text(job \t education \t salary[0] \t salary[1])); } }这里有个细节值得多说一句CSV文件里如果一个字段值本身包含逗号正规CSV会用引号包起来比如北京,朝阳。直接用line.split(,)会把字段劈开导致字段错位。实训数据通常没这么刁钻但真实场景里一定要处理引号包裹的情况。我看到很多人的清洗代码只做两步split和filter完全不考虑“清洗后的数据输出给谁用”。在MapReduce任务里Map输出格式直接决定了Reduce能不能方便地做聚合。上例中Map输出的value是job\teducation\tminSalary\tmaxSalary拼接串Reduce端只需要再拆一次就能统计。如果你Map输出的是原始CSV行的原样valueReduce端还得重复解析一遍清洗逻辑代码冗余且容易被评测的格式差异搞崩。4.3 Reduce阶段的聚合与输出清洗完成之后Reduce阶段要做的是聚合统计。最常见的需求是按城市统计岗位数量、按学历统计平均最低薪水。按城市统计岗位数量逻辑很简单public static class CityReduce extends ReducerText, Text, Text, Text { Override protected void reduce(Text key, IterableText values, Context context) throws IOException, InterruptedException { int total 0; for (Text val : values) { total; } context.write(key, new Text(String.valueOf(total))); } }按学历统计平均最低薪资呢就不能在Map端用Combiner做局部聚合了原因前面说过——Combiner不能用于求平均值。Reduce端必须把每个学历对应的所有最低薪资收集起来统一求和再除个数数。如果数据量大到单台Reduce扛不住就得用二次聚合的思路第一轮MapReduce统计每个学历的总薪资和记录数第二轮再算平均值。实训数据一般不大直接一个Reduce搞定。输出格式是评测判分的命门。绝大多数实训关卡用TextOutputFormat它输出key\tvalue。如果你的统计结果是“城市加上岗位数”输出应该是北京\t123。多一个空格、多一个tab、换行用了\r\n都可能被评测脚本判为“结果不一致”。我建议提交之前先看一眼输出文件的前几行把肉眼能发现的格式问题提前扼杀掉。4.4 适配实训平台评测的细节实训平台和真实集群最大的区别在于它有一套固定的评测脚本比对的是“你的输出文件”和“标准答案”的差异而不是“你有能力解决问题”就行。这意味着你的输出格式必须和评测脚本严格匹配哪怕逻辑是对的格式错了也是零分。我在头歌这类平台上刷题总结了几个细节。第一不要在你的Map或Reduce代码里用System.out.println()打日志这些输出会混进标准输出流评测脚本解析时可能把它当成结果的一部分。调试信息请用System.err.println()或者通过Context的Counter机制记录。第二Reduce输出的key和value类型要跟评测期望对上比如期望输出Text, Text你非要输出Text, IntWritable虽然文本渲染出来都是数字但类型不同在某些评测逻辑下会报错。第三文件名和路径要严格按题目要求来输出目录是指定的/output你多创建一个/output2也没意义。还有一个小技巧本地跑通逻辑之后在提交评测前手动把Reduce输出的一行样例复制出来比对一下题目的示例输出格式。“北京 123”和“北京\t123”在肉眼上可能没什么区别但在评测脚本眼里天差地别。5. 常见问题与排查技巧实录5.1 实训中最高频的十个报错我把这些年带学生、刷题时踩过的坑整理成了故障速查表每一个都对应真实的报错场景。遇到问题先查表大部分情况能帮你少走两小时弯路。现象可能原因解决思路作业提交时报“Path already exists”输出目录已存在删除输出目录或换个名字Hadoop拒绝覆盖已有输出目录运行报ClassCastExceptionMap输出类型与Reduce输入类型不一致检查map输出与reduce输入的类型必要时用setMapOutputKeyClass自定义分组不生效Driving里没调setGroupingComparatorClass在主类里显式配置并重新打JAR再提交排序结果还是按字典序key是Text类型而非数值类型换IntWritable或自定义WritableComparableReduce数量为1跑得极慢没设NumReduceTasks或者所有key被分区到同一个Reduce检查Partitioner和setNumReduceTasks配置输出文件每行末尾多了个空格代码里拼接value时手滑加了多余空格检查所有字符串拼接逻辑按“\t”精确拼接倒排索引里取不到文件名直接强转getInputSplit失败确认没使用CombineFileInputFormat普通TextInputFormat可强转FileSplit有中文乱码输入文件编码不是UTF-8统一转成UTF-8或在Configuration里设置合适的字符集某个Reduce数据特别多数据倾斜key分布不均考虑加盐或二次聚合把热点key打散本地跑通集群跑挂本地默认是LocalJobRunner集群有资源限制和文件权限问题看日志排查常见于权限、依赖JAR缺失、路径权限5.2 手工排查三件套每次作业跑挂别急着改代码重试。先做三件事看日志看计数器跑小数据。日志怎么看YARN集群上执行yarn logs -applicationId app_id可以拿到ApplicationMaster日志和所有Container日志。日志里一般会直接打出异常栈顶最常见的NullPointerException、ClassCastException一眼就能定位到类和行号。本地模式调试时IDE控制台会直接打印Hadoop运行日志比集群上更直观。计数器是MapReduce自带的一个统计机制你可以在代码里给某个逻辑挂计数器context.getCounter(clean, emptyCompany).increment(1);然后在作业跑完后从Client端打印或者Web UI上看这个计数器的值。比如你Map端清洗时过滤掉了1000条字段为空的记录这1000就是Counter的输出。计数器不会影响输出格式是排查清洗逻辑的利器。我写数据清洗题时比如“脏数据有多少条被过滤掉了”“面议薪资占多少比例”都在Counter里临时打一版确认数量合理再删掉。最后一个小数据本地跑技巧输入只放几十行数据代码里加断点或日志快速验证逻辑正确性。本地跑法很简单把FileInputFormat输入路径指向本地文件系统的一个小文件输出目录也指到本地目录不需要启动任何集群。这样能跳过YARN调度的干扰把问题聚焦在数据处理逻辑上。5.3 关于MapReduce学习的一点心得讲真MapReduce在今天已经算不上最前沿的大数据计算引擎了Spark、Flink这些内存计算框架早就把“快”这个字做到了另一个维度。但我不觉得学MapReduce是浪费时间。恰恰相反MapReduce是理解分布式计算“为什么难”的最小样本数据怎么切、怎么跨节点传输、怎么容错、怎么排序、怎么聚合这些基本问题在任何分布式框架里都会重演。你弄懂了MapReduce的自定义key、分区器、分组比较器再去学Spark的partitionBy、groupByKey、reduceByKey会发现大量概念能一一对应上。如果你正在实训平台刷MapReduce的排序关卡别只盯着“怎么把这个用例跑通”。试着问问自己如果数据量放大一万倍我的方案还能不能跑如果同一个key的数据分布极不均我该怎么办把这些问号一个个解开你才算真正打通了分布式计算的任督二脉。最后说一个我个人的习惯每做完一个题目我都会把reduce输出结果手动算一遍拿小样本按业务逻辑手推再对比程序输出。这一步看似笨拙却能精准揪出“排序规则写反了”“分组维度搞错了”“类型转换丢精度”之类的高频错误。所谓经验很大程度上就是把错误一个个踩过去之后沉淀下来的判断力没有捷径只能靠题量和对细节的较真慢慢堆。
返回列表