ARTICLE DETAIL

资讯详情

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

Hadoop伪分布式实战:从HDFS读写到MapReduce词频统计与ZooKeeper整合

Hadoop伪分布式实战:从HDFS读写到MapReduce词频统计与ZooKeeper整合 简介这份Hadoop简单应用案例面向大数据入门与进阶学习者围绕分布式存储与计算的核心链路帮助读者在真实项目中理解MapReduce、HDFS、Zookeeper与Hive的协作方式。资源以HadoopDemo-master项目为主体包含单词统计、Web日志分析、PageRank等典型实验场景并配套Hive建表查询与Zookeeper集群协调的实践内容。压缩包共89个文件以38个jar依赖、30个java源码、13个csv数据集为主另有gz压缩数据与工程配置文件整体约30MB目录按源码、数据集与依赖库分层组织便于按模块查阅与二次开发。目前已有1822人学习下载。通过该案例读者可掌握MapReduce编程模型、HDFS文件操作、日志指标提取以及Hive离线分析的基本流程积累从环境依赖到任务运行的排错经验适合作为大数据课程实验或自学练手项目。1. 从一台笔记本到伪分布式Hadoop简单应用案例到底能跑出什么很多人第一次接触 Hadoop卡在“装完不知道拿它干嘛”。伪分布式搭建教程看了一堆jps也能看到五个进程但真让你写个 MapReduce 跑一遍或者把本地文件塞进 HDFS 再读出来手就停了。这篇笔记就干一件事用一台普通开发机把 Hadoop 从安装配置到跑通一个完整小应用的全流程走一遍包括 HDFS 读写、MapReduce 词频统计、和 ZooKeeper 做一次简单整合最后说清楚哪些参数不能乱动、哪些坑我踩过。适合正在做课程设计、准备面试、或者想给日志处理找个入门方案的后端和运维同学。全程不需要多台机器伪分布式足够把核心机制跑明白。2. 伪分布式环境搭建从下载安装到五个进程全部起来2.1 为什么先跑伪分布式而不是真集群真集群多节点当然更接近生产但对“简单应用案例”这个目标来说伪分布式有三个不可替代的好处。第一NameNode、DataNode、ResourceManager、NodeManager、SecondaryNameNode 全在一台机器上你能用jps一眼看到谁死了谁活着排查成本极低。第二HDFS 的块副本机制、MapReduce 的 shuffle 过程在伪分布式下和真集群完全一致只是数据量小。第三课程设计和面试题里问的“Hadoop HA”“集群搭建”底层依赖的配置文件项在伪分布式里已经全部出现过后面扩成真集群只是复制粘贴加改主机名。我一般建议先用伪分布式把 HDFS 写读、MapReduce 提交、ZooKeeper 协调这三件事跑通再去碰多节点。顺序反了你会把大量时间花在 SSH 免密和防火墙这种和 Hadoop 本身无关的事情上。2.2 安装与配置的最小命令集假设你用一台 Linux 机器Ubuntu 或 CentOS 都行JDK 用 8 或 11。Hadoop 下载安装教程网上很多核心就几步解压、配环境变量、改四个 XML 文件、格式化、启动。# 解压到 /opt 下目录名不要带版本号以外的空格 tar -zxvf hadoop-3.x.tar.gz -C /opt/ mv /opt/hadoop-3.x /opt/hadoop # 配置环境变量追加到 ~/.bashrc echo export HADOOP_HOME/opt/hadoop ~/.bashrc echo export PATH$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin ~/.bashrc echo export JAVA_HOME/usr/lib/jvm/java-8-openjdk-amd64 ~/.bashrc source ~/.bashrc # 验证 hadoop version上面这段的逻辑是Hadoop 的启动脚本依赖HADOOP_HOME和JAVA_HOME两个都不能少。hadoop version能输出版本号说明环境变量生效。如果报 “JAVA_HOME is not set”先确认echo $JAVA_HOME有没有值再检查hadoop-env.sh里是否显式指定了 Java 路径——这是新手第一个高频翻车点。接下来改四个核心文件路径在$HADOOP_HOME/etc/hadoop/下。core-site.xmlconfiguration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property property namehadoop.tmp.dir/name value/opt/hadoop/data/tmp/value /property /configurationfs.defaultFS告诉客户端 NameNode 在哪伪分布式就写 localhost。hadoop.tmp.dir是所有临时数据的根目录必须手动创建否则格式化时会报目录不存在。这个参数很多人忽略结果 NameNode 一重启数据就丢。hdfs-site.xmlconfiguration property namedfs.replication/name value1/value /property property namedfs.namenode.name.dir/name value/opt/hadoop/data/namenode/value /property property namedfs.datanode.data.dir/name value/opt/hadoop/data/datanode/value /property /configuration伪分布式只有一台机器副本数必须设成 1设成 3 会一直卡在副本不足的告警里。dfs.namenode.name.dir和dfs.datanode.data.dir分别指定元数据和块数据的存放路径同样要提前建好目录。mapred-site.xml和yarn-site.xml分别指定 MapReduce 跑在 YARN 上、以及 ResourceManager 的主机名。这两个文件的内容比较固定按官方模板改yarn.resourcemanager.hostname为 localhost 即可。配置完成后格式化并启动hdfs namenode -format start-dfs.sh start-yarn.sh jpsjps应该看到 NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager 五个进程。少一个就去$HADOOP_HOME/logs/下看对应日志90% 的问题是目录权限或端口占用。提示格式化只能做一次。如果反复格式化DataNode 的 clusterID 会和 NameNode 对不上表现为 DataNode 启动后立刻退出。解决办法是删掉所有 data 目录重新格式化或者手动同步 clusterID。3. HDFS 读写与 MapReduce 词频统计一个能交作业的完整案例3.1 用 HDFS 命令完成一次上传和读取装好之后第一个应用案例我建议从 HDFS 开始因为它最直观。目标把本地一个文本文件上传到 HDFS再读出来确认内容一致。# 在 HDFS 上建目录 hdfs dfs -mkdir -p /user/root/input # 上传本地文件 hdfs dfs -put ./words.txt /user/root/input/ # 查看文件列表和内容 hdfs dfs -ls /user/root/input/ hdfs dfs -cat /user/root/input/words.txt # 从 HDFS 下载回本地 hdfs dfs -get /user/root/input/words.txt ./words_copy.txt这几条命令覆盖了 HDFS 最常用的操作。-mkdir -p支持多级目录-put上传-get下载-cat直接输出内容。注意-put和-copyFromLocal在大多数场景下等价但-put可以从标准输入读更灵活。这里有个容易踩的坑如果你在core-site.xml里把fs.defaultFS配成了hdfs://localhost:9000那hdfs dfs -ls /访问的就是 HDFS 根目录如果你临时想操作本地文件系统要写hdfs dfs -ls file:///。不写file://前缀命令会默认走 HDFS找不到文件时你会怀疑人生。3.2 MapReduce 词频统计代码、打包、提交三步走词频统计是 Hadoop 的 “Hello World”但很多人只抄过代码没自己提交过。完整流程是写 Mapper 和 Reducer、打包成 jar、用hadoop jar提交。// WordCountMapper.java import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; import java.io.IOException; import java.util.StringTokenizer; 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 { // 按空格切分每一行每个词输出 (word, 1) StringTokenizer itr new StringTokenizer(value.toString()); while (itr.hasMoreTokens()) { word.set(itr.nextToken()); context.write(word, one); } } }Mapper 的四个泛型分别是输入 key、输入 value、输出 key、输出 value。输入 key 是行偏移量LongWritable输入 value 是行内容Text。context.write每调用一次就产生一个中间键值对这些键值对会经过 shuffle 按 key 分组后送给 Reducer。// WordCountReducer.java import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; import java.io.IOException; 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; // 同一个 key 的所有 value 累加 for (IntWritable val : values) { sum val.get(); } result.set(sum); context.write(key, result); } }Reducer 收到的 value 是一个迭代器里面是所有相同 key 的 1。累加后输出(word, 总次数)。注意 Reducer 的输入类型必须和 Mapper 的输出类型一致否则运行时会报类型不匹配。驱动类里设置 Job 的输入输出路径、Mapper 和 Reducer 类然后打包# 编译打包假设 Hadoop 的 jar 包已在 classpath 中 javac -classpath hadoop classpath -d classes *.java jar -cvf wordcount.jar -C classes/ . # 提交到 YARN 运行 hadoop jar wordcount.jar WordCountDriver /user/root/input /user/root/output # 查看结果 hdfs dfs -cat /user/root/output/part-r-00000hadoop classpath会自动把 Hadoop 所有依赖 jar 拼成 classpath省去手动找包的麻烦。提交时输出目录必须不存在否则 Job 会直接失败——这是 MapReduce 的保护机制防止你覆盖上一次的结果。跑完后part-r-00000就是结果文件前面的part-r-是固定前缀后面的数字是 reduce 任务编号。注意如果 Job 卡在 map 0% reduce 0% 不动先看 ResourceManager 的 Web UI默认 8088 端口里应用状态。常见原因是yarn-site.xml里yarn.nodemanager.aux-services没配成mapreduce_shuffle导致 shuffle 阶段无法启动。4. Hadoop 和 ZooKeeper 整合实战用分布式锁协调两个客户端4.1 为什么 Hadoop 应用里会用到 ZooKeeperHadoop 本身的高可用HA就依赖 ZooKeeper 做 Active/Standby NameNode 的选举。但在“简单应用案例”层面更常见的需求是多个客户端同时往 HDFS 写同一个目录时怎么保证不冲突。比如两个采集程序都要往/user/root/logs/下写文件文件名如果都按时间戳生成同一秒就可能撞车。这时候用 ZooKeeper 做一个分布式锁谁先拿到锁谁先写写完释放另一个再写。ZooKeeper 的临时顺序节点天然适合做这个。每个客户端在锁目录下创建一个临时顺序节点然后判断自己是不是序号最小的那个是就获得锁不是就监听前一个节点。这个模式比轮询优雅得多也是面试题里常问的“ZooKeeper 实现分布式锁”的标准答案。4.2 整合步骤与关键代码先确保 ZooKeeper 已经跑起来单机模式即可然后在 Java 项目里引入zookeeper和hadoop-client依赖。核心逻辑分两步抢锁、写 HDFS。// DistributedLock.java 核心片段 import org.apache.zookeeper.*; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; import java.io.IOException; import java.util.Collections; import java.util.List; public class DistributedLock { private ZooKeeper zk; private String lockPath /hadoop-lock; private String currentNode; public DistributedLock() throws IOException { // 连接 ZooKeeper会话超时 3000ms zk new ZooKeeper(localhost:2181, 3000, event - {}); } public void acquire() throws Exception { // 创建临时顺序节点 currentNode zk.create(lockPath /lock-, new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL_SEQUENTIAL); while (true) { ListString children zk.getChildren(lockPath, false); Collections.sort(children); String smallest children.get(0); if (currentNode.endsWith(smallest)) { return; // 拿到锁 } // 监听前一个节点 String prev children.get(Collections.binarySearch(children, currentNode.substring(lockPath.length() 1)) - 1); zk.exists(lockPath / prev, true); Thread.sleep(100); } } public void release() throws Exception { zk.delete(currentNode, -1); zk.close(); } }create的第四个参数EPHEMERAL_SEQUENTIAL表示临时顺序节点会话断开自动删除且节点名带自增序号。getChildren拿到所有子节点后排序序号最小的持有锁。zk.exists注册的 watcher 会在前一个节点被删除时触发当前客户端被唤醒后重新检查自己是不是最小。拿到锁之后写 HDFS 就简单了Configuration conf new Configuration(); conf.set(fs.defaultFS, hdfs://localhost:9000); FileSystem fs FileSystem.get(conf); Path out new Path(/user/root/logs/ System.currentTimeMillis() .log); FSDataOutputStream os fs.create(out); os.writeBytes(some log content\n); os.close(); fs.close();FileSystem.get(conf)会根据fs.defaultFS返回对应的文件系统实例。fs.create创建文件时如果父目录不存在会报错所以要么提前fs.mkdirs要么确保目录已存在。写完记得关流和关文件系统否则连接泄漏跑久了会报 “Too many open files”。提示ZooKeeper 的会话超时时间不要设太短。3000ms 在本地够用但如果你的应用有 GC 停顿可能还没释放锁会话就过期了导致锁被提前释放。生产环境一般设 10 到 30 秒。5. 避坑与排查伪分布式和整合阶段最容易翻车的五件事现象一jps看不到 DataNode但日志里没有明显报错。原因多次执行hdfs namenode -format后NameNode 的 clusterID 变了DataNode 启动时发现自己的 clusterID 对不上直接退出。 解决停掉所有进程删除dfs.namenode.name.dir和dfs.datanode.data.dir下的所有内容重新格式化一次再启动。以后记住格式化只做一次。现象二MapReduce 任务一直卡在 map 0%ResourceManager 页面显示应用处于 ACCEPTED 状态。原因YARN 的 NodeManager 没有正常注册或者yarn.nodemanager.resource.memory-mb设得比机器实际内存还大导致容器无法分配。 解决检查yarn-site.xml里的内存配置伪分布式下设成 2048 或 4096 即可。同时确认mapred-site.xml里mapreduce.framework.name是yarn。现象三HDFS 上传文件时报 “Could only be replicated to 0 nodes”。原因DataNode 没起来或者dfs.replication设成了大于实际 DataNode 数量的值。 解决先jps确认 DataNode 在再把dfs.replication改成 1。伪分布式下副本数只能是 1。现象四ZooKeeper 锁偶尔失效两个客户端同时写入了。原因客户端 GC 停顿超过会话超时时间ZooKeeper 认为会话已死临时节点被删除锁自动释放但客户端自己不知道还在写。 解决在写 HDFS 前再检查一次自己是否还持有锁比如判断节点是否存在或者把会话超时调大。更稳妥的做法是用 Curator 框架的InterProcessMutex它帮你处理了这些边界。现象五hadoop jar提交时报 ClassNotFoundException。原因打包时没有把依赖的 class 打进去或者驱动类的全限定名写错了。 解决用jar -tf wordcount.jar确认驱动类在包里提交时写全限定类名如com.example.WordCountDriver。如果依赖第三方库用 Maven 的 shade 插件打 fat jar。6. 进阶技巧用 distcp 做一次跨目录迁移顺便验证集群健康度跑通词频统计和 ZooKeeper 锁之后我建议你拿distcp做一次收尾练习。distcp是 Hadoop 自带的分布式复制工具底层就是 MapReduce 任务能并行拷贝大量数据。用它有两个好处一是熟悉常用参数二是如果distcp能跑通说明你的 HDFS 和 YARN 都是健康的。最基础的用法hadoop distcp hdfs://localhost:9000/user/root/input \ hdfs://localhost:9000/user/root/backup这条命令把 input 目录整个复制到 backup 下。注意目标路径如果不存在distcp 会创建如果存在同名目录默认会报错除非加-overwrite。几个我常用的参数参数作用什么时候用-m 2指定最多 2 个 map 任务小集群防止任务过多抢资源-update只拷贝源和目标不一致的文件增量同步-delete删除目标端源端没有的文件保持两端完全一致-log /tmp/distcp.log记录日志排查失败任务比如做增量同步hadoop distcp -update -m 2 \ hdfs://localhost:9000/user/root/input \ hdfs://localhost:9000/user/root/backup-update会对比文件大小和修改时间只传变化的。-m 2限制并发 map 数伪分布式下别设太大否则内存不够会 OOM。验证方法跑完distcp后用hdfs dfs -ls -R对比两个目录的文件列表和大小再用hdfs dfs -cat抽查几个文件内容。如果distcp报 “Failed to get block locations”说明 NameNode 或 DataNode 有问题回头查jps和日志。我自己的习惯是每次改完 Hadoop 配置先跑一遍distcp小目录再跑词频统计。这两个都过了才认为环境是稳的。这个习惯帮我省了很多“以为配好了结果提交任务才发现问题”的时间。希望帮到你。本文还有配套的精品资源点击获取
返回列表