Hadoop WordCount实战:从环境搭建到MapReduce调优全指南

Hadoop WordCount实战:从环境搭建到MapReduce调优全指南 简介面向Hadoop初学者的完整词频统计MapReduce实现基于Hadoop 2.2.0可直接用于学习MapReduce编程模型与词频统计任务的完整开发流程。压缩包共17个文件以7个Java源码与7个编译后的class文件为主附带1个十万单词级别的TXT测试文件以及.project、.classpath工程配置文件整体仅154KB。浏览人数已超5800人是经典的Hadoop入门参考资料。借助这份材料读者能对照源码与class理解Map、Reduce阶段的逻辑拆分掌握Job参数配置、文件输入输出及工程目录结构十万单词测试语料也便于直接运行并验证统计结果适合正在学习分布式计算、准备大数据实验或面试复习的开发者。1. 为什么拿词频统计开刀项目价值与运行环境准备先别急着敲代码。我见过太多人一上来就复制WordCount示例结果跑出来一个空文件或者压根提交不到集群上然后一头扎进报错堆里半天出不来。我自己当初学Hadoop时也踩过同一批坑所以这篇打算用词频统计这条主线把从环境搭建到代码编写、打包、提交、调优的完整链路串一遍。词频统计在Hadoop生态里的地位和编程语言里的Hello World差不多。它的核心目的是让你理解MapReduce的两个阶段Map阶段把输入数据拆成键值对Reduce阶段把相同键的值聚合起来。听起来简单但背后牵扯到分片、排序、分区、洗牌这些MapReduce框架的核心机制把这些搞明白后面再接触Hive、Spark这类上层计算引擎理解成本会低很多。1.1 版本选择与安装清单先说版本。JDK建议用1.8Hadoop用3.2.4或3.3.x系列。Hadoop 3.x已经把默认端口从50070换成了9870网上很多老教程还在用旧端口写访问地址你如果照着配完发现打不开Web界面先看看是不是端口写错了。安装模式我建议分两种情况选。如果只是做课程设计或者学习验证伪分布式模式就完全够用一台机器同时跑NameNode、DataNode、ResourceManager、NodeManager配置比完全分布式简单很多排错也方便。如果是想练手集群运维再上三节点的完全分布式。伪分布式的核心配置集中在四个文件里$HADOOP_HOME/etc/hadoop/hadoop-env.sh $HADOOP_HOME/etc/hadoop/core-site.xml $HADOOP_HOME/etc/hadoop/hdfs-site.xml $HADOOP_HOME/etc/hadoop/mapred-site.xml $HADOOP_HOME/etc/hadoop/yarn-site.xmlhadoop-env.sh里把JAVA_HOME指向你实际的JDK路径。网上最常出现的问题就是这里没改或者改成了export JAVA_HOME/usr/lib/jvm/java-8-openjdk-amd64这种不存在的路径导致启动HDFS时直接报错。core-site.xml配置NameNode地址hdfs-site.xml配置副本数和NameNode数据存储目录伪分布式副本数设为1就够了configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property /configurationconfiguration property namedfs.replication/name value1/value /property /configurationmapred-site.xml指定用YARN来调度MapReduce任务yarn-site.xml配置ResourceManager和NodeManager的地址。这一套配完后先执行hdfs namenode -format格式化NameNode再用start-dfs.sh和start-yarn.sh启动服务。用jps命令检查进程能看到NameNode、DataNode、ResourceManager、NodeManager这四类进程就说明环境基本就绪了。1.2 准备输入数据在HDFS上建目录环境起来之后下一步是往HDFS里放数据。很多人习惯在本地随便建个文件就往Hadoop提交实际上MapReduce默认从HDFS读数据你得先把文件传上去。# 创建HDFS目录 hdfs dfs -mkdir -p /input # 把本地文件传到HDFS hdfs dfs -put /opt/data/wordcount.txt /input/ # 查看确认 hdfs dfs -ls /input这里我建议准备的数据文件别太小也别只有一两行。用一段几百词、包含重复单词的英文文本最好这样方便观察Map和Reduce的真实处理过程。我实战时习惯专门在文本里放几个大小写不同的相同单词比如Hadoop和hadoop还会放标点紧贴单词的情况用来检验代码对数据清洗的处理能力。2. 手写WordCountMapper、Reducer、Driver三段式拆解接下来是核心代码。这一步的目标不光是让程序跑通更重要的是理解每一行代码在框架里干了什么。2.1 Mapper类的实现与数据流理解import java.io.IOException; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; public class WordCountMapper 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 { 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); } } }注意四个泛型参数LongWritable是输入键的类型表示行偏移量Text是输入值的类型表示一行文本后面的Text和IntWritable是输出键值对类型。这里有个很容易被忽略的细节——Hadoop的序列化体系。你不用Java自带的String和Integer而是用Text和IntWritable是因为它们实现了Writable接口序列化效率更高能直接在网络间传输。Mapper的输入是HDFS文件里的每一行框架会逐行调用map()方法。split(\\s)是按空白字符切分这里的\\s匹配空格、制表符、换行等。实际项目中文档往往有大量标点符号直接按空格切会让Hadoop,这种带逗号的词和Hadoop变成两个不同的词所以如果要做清洗可以在切分后加一句w w.replaceAll([^a-zA-Z], )。2.2 Reducer类的实现聚合逻辑import java.io.IOException; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.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); } }Reducer的输入是(单词, [1, 1, 1...])的形态框架把所有相同键的值聚合到一个迭代器里。这个Iterable看起来像集合但底层可能是因为数据量太大边取边释放所以它只能遍历一次不能反复遍历。有些刚入门的人会尝试在循环外再取一次values结果取不到任何东西就是这个原因。2.3 Driver类作业提交的入口与参数设置import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; public class WordCountDriver { public static void main(String[] args) throws Exception { if (args.length ! 2) { System.err.println(Usage: WordCountDriver input path output path); System.exit(-1); } Configuration conf new Configuration(); Job job Job.getInstance(conf, word count); job.setJarByClass(WordCountDriver.class); job.setMapperClass(WordCountMapper.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这行很多人不明白作用它是让框架知道去哪个jar包里找类。你用IDE开发时class文件散落在target目录提交到集群时不可能把这些散落的class文件一个个传过去Hadoop要求把整个作业打成jar包setJarByClass就是指定这个jar包的位置。打包时注意不要把hadoop依赖打进去Hadoop运行环境自带了这些类库。3. 编译打包的坑从Maven配置到jar包提交代码写完后第一次运行大概率会遇到问题。下面是最容易踩的几个地方。3.1 pom.xml最小配置与maven-shade-plugin用Maven构建时pom.xml里最省心的配置方式是这样的依赖只声明hadoop-client版本和集群保持一致打包插件选用maven-shade-plugin顺便把Main-Class指定好project modelVersion4.0.0/modelVersion groupIdcom.example/groupId artifactIdwordcount/artifactId version1.0/version packagingjar/packaging properties maven.compiler.source1.8/maven.compiler.source maven.compiler.target1.8/maven.compiler.target /properties dependencies dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version3.2.4/version /dependency /dependencies build plugins plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-shade-plugin/artifactId version3.2.4/version executions execution phasepackage/phase goals goalshade/goal /goals /execution /executions configuration transformers transformer implementationorg.apache.maven.plugins.shade.resource.ManifestResourceTransformer mainClasscom.example.WordCountDriver/mainClass /transformer /transformers /configuration /plugin /plugins /build /project这样配置之后mvn clean package会生成一个可执行jar包java -jar能直接找到入口类。如果你是手工用javac编译需要把所有依赖class文件一起打进去比较麻烦我建议直接上Maven。3.2 提交命令与两个高频报错打包完成后提交命令长这样hadoop jar target/wordcount-1.0.jar /input /output注意第二行运行路径中的核心报错提示——jar does not exist or is not a normal file: /usr/local/hadoop/share/hadoop/m...这是网上抱怨最多的问题之一。这个报错的原因是把CLASSPATH相关参数误传给了hadoop jar命令。正确命令应该是hadoop jar 包路径 主类全路径 参数但有人会在主类位置误写成Hadoop自带库路径或者在命令里多了-libjars相关参数但jar包路径没写对。还有情况是大小写或路径不对/usr/local/hadoop/share/hadoop/mapreduce下的jar包文件确实存在但你输入的是/usr/local/hadoop/share/hadoop/m路径被截断了当然找不到。这个改动其实很容易排查先看包路径是否完整再看主类有没有写全类名最后确认jar包在hadoop用户下有读取权限。把这三步都走一遍99%的报错都能解决。第二个高频报错是Output directory hdfs://localhost:9000/output already exists。HDFS为安全考虑不允许覆盖已有输出目录。# 删除已有输出目录 hdfs dfs -rm -r /output我实战时习惯每次运行前先看下输出目录在不在在就删掉否则反复提交作业会一直报错。3.3 提交后怎么看日志作业提交后控制台会刷出一堆日志。重点看两部分一是INFO mapreduce.Job: Running job: job_xxxxxxxxxx这一行记录了作业ID二是作业完成时的计数器汇总。# 查看YARN上运行的作业列表 yarn application -list # 查看某个作业的日志 yarn logs -applicationId application_xxxxxxxxxx如果跑挂了先看控制台有没有FAILED字样。我遇到过好几种情况是Container killed、Java heap space、Disk out of space这些是集群资源问题和代码没关系。训练阶段如果真的跑在真集群上注意别让多个作业同时抢资源。如果只是伪分布式一般不会有资源问题。4. 运行结果解读从控制台到HDFS输出文件作业跑完去HDFS上看输出文件。hdfs dfs -ls /output hdfs dfs -cat /output/part-r-00000 | head -20输出目录里通常有两个东西_SUCCESS空文件和part-r-00000数据文件。part是分区前缀r表示这个文件来自Reducer阶段如果是m则表示Mapper阶段直接输出。part-r-00000是第一个分区的结果默认只有一个Reduce任务时就只有一个文件。如果你在结果文件里看到类似下面的内容Hadoop 3 is 2 hadoop 1先别急着认为程序有问题。Hadoop和hadoop被当成两个词是因为Map阶段没有做转小写处理。is前面有点奇怪查一下原始数据多半是用了\t制表符切分或者原文确实带不可见字符。这些不是框架问题是数据清洗的问题。从网上搜索词来看很多人也在搜“hadoop hdfs hive安装”“hadoop和spark hive结合使用”“如何将数据存入hadoop平台”之类的话题。想说明的是词频统计跑通后再往上走就是Hive建表做SQL统计或者Spark读HDFS做分布式计算都是同一个数据源、不同计算引擎的演进路径。词频统计在你理解了MapReduce之后再回头看会觉得它很简单但它是把这些上层工具的原理串起来的关键一环。5. 一次完整排查记录从“作业一直pending”到“Container失败”这部分内容是网上教程里最难找到的但对实战最有用。我拿自己真实踩过的一个坑来拆解。5.1 表面现象作业提交成功控制台显示Running job但等了五分钟还没进Map阶段再过一会儿直接报Application application_xxx_0001 failed 2 times due to AM Container for appattempt_xxx_0001_000002 exited with exitCode: 15.2 排查链路第一步先查YARN日志yarn logs -applicationId application_xxx_0001翻到堆栈最底部看到的是java.io.IOException: No space left on device。第二步查磁盘空间。df -h发现根分区满了。清理了一下本地的/tmp目录问题依旧。第三步意识到这不是本地磁盘是DataNode的存储目录。默认配置下DataNode数据目录在$HADOOP_HOME/tmp/dfs/data而HDFS小文件过多会占满inode也会报No space left。用hdfs dfsadmin -report查看各个DataNode存储情况看到剩余空间为零才确认是HDFS存储满了。第四步清理HDFS上无用的旧输出目录和日志文件同时把fs.trash.interval设成0让删掉的文件不进回收站property namefs.trash.interval/name value0/value /property重启HDFS后再跑作业秒过。这个坑的启发是看到Container失败别先改代码优先看资源。磁盘满、内存不足、CPU配额超限都是大数据的常见死法和业务代码关系不大。5.3 伪分布式最常被忽略的内存配置在笔记本上跑伪分布式默认的YARN内存配置经常把作业卡死。原因很简单NodeManager默认分配8G内存但虚拟机或小内存机器根本没那么多可给。你在yarn-site.xml里加上这些配置把内存调小property nameyarn.nodemanager.resource.memory-mb/name value4096/value /property property nameyarn.scheduler.maximum-allocation-mb/name value4096/value /property property nameyarn.nodemanager.vmem-check-enabled/name valuefalse/value /propertyvmem-check-enabled这个参数尤其值得注意。虚拟机里的虚拟内存和物理内存比例可能不满足默认要求导致Container频繁被杀掉关掉这个检查是伪分布式场景下最有效的保命手段。6. 词频统计还能怎么扩展跑通基础版不算完。我建议你沿着下面这几个方向做延伸每做一步对MapReduce的理解就深一层方向一Combiner优化。在Driver里加一行job.setCombinerClass(WordCountReducer.class)让每个Map任务先在本地做一次局部求和减少Shuffle阶段的数据传输量。注意Combiner不是万能药它要求元素的合并操作满足交换律和结合律。词频统计的加法满足所以可以复用Reducer类做Combiner。很多面试题会问“Combiner和Reducer的区别”你实际跑一遍对比一下运行时间比死记硬背强得多。方向二自定义数据清洗。在Mapper里加正则去标点、转小写、过滤停用词做一个更接近真实场景的文本处理逻辑。这能帮你建立“数据质量决定结果质量”的意识毕竟真实业务里的数据比这脏多了。方向三改用Hive跑同一份数据。把文本文件加载成Hive表用一段SQL完成同样的统计CREATE EXTERNAL TABLE wordcount(text_line STRING) ROW FORMAT DELIMITED LINES TERMINATED BY \n LOCATION /input; SELECT word, COUNT(*) AS cnt FROM ( SELECT explode(split(text_line, \\s)) AS word FROM wordcount ) t GROUP BY word ORDER BY cnt DESC;Hive底层翻译后也是MapReduce任务。你在Hive里写个SQL就能出结果的体验和你手写Java实现Task的体验一对比就会明白“大数据开发为什么要分层”。方向四换成Spark跑。用Scala或PySpark写词频统计整个代码可以控制在十行以内from pyspark import SparkContext sc SparkContext() text_file sc.textFile(hdfs://localhost:9000/input) counts (text_file.flatMap(lambda line: line.split( )) .map(lambda word: (word, 1)) .reduceByKey(lambda a, b: a b)) counts.saveAsTextFile(hdfs://localhost:9000/output-spark)你会直观感受到内存计算和磁盘计算的差异。这也是网上的热词里为什么“hadoop spark hive”经常一起出现的原因——它们是三条互补的技术路线。关于Hadoop的Docker镜像我也说两句。很多人为了省去环境配置的时间喜欢直接拉一个Hadoop镜像来跑。这确实能快速起环境但要注意网络模式和端口映射。用Docker跑伪分布式容器一重启NameNode元数据如果没做持久化挂载就得重新格式化这个坑我已经踩过一次了。学习阶段我还是建议在实体机或虚拟机上手动装一次配置过程本身就是最好的学习材料。最后分享一个我测试时的小技巧把输入文件分成小份多次测每次改一行代码就跑一次观察计数器变化再对照日志理解框架行为。比如你可以在Mapper里故意不设setOutputValueClass看看会不会报错报什么错可以把Reducer的输出改成切分到多个文件观察分区逻辑。这样主动制造问题、定位问题的过程比照着教程抄十遍代码都有用。本文还有配套的精品资源点击获取