MapReduce核心原理与实战:从环境搭建到数据清洗、YARN调度全解析
发布时间:2026/9/30 15:09:37来源:尧图网络
做大数据开发这些年MapReduce一直是个绕不开的话题。尤其当你去面试Hadoop相关岗位或者在公司里接手离线计算任务时不管是Presto、Hive还是Spark底层都脱离不了MapReduce这套分而治之的思想。今天这篇内容我想从理论机制到代码落地把MapReduce核心原理、环境搭建、排序实战、数据清洗案例一直到YARN提交流程和常见报错排查一次性讲透。整个内容适合三类人正在学大数据准备课程设计的学生、刚转行做数据开发的工程师、以及想要巩固底层原理去应对面试的朋友。这篇博文的内容规划我会按照我自己当年踩坑过来的顺序来写先理解MapReduce为什么是这个设计再动手搭伪分布式环境然后从基础编程到排序专题再到综合案例的数据清洗练习最后补上YARN、ZooKeeper集群整合和面试高频题。文中的所有代码、配置、报错信息都来自我实际测试过的环境可以直接拿来当参考。1. MapReduce核心机制拆解它到底在解决什么问题1.1 从分而治之说起MapReduce的诞生逻辑MapReduce不是某个公司发明的概念Google在2004年发表的论文《MapReduce: Simplified Data Processing on Large Clusters》把它体系化了后来Hadoop做了开源实现。它的核心思想只有四个字分而治之。举个例子你要统计一本1000页书里每个单词出现的次数一个人从头翻到尾可能要一天如果分成10个人每人负责100页各统计各的局部结果最后汇总10份结果时间就能压缩到十分之一左右。MapReduce就是把这件事抽象成两个阶段Map阶段负责并行处理局部数据Reduce阶段负责汇总局部结果。这里有一个关键点往Map输入数据、从Reduce输出结果输入输出都是键值对的形式。整个框架把并发、容错、数据分发、负载均衡全部封装起来了程序员只需要关心两条业务逻辑——map函数怎么写、reduce函数怎么写。这也是MapReduce能火起来的原因它把一个大数据计算问题的高度从分布式计算降到了单机函数逻辑门槛一下子低了很多。但是这里面有个容易误解的地方分而治之的前提是数据能被切分。如果业务逻辑强依赖全局数据比如要计算全量数据的某个分位数单纯靠Map和Reduce这个过程是做不到的这也是为什么后来Spark的RDD、DataFrame能提供更丰富算子的原因。MapReduce的表达能力有限但它胜在简单稳定设计目标就是面向TB甚至PB级的离线批处理。1.2 Map、Shuffle、Reduce三阶段的数据流转细节一个完整的MapReduce作业数据要经历的路径是这样的HDFS上的输入分片InputSplit被读取后框架调用InputFormat去解析记录每一条记录会交给map函数处理map输出的键值对会先写入环形缓冲区在缓冲区溢出到磁盘时做一次预排序然后进入Shuffle阶段。Shuffle是Map和Reduce之间最关键的枢纽它分Map端和Reduce端两边来看。Map端Shuffle的具体动作map输出先放到一个100MB可以通过mapreduce.task.io.sort.mb配置的环形内存缓冲区里当缓冲区使用率达到80%时后台线程开始将数据溢写到本地磁盘。溢写过程中会执行分区Partitioner决定数据去哪个Reduce和分区内排序。这里要注意默认的HashPartitioner是按key的hashcode模上reduce Task数量的所以理论上相同key一定会进入同一个Reduce但数据不一定均匀这就产生了数据倾斜的可能。Reduce端Shuffle的动作Reduce Task启动后会启动Fetcher线程从各个Map Task所在节点拉取属于自己分区的数据。拉取到的数据先放内存缓冲区如果数据量大会先合并排序再落盘。等到所有Map Task的数据都拉取完毕后Reduce进入Sort阶段对所有key做归并排序。这里排序不是Reduce干的而是Map端溢写时就做了一遍区内排序Reduce端主要是做归并这就保证了整个框架里相同key的数据一定是连续出现的。最后才调用reduce函数逐key处理。如果你仔细看这个过程会发现在纯MapReduce里很难做全局排序——因为数据进了不同分区每个分区内部有序但分区与分区之间没有全局顺序。真要全局排序就要用单一的Reduce Task但那样并行度就废了。所以工程上一般取巧先把数据按总排序范围划分区间再配合TotalOrderPartitioner做全局有序。1.3 移动计算而非移动数据这个设计对性能的影响MapReduce有个著名的口号移动计算比移动数据更便宜。在HDFS里数据默认有三个副本分布在不同的节点上如果一个100GB的数据集和一个计算程序同时在一个集群里框架会把计算任务调度到数据所在的节点上执行而不是把100GB的数据从存储节点传到计算节点。这个思想在大数据领域后来被广泛延续Spark里的数据本地性Data Locality也是同一个逻辑。具体到实现层面JobTracker/ResourceManager在分配任务时会优先选择数据本地性最好的节点数据在本地就选本地NODE_LOCAL如果有跨机架的数据那就选同机架的节点RACK_LOCAL最后才是任何可用的节点OFF_SWITCH。如果我们写代码的时候没有意识到这一点比如频繁在Map里通过网络读取HDFS上其他文件的数据就会破坏数据本地性原则整个作业的IO开销会成倍放大。我实际测试过一个任务在完全相同的输入数据下把数据本地性从NODE_LOCAL调整到OFF_SWITCH作业时间差了将近2倍。所以这不仅仅是一个概念是直接影响作业调度速度的关键要素。2. 开发环境从零搭建伪分布式与开发工具选型2.1 伪分布式vs集群学习阶段该怎么选很多搜索热词都在问hadoop伪分布式搭建hadoop集群搭建和从零开始安装hadoop说明大家第一道坎就是环境。我的建议是如果机器配置一般4G内存以下第一次学习用伪分布式就够了。伪分布式的本质是所有守护进程NameNode、DataNode、ResourceManager、NodeManager都在同一台机器上运行你看到的进程是完整的配置文件也是按集群模式写的只是节点数量是1。集群模式一般出现在课程设计或者公司环境至少要3台机器需要配置SSH免密登录、主节点和从节点的JAVA_HOME、masters与workers文件旧版本叫slaves、core-site.xml里的namenode地址、hdfs-site.xml里的副本数还有yarn-site.xml里的ResourceManager地址。集群多起来的第一个坑就是节点之间时间不同步会导致心跳异常所以生产环境一般都要配NTP。学习阶段没必要一上来就玩这种复杂度的东西先在伪分布式上把MapReduce原理跑通再上集群会顺手很多。另外提一句Docker。现在很多教程用Hadoop的Docker镜像来搭建开发环境这个做法对折腾配置的人来说很友好——镜像拉下来启动容器Hadoop的启动脚本一敲就完事。但Docker网络和端口映射有时候会带来额外的排除难度我在第7节会专门写Docker环境下的坑。如果你对Linux不熟建议还是老老实实按虚拟机或者云服务器的方式走一遍。2.2 Ubuntu伪分布式搭建完整步骤与关键配置我在Ubuntu 20.04上从零搭过一套Hadoop 3.3.x伪分布式环境下面把关键步骤和踩坑点整理出来。第一步准备Java环境。Hadoop 3.x要求JDK 8以上新版3.4.x甚至已经必须JDK 8u202。安装完java后记得设置JAVA_HOME在/etc/environment里加上JAVA_HOMEPATH并在/etc/profile.d/下写一个hadoop.sh方便以后改版本。Hadoop本身用的是自己的JVMMapReduce任务也是在独立的JVM里跑的所以JAVA_HOME如果没写对启动NameNode时可能不报错但跑MapReduce作业会直接报找不到Java环境。第二步下载和解压Hadoop。建议去Hadoop官网下载稳定版二进制包不要下源码包——网上很多教程给的是编译好的版本版本差异可能导致libhadoop.so不匹配跑本地库报错。解压后重点修改五个文件core-site.xml里最关键的是fs.defaultFS。伪分布式通常配成hdfs://localhost:9000。hdfs-site.xml里是副本数dfs.replication配1因为伪分布式只有一个DataNode默认3副本会一直触发副本冗余警告。还要把dfs.namenode.name.dir和dfs.datanode.data.dir指向你创建的数据目录不要用默认的/tmp因为/tmp在系统重启后会被清掉一次重启你整个HDFS数据就没了这个坑我踩过。mapred-site.xml需要指定使用YARN框架mapreduce.framework.name配置为yarn。yarn-site.xml里要配置yarn.nodemanager.aux-servicemapreduce_shuffle同时因为单机内存有限要把yarn.nodemanager.resource.memory-mb调低到1024或2048yarn.scheduler.maximum-allocation-mb也同步调低否则作业一跑就容易被NodeManager杀掉。第三步是格式化NameNode。第一次启动之前要执行hdfs namenode -format只有第一次需要后面如果再手贱执行一次NameNode的namespaceID会变DataNode注册不上就会报Incompatible clusterIDs。格式化完成后用start-dfs.sh和start-yarn.sh启动jps命令能看到NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager五个进程就说明环境OK了。启动HDFS之后浏览器打开http://localhost:9870Hadoop 3.x的NameNode Web端口是9870老版本是50070你应该能看到DataNode活着。打开http://localhost:8088能看到YARN的资源管理界面。这个验证步骤不要省因为很多环境问题在跑作业前是不会有任何提示的。2.3 Windows下使用IDEA搭建Hadoop开发环境搜索热词里有一长串都是在问Windows下怎么用IDEA搭Hadoop开发环境这部分确实值得单独讲。Windows本机不装Hadoop守护进程但MR的开发调试完全可以在IDEA里搞定。先说最简单的方案在Windows本地装一个与远端Hadoop版本一致的Hadoop解压包配置HADOOP_HOME环境变量然后把winutils.exe和hadoop.dll放到bin目录。这一步是很多教程容易忽略的——Hadoop在Windows上运行本地库时需要winutils.exe来模拟Unix权限和文件操作。如果缺了它跑本地模式的MR程序会直接报错或者卡在native code上。然后我强烈建议你不要用默认的LocalJobRunner跑生产级作业你可以在IDEA的resources目录放一个core-site.xml把fs.defaultFS指向远端的hdfs://IP:9000再放一个hdfs-site.xml配上远端NameNode的地址和RPC端口。这样在IDEA里直接运行main方法作业会以本地客户端YARN远端集群的方式提交——本地只负责产生JobConf、打包jar和上传资源计算全在集群跑。在这个模式里你需要有一个hadoop-client的Maven依赖。版本要和集群保持一致我测试过Hadoop 2.7的代码去连3.x集群会出现RPC协议不兼容的报错建议直接用3.3.x。dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version3.3.6/version /dependency在Windows的IDEA里跑MapReduce还有一种方式把Job的FileInputFormat输入路径指向本地文件系统job.setJarByClass指向当前类不配置任何远端地址这时候MapReduce会以local模式LocalJobRunner运行所有Task都在本地JVM里模拟执行。对于纯学习MapReduce API和调试逻辑阶段来说这种方式最轻量不需要启动任何Hadoop进程跟写普通Java程序一样。但要注意local模式下Shuffle是内存里直接搞的网络开销、磁盘溢写完全模拟不出来所以你调出来的性能结论在真实集群里不一定成立。2.4 验证环境跑第一个自带示例环境搭好之后不要急着写代码先用Hadoop自带示例验证环境。在Hadoop目录下执行hadoop jar share/hadoop/mapreduce/hadoop-mapreduce-examples-3.3.6.jar wordcount /input /output其中/input是hdfs上你准备的一个包含若干文本文件的目录/output是输出目录必须不存在。如果代码路径上有我的测试文件正常跑完后你会看到终端打出Map和Reduce的进度百分比最后在/output目录里出现part-r-00000文件。我第一次跑这个示例的时候花了很长时间排查一个问题输出目录一直显示FileAlreadyExistsException。原因是Hadoop的输出分区目录会在配置阶段就被校验如果目录已经存在无论里面有没有数据都会直接报错。这是一个非常容易踩的坑后续我们写所有MR Job之前最好在main方法里写一个删除输出目录的工具方法后面关于这个问题我再细讲。3. 基础编程实战从WordCount到自定义Writable3.1 MapReduce三件套Mapper、Reducer与Driver逐行拆解初学MapReduce编程最经典的入口就是WordCount。它麻雀虽小但五脏俱全一个最基础的程序里包含了MapReduce Job的所有必填配置项。Mapper的基类是org.apache.hadoop.mapreduce.Mapper里面最核心的方法是map(K key, V value, Context context)。WordCount里输入是偏移量LongWritable作为key、一行的文本内容Text作为value。map把一行按空格切分每个单词输出一个(word, 1)。public class WordCountMapper 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[] words value.toString().split(\\s); for (String w : words) { if (w.length() 0) continue; word.set(w); context.write(word, one); } } }有个细节值得注意Hadoop的网络传输、序列化都是基于Writable接口的所以它没有用Java自带的String和Integer而是用了Text和IntWritable。Text不是一个简单的String包装它内部维护字节数组set方法会复用内部数组频繁set比每次都new一个Text要高效得多。在数据量巨大的时候这个差别能显著减少GC压力和内存分配。Reducer的基类是org.apache.hadoop.mapreduce.Reducer核心方法是reduce(K key, Iterable values, Context context)。这个接口很讲究values不是一次性全部加载进内存的List而是一个迭代器框架会逐个value地取。如果你在reduce里想多次遍历同一个key的values必须先把values拷到一个List里否则第二次迭代永远是空。很多新手在reduce里做两层循环或多次聚合时踩到这个坑原因就在这里。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); } }Driver部分是我们提交作业的入口也是最容易被忽略的业务代码。很多网上教程里Job.setJarByClass(WordCount.class) 直接填当前类但如果你的类在多个jar包里或者主类不在jar包里这里就需要填一个真正被打进jar包的类名。如果你是用IDEA的Maven管理依赖写的用maven-shade-plugin把依赖打成一个肥jar后提交会省去一堆ClassNotFound的麻烦。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(WordCountMapper.class); job.setCombinerClass(WordCountReducer.class); job.setReducerClass(WordCountReducer.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(WordCountReducer.class)这一行。Combiner是Map端的本地聚合器它做的事情和Reducer一样但它运行在每个Map Task节点上。WordCount场景下一个Map Task可能输出几万条(word,1)如果直接在Map端先做一次求和再发给Reduce网络传输量能减少90%以上。但Combiner不是所有场合都能用——它要求合并函数满足结合律和交换律否则结果会错。比如求平均值的操作就不能直接用Reducer当Combiner这就是面试里经常问的Combiner和Reducer的区别。3.2 自定义序列化对象实现Writable接口的细节实际业务中Map的输出往往不是简单的Text和IntWritable而是一个包含多个字段的对象比如日志里有用户ID、访问时间、页面URL。你可以用Text把多个字段拼成一个长字符串再在Reduce端解析但这样做的可维护性和性能都很差。更好的办法是自定义一个Writable类。自定义Writable要继承Writable接口实现write和readFields两个方法。write方法负责把对象序列化到DataOutput里readFields负责从DataInput反序列化。这两个方法的字段顺序必须完全一致。下面是流量统计场景里的一个Bean示例public class FlowBean implements Writable { private long upFlow; private long downFlow; private long sumFlow; public FlowBean() {} public FlowBean(long upFlow, long downFlow) { this.upFlow upFlow; this.downFlow downFlow; this.sumFlow upFlow downFlow; } Override public void write(DataOutput out) throws IOException { out.writeLong(upFlow); out.writeLong(downFlow); out.writeLong(sumFlow); } Override public void readFields(DataInput in) throws IOException { this.upFlow in.readLong(); this.downFlow in.readLong(); this.sumFlow in.readLong(); } // getter、setter、toString略 }必须提供无参构造因为Hadoop反射创建对象的时候会调用无参构造器。如果你只自己写了有参构造框架在反序列化时就会直接抛异常。另外toString方法一定要写了它会影响输出的格式化结果如果你不重写toString最后输出到文件里的是一串对象地址排查数据的时候会特别痛苦。FlowBean被用作出参或入参时如果它作为Map输出的Value那只需要实现Writable。但如果它作为Map输出的Key比如你要把整个Bean作为排序的Key就必须实现WritableComparable接口这在排序实战部分会细讲。3.3 控制ReduceTask并行度的三个参数Map并行度是文件分片数量决定的一个InputSplit对应一个Map Task。Reduce并行度则完全由代码控制用job.setNumReduceTasks(int n) 设置。默认情况下如果代码里不设置Hadoop会用1个Reduce Task。但它不是随便设置越大越好——ReduceTask数量太多每个Task处理的数据量就少且启动、调度开销剧增太少数据都堆在一个节点上整体资源利用率低。另外分区数量必须和Reduce Task数量配套。默认的HashPartitioner会把每个key映射到0到numReduceTasks-1的分区编号如果ReduceTask设为1不管Hash值怎么算所有数据都进同一个分区。所以如果你想对特定key做特殊处理比如让某个重要key全部进一个Reduce可以自定义Partitioner在getPartition方法里写业务判断逻辑。public static class MyPartitioner extends PartitionerText, FlowBean { Override public int getPartition(Text key, FlowBean value, int numPartitions) { if (key.toString().startsWith(vip)) { return 0; } return 1 % numPartitions; } }如果你设置了numReduceTasks等于2但Partitioner返回了编号2或3作业提交时会直接报错。这回导致你在YARN界面里看到作业是FAILED但日志里看不到任何Reduce端异常因为错误实际发生在提交阶段。4. 排序专题实战自定义排序、分组排序与倒排序索引4.1 自定义排序实现WritableComparable接口的两种方式排序是MapReduce里最经典也最考验细节的考点之一。默认情况下Map输出的key会按字典序或数值序自动排序但这个自动排序依赖key的类型实现了Comparable接口。IntWritable的compareTo就是按int值比大小Text的compareTo就是按UTF-8字节序比字典序这是MapReduce框架天然就有的能力。想要按自定义规则排序核心是让key实现WritableComparable接口。这个接口同时继承了Writable和Comparable所以实现它不仅要写序列化方法还要写compareTo方法。以FlowBean为例需求通常是对流量做个排序先按上行流量升序上行相同再按下行流量降序。compareTo的写法决定了Map端溢写排序和Reduce端分组排序的结果Override public int compareTo(FlowBean o) { if (this.upFlow ! o.upFlow) { return (int) (this.upFlow - o.upFlow); } return (int) (o.downFlow - this.downFlow); }注意compareTo返回值的含义返回负数代表当前对象排在前面正数代表排在后面0代表相等。如果你把它倒过来全表顺序就会完全翻转。还有一种方式更轻量不需要改Bean直接在Driver里用job.setSortComparatorClass指定一个自定义的RawComparator类。这个类直接对序列化后的字节做比较不需要反序列化成对象效率更高且可以同时解决版本兼容的问题。大部分场景用实现WritableComparable就够了。4.2 分组排序二次排序的实现与常见误区在Hadoop MapReduce核心原理相关的热词搜索里分组排序出现得特别多对应的头歌平台作业、课程设计也特别多。分组排序本质上不是排序它要解决的是reduce一次调用处理多少个value的问题。默认情况下框架对进入Reducer的数据按key分组相同key的所有value会放在同一个Iterable里reduce函数被调用几次取决于有多少个不同的key。但如果你自己定义了分组规则可以让key不同但业务上相关的value进入同一个reduce调用。经典场景是按订单号为key做二次排序同一个订单的所有商品明细按照时间排序然后reduce一次处理完一个订单的所有明细。实现二次排序需要三个步骤第一用组合key比如订单号时间构成一个OrderKey对象实现WritableComparablecompareTo先比订单号订单号相同再比时间。第二自定义Partitioner保证订单号相同的数据分到同一个ReduceTask里。第三自定义GroupingComparator让Reduce在分组时只比较订单号忽略时间字段。这样框架先对数据按订单号时间排序保证同一订单内时间有序分组时又按订单号合并reduce拿到的是一个订单号下的有序时间明细列表。public class OrderGroupComparator extends WritableComparator { protected OrderGroupComparator() { super(OrderKey.class, true); } Override public int compare(WritableComparable a, WritableComparable b) { OrderKey oa (OrderKey) a; OrderKey ob (OrderKey) b; return oa.getOrderId().compareTo(ob.getOrderId()); } }这里有一个特别容易踩的坑使用WritableComparator时super构造方法里的第二个参数必须传入true表示框架会创建真实的对象实例来比较否则compare方法里强转类型时会得到空对象或者直接报类转换异常。还有一点要注意的是分组排序不是replace默认的分组而是重写。如果你自定义了GroupingComparator却没有自定义Partitioner那相同订单号的数据可能被分到多个ReduceTask里此时分组逻辑就算对了也没用因为数据根本不在同一个Reduce里。这三件套——组合Key、Partitioner、GroupingComparator——是配套使用的缺一不可。4.3 倒排序索引的实现思路倒排序索引Inverted Index是搜索引擎的基础数据结构在MapReduce课程设计里也是高频题目。需求一般是这样给你一堆文档输出每个单词出现在了哪几个文档里以及每个文档中出现过多少次。直接用一个MapReduce作业做倒排索引是可行的但输出会遇到一个问题Reduce端输出的value是一个列表如果文档数量很多这个列表可能非常大超过单条记录的合理大小。我看到很多课程设计给的方案是Map阶段针对每篇文档的每个单词输出(word, doc1:1)Reduce阶段直接拼接所有文档信息。这个方案的问题在于没有做Combine成倍放大了Shuffle的数据量。更合理的做法是分两阶段或者说用两次MapReduce串联。第一个作业做常规的WordCount但额外记录文档名输出(word, docName:count)第二个作业再按word聚合所有docName:count。第一个作业的Map端Combiner可以合并同一个文档里的重复单词减少中间数据。第二次MapReduce开始前可以在第一个Reduce的输出中先做一个全排序便于最终的索引按字典序展示。这个例子说明了一个道理实际生产环境中的业务很少有单个MapReduce就能解决的大多数需要串行多个Job所以理解Job之间的衔接比单纯会写一个Mapper更重要。4.4 合并去重如何用MapReduce做数据裁剪顺带说一个MapReduce在ETL里最常见的用途——合并去重。因为MapReduce天然会把相同key聚在一起去重就变得非常简单Map阶段把需要去重的字段作为keyvalue随便写Reduce阶段直接取第一个value输出即可。也可以用MapReduce自带的机制如果你只需要key去重Map输出后的key会自动去重但你如果还想要其他字段就必须在Reduce里处理。我在做网约车数据清洗时用过一个技巧利用TextOutputFormat的key和value分隔符把多列拼成去重主键把几列拼接成比如city_id|driver_id|order_time作为Map输出的key然后在Reduce里输出原始行数据。这样做有两个好处一是去重逻辑完全由框架的shuffle保证不需要自己写额外的HashSet——因为数据是从分布式并行节点扫的二是当数据量极大时你自己在Map端维护一个全局HashSet是不可行的内存会直接溢出。5. 综合案例实战招聘数据清洗与网约车统计5.1 招聘数据清洗脏数据过滤的MapReduce思路很多搜索词里反复出现实验4 mapreduce综合应用案例 — 招聘数据清洗和网约车大数据综合项目——基于mapreduce的数据清洗这类实训题目的本质都是ETL清洗。招聘数据清洗的典型场景是输入是一堆CSV招聘记录含公司名、岗位、城市、薪资区间、学历要求、经验要求、发布日期等字段。数据里普遍存在空值、格式不统一、薪资字段写成10k-20k或者面议等非结构化文本。处理思路一般分两级。Map阶段做第级清洗解析CSV分隔符时注意转义字符把脏数据用continue跳过不符合字段数量、空值比例过高的记录直接丢弃。但注意Association的一个问题是出错的记录到底该丢弃还是保留需要业务方定夺——比如核心字段公司名、岗位是空的一般直接丢弃如果只是次要字段如福利标签为空可以给个默认值补齐。第二级清洗在Reduce阶段做按公司或岗位分组后清洗组内的记录字段格式比如薪资解析成统一的数字单位。如果某些字段存在一个主键多条记录的情况还可以在Reduce里做最新记录去重。我实际做这类项目时有三个经验第一Map阶段一定要加一个计数器Context.getCounter(ETL, invalid_lines).increment(1)这样你最终跑完能看到清理了多少条不用再去数文件。第二TextInputFormat的输入默认是按Text类型解析的如果你的字段里有中英文逗号混用直接用split(,)会切出很多空字符串最好先用正则或者CSV解析库统一清洗一遍再进Map。第三输出不要写到HDFS上生成本地文件直接把清洗结果输出成Parquet或ORC格式的列式存储文件后续统计和查询都会快很多。5.2 网约车数据清洗与多维度统计网约车数据清洗这个项目我在之前文章里聊过一次这里把MapReduce的实现思路再展开说一下。数据集一般包含订单号、司机ID、乘客ID、上下车时间、经纬度、费用、金额等字段。清洗的重点是过滤上下车时间异常下车早于上车、时间差为负或者超过24小时、过滤经纬度在合理范围之外的单子、过滤费用为0或者异常负数的情况。清洗完之后往往紧接着就是统计。同一个作业里Map阶段的输出key可以设计成日期|城市|司机IDReduce阶段直接算每日各司机流水、接单量、平均客单价。这里有个MapReduce经典性能痛点——数据倾斜。网约车场景里热门城市的订单量可能比冷门城市高一个数量级如果我们按城市聚合轮到热门城市的ReduceTask就要处理大量的数据其他ReduceTask却闲得没事。解决思路有两个一是加盐salted key把热点key先打散到多个临时Reduce扛住压力再对中间结果做第二次汇总二是使用自定义Partitioner把热点key均匀分到多个分区避免单点堆积。我第一次实施时选的是加盐方案在Map输出key的前面加一个随机数后取模N第一轮Reduce把流量打散第二轮MapReduce再按真实key合并效果非常明显。5.3 交通信息分析系统一个完整的课程设计参考如果要做一个基于Hadoop的交通信息分析系统的设计与实现这类课程设计除了代码之外有两点值得提前规划数据的持久化和可视化。数据持久化是指清洗后的结果不能直接丢给前端应该把HDFS上的聚合结果导出到一个关系数据库或ES里。可视化则是课程设计展示加分项。一般来说整个系统的MapReduce部分只需要三个作业串起来就可以覆盖大部分功能点第一按卡口和日期统计车流量输出到MySQL第二按车辆号牌去重统计在途车辆数输出到CSV第三用倒排序索引实现车牌-卡口-时间的关联查询。你还可以在里面加入MapReduce自带的Combiner来优化第二和第三步的时间这部分是答辩提问时一定会被问到的点。还有一个容易被忽略的点是集群参数调优。课程设计的演示机往往只有一两台机器输入数据量也就几万条这种情况下跑MapReduce时间主要花在框架启动和JVM初始化上而不是真正的计算。所以如果发现你的程序在小数据量下跑得比单机遍历还慢那是正常的——MapReduce从启动到调度到结束大约有20到30秒的固定开销这也侧面应了没有足够大数据量就没必要上MapReduce这句话。6. 进阶关键链HDFS、YARN调度与ZooKeeper整合6.1 HDFS与MapReduce的关系为什么输入输出离不开HDFSHDFS是MapReduce的数据底座。MapReduce从HDFS上读InputSplit计算结果写到HDFS的part文件里中间数据通过YARN的分布式缓存分发到节点上。二者之间最重要的约定是数据分块与持久化。HDFS默认把文件切成128MB的块Hadoop 3里dfs.blocksize默认128MB旧版本是64MB每个块有3个副本。MapReduce的InputSplit默认对齐块边界一个块对应一个MapTask这样Map就在数据本地执行不用跨节点拉数据。如果你把HDFS的块大小调小会导致MapTask数量变多、每个Task处理量变少框架的调度开销就会变大但如果把块调得过大一个MapTask要处理的数据就过多负载不均会更明显。实际经验是在没有特殊压缩格式的情况下输入规模在GB以下时块大小保持默认即可不要调。MapReduce写输出时还会涉及HDFS一个非常重要的机制——推测执行。如果一个节点的MapTask跑得比同批任务慢很多YARN会在另一个节点上启动同一个任务的副本以快的那一份结果为准慢的那个会被kill掉。这个机制能提升集群整体吞吐但在你调试代码阶段会让日志非常难跟踪经常会在YARN日志里看到Task attempt_xxx_0002_m_000001_0_1 finished这类信息不要慌它只是推测执行生效了。6.2 作业提交到YARN的完整流程Hadoop作业提交到YARN的流程是面试里的高频题也是理解MapReduce和YARN关系的钥匙。整个流程可以这样走一遍第一步Job Client提交作业到ResourceManagerJob.waitForCompletion会先上传jar包、配置和依赖到HDFS上的一个临时目录默认/tmp/hadoop-yarn/staging。第二步ResourceManager收到请求后分配一个Application ID并启动一个ApplicationMasterAM进程AM是一个独立的容器负责调度Mapper和Reducer。第三步AM向ResourceManager申请资源拿到资源后在NodeManagerNM上启动Map Task每个Map Task是一个独立的JVM。第四步Map Task跑完后AM再启动Reduce TaskReduce Task到Map Task所在的节点拉取数据。第五步整个作业完成后AM向ResourceManager注销自己同时清掉临时目录里的staging文件。这里面有一个面试官很喜欢问的细节Map和Reduce的Task执行在哪台节点上答案是由容错和本地性决定的——Map Task尽量调度到数据所在节点Reduce Task则没有这个限制只看资源是否充足。另一个细节是ApplicationMaster如果挂了怎么办ResourceManager会自动重启AM重启后会重新申请已经失败的TaskContainer整个过程对客户端透明所以你可以看到作业状态一直切换却还跑着。如果你要用代码控制这个流程可以在Job配置里加上job.setQueueName(hadoop_queue)把作业提交到指定队列。生产环境上队列管理是YARN里常见的资源隔离手段面试问到多租户资源隔离时经常涉及这个。6.3 Hadoop与ZooKeeper整合HA集群的基石什么是Hadoop和ZooKeeper整合核心场景是HDFS的HA模式High Availability。在单NameNode架构里NameNode挂掉整个集群就不可用。HA模式部署两个NameNode一个是Active一个是StandbyZooKeeper负责管理和切换主备关系。整合过程要做的事情有第一在core-site.xml配置ha.zookeeper.quorum指向三个ZooKeeper节点比如node1:2181,node2:2181,node3:2181。第二配置dfs.nameservices以及dfs.ha.namenodes同时配置dfs.namenode.rpc-address和http-address。第三在hdfs-site.xml里配置dfs.ha.automatic-failover.enabledtrue启用自动故障转移。第四启动顺序有讲究先启动ZooKeeper集群然后启动JournalNode用于Active和Standby之间同步edits日志最后才格式化NameNode并启动HDFS。我第一次搭HA集群时遇到过Failover failed的报错排查了很久才发现是ZooKeeper节点的选举端口没开。ZooKeeper之间通信默认需要2888和3888端口如果防火墙只开了2181NameNode的HA状态同步永远无法成功。所以如果你在云服务器上搭HA一定要把安全组的端口范围放开。ZooKeeper在Hadoop生态里还负责另一个事——MapReduce框架的ResourceManager HA。ResourceManager也依赖ZooKeeper做active和standby状态存储配置在yarn-site.xml里的yarn.resourcemanager.ha.enabled。如果你看到RM服务在YARN UI上显示standby或者作业提交后长时间卡在ACCEPTED状态大概率是RM没有成功获取Active权限这时候先去检查ZooKeeper状态。6.4 集群模式主节点配置要点现在热词里还有一组高频词hadoop集群模式主节点配置hadoop集群搭建这套配置在伪分布式基础上多了几件关键事。首先是masters和workersHadoop 3.x叫workers2.x叫slaves文件。masters文件里写上主节点hostnameworkers文件里写入所有从节点hostname每行一个注意不能有多余空格。其次所有节点之间要配置SSH免密登录最简单的方式是主节点生成密钥执行ssh-copy-id把公钥分发到所有节点。第三各节点的时间要同步否则心跳时间错乱DataNode会被误判为宕机。第四hdfs-site.xml里的副本数配置通常设成小于等于DataNode数量比如三个DataNode设为3五个DataNode可以设为3或2。集群模式比伪分布式更容易在启动时报错最常见的就是Network is unreachable这类。我第一次搭3节点集群时所有进程都能看到但DataNode卡在in safe mode后来才发现是secondaryNameNode没有配好导致checkpoint进程每次都在重试连接一直占着资源。安全模式的原理是NameNode启动初期会做元数据恢复和块报告收集只有当DataNode上报的块达到阈值默认98%后才会自动离开安全模式。如果集群数据本来就很少且是刚格式化的日志一直提示in safe mode属于正常现象等几秒就好了但如果一直不退出就要检查有没有哪个DataNode起不来或者dfs.replication设置和实际副本数不一致。集群配置还有一个容易踩的点hdfs-site.xml里如果单独给某个DataNode设置了dfs.datanode.data.dir以外的自定义目录目录权限不对也会导致块写入失败在YARN日志和DataNode日志里会看到Permission denied。Hadoop的目录权限模型类似Unix但它是走Java的FileSystem API所以直接用chmod -R 777 /data并行目录来保证权限是常见做法。7. 常见问题排查与面试考点整理7.1 报错速查表从内存溢出到端口冲突我把实际工作里碰到的MapReduce报错整理成了一个速查表下面的每一条都是我或我的同事在真实集群里遇到并排查过的。报错信息可能原因解决办法Container killed by YARN for exceeding memory limitsMapper或Reducer内存超限调大yarn.nodemanager.resource.memory-mb同时调大mapreduce.map.memory.mbJob failed due to stage FAILED: Task Failed通常为代码逻辑或资源不足先取日志看具体异常排除ClassNotFound、NullPointer后轮询YARN UIRemoteException: FileAlreadyExistsException输出目录已存在提交前先delete输出目录No such file or directory本地或HDFS路径不对检查输入路径首尾空格用hdfs dfs -ls确认目录存在NameNode is in safe modeHDFS刚启动或异常恢复安全模式会自动退出若一直卡住检查DataNode在线数量和块报告Bad connect to dfs.namenodeNameNode进程没起或端口被占用执行jps确认NameNode进程netstat查端口9000/9870状态~/.ssh/config: line 1 errorSSH配置语法错误检查SSH的known_hosts和config文件Packet failed from node: java.io.IOException: Connection reset by peer数据节点间网络异常或DataNode宕机检查网络丢包和无响应节点看DataNode日志7.2 实际问题排查实录一次数据倾斜的分析过程我在做网约车数据清洗项目时遇到过最头疼的问题就是数据倾斜。当时统计城市维度订单量作业总共10个ReduceTask结果有一个Task跑了2小时还没结束其他9个Task早就成功。排查思路是这样先在YARN的Counter页面看每个Task的输入记录数大概率发现某个Task读了几亿记录其他Task只有几百万。再用hdfs dfs -count看热点key的数据分布。确认是数据倾斜后我给key加盐把热点key先拆分成多个带随机后缀的伪key第一轮Reduce做局部聚合第二轮再合并。加盐还有一种变体直接用CompositeKey按城市名业务参数组合成多个虚拟分区同样能缓解数据倾斜。数据倾斜是MapReduce运维排查里最常被问到的点之一。它的根本原因大多是业务本身的长尾特征比如个别key的基数极大。解法无非三类加盐、自定义Partitioner使热点key均匀分布、Map端Join替代Reduce端Join。如果是超大key本身无法分割或者倾斜发生在大量空Key上处理时必须先判断这些空Key是否有业务含义不要盲目过滤。7.3 面试高频考点MapReduce系列题目与答题思路热点词里搜hadoop面试题的人很多我把MapReduce方向最常考的几道核心题目整理成了一份简版题库同时写一下答题思路。第一问MapReduce的数据流是什么答题思路分两端客户端提交JobResourceManager分配ApplicationMaster并调度TaskMap端数据切分→环形缓冲区→溢写→分区排序→合并Reduce端拉取→合并排序→分组→reduce函数→输出。第二问为什么说Shuffle是MapReduce的核心因为Shuffle决定了数据如何在不同节点间高效流转是影响作业性能的第一因素Shuffle做得好排序、聚合和分组效率和容错性就能保证。第二问MapReduce无法处理什么样的问题流式计算、实时计算、机器学习迭代计算如迭代模型训练需要反复扫描数据以及强依赖全局状态的算法比如PageRank……这也是Spark存在的原因。第三问说一下Combiner、Partitioner和GroupingComparator的作用和区别。这题答得好代表你真的会写代码建议画图并用Map端局部聚合、路由分发、Reduce端分组这个口径作答。第四问你在实际项目里是怎么调优的如果简历里写了MapReduce项目这个几乎是必问的。调优维度无非减小数据量压缩、列剪枝、合适的文件格式、提高并行度、合理设置ReduceTask个数、避免数据倾斜、开启Combiner、设置推测执行。7.4 石墨文档之外的效率工具与后续扩展最后想分享一点关于效率的经验。MapReduce本身作为编程模型很基础但在实际开发里如果频繁手写成本确实高。现在生态里写离线计算一般直接用Hive、SparkSQL或Flink底层引擎在做优化时很多还是沿用了MapReduce的思路。你学MapReduce不意味着以后每天都要写Mapper和Reducer但它的诸多概念——分片、分区、Shuffle、Combiner、推测执行、数据本地性——仍然像地基一样支撑着你理解上层框架。另外有一点非常实用的建议如果需要跑多个MapReduce作业串联比如清洗后再统计、统计后再建索引不要一个一个手动submit用Oozie或Azkaban这类工作流调度工具把它们串起来。这样可以省掉大量等待和手动管理作业生命周期的痛苦也能在重复执行时快速重跑。写在最后的实操经验写到这里MapReduce从理论到落地的主干和细节基本覆盖了。最后分享几个我这几年来实操中沉淀下来的经验算是帮大家避坑。第一个经验是调优时先调数据、再调资源、最后才调代码。很多人一上来就调内存、调并行度但如果你输入数据是纯文本、没有压缩那么哪怕资源调到极致IO开销依然巨大。先改成SequenceFile或Parquet再上Snappy压缩往往能让作业时间直接减半。资源调优反而是三板斧里最不太需要花太多时间的那一个。第二个经验是排查问题要从YARN日志和Counter下手不要只盯终端输出。客户端打印的log是汇总信息真正的Task错误细节在NodeManager的日志目录container日志里。我记得有一次跟踪作业失败原因怎么都看不出问题后来去AM日志才看到Map端自定义类里的一个静态字段被多线程共用导致数据错乱。日志是唯一的线索别省这一步。第三个经验是如果你是学生课程设计想选基于Hadoop的XX系统这类题一定不要只交一个MapReduce统计脚本。老师更看重系统完整性——你最好把整个流程串起来从数据采集、清洗、入库、调度再到可视化大屏。这部分才是真正的加分项。把这篇文章里自己跑几遍你大概率就可以应对Hadoop MapReduce编程实战、课程设计和大部分面试题了。有问题随时在评论区留言我看到都会回。
网站建设高端定制企业官网