基于Hadoop的网站日志分析:MapReduce全流程项目拆解
发布时间:2026/9/26 2:52:52来源:尧图网络
简介基于Hadoop的网站日志分析程序是一份面向大数据初学者与人工智能应用开发者的实战源码包演示如何借助Hadoop MapReduce对网站日志进行分布式处理与分析。资源共14个文件以7个Java源程序与7个编译后的class文件组成包含完整的Map与Reduce逻辑便于直接阅读或运行验证压缩包仅16KB非常轻量。已有170人学习下载。项目完整覆盖日志解析、数据清洗、访问统计、用户行为挖掘等典型环节从日志格式解析、无效数据过滤到访问频次统计与用户路径还原各阶段均有清晰代码对应。结合人工智能场景分析得到的用户偏好与访问特征可作为训练推荐系统或构建用户画像的数据基础项目还展示了如何通过分区与Combiner优化作业性能适合希望从代码层面掌握大数据分析与AI数据预处理流程的读者。1. 网站日志分析为什么值得用 Hadoop 专门做一遍很多人在学 Hadoop 时都会卡在同一个问题官方 WordCount 跑完了HDFS 命令也会敲了但真要自己写一个分析程序却不知道从哪下手。这份“基于 Hadoop 的网站日志分析程序”正好补上这段空白——它是一个完整的、能直接跑通的日志分析项目覆盖了从日志采集、清洗、统计到结果导出的全流程而不是那种只演示 API 的碎片 Demo。如果你正在准备课程设计、毕业设计或者想搞清楚 MapReduce 在真实场景里到底怎么组织代码这个资源值得花时间拆一遍。我拆完这个项目后最大的感受是日志分析是最适合入门 Hadoop 的业务场景。日志是天然的大数据来源格式半结构化、数据量大、处理逻辑明确MapReduce 的优势能完整发挥。这份资源里包含了完整的 Java 工程代码、模拟日志生成器和详细的环境配置说明拿到手可以直接导入 IDEA 运行也可以根据自己的数据改改正则就能用对新手和想快速上手 Hadoop 开发的人都适用。2. 日志数据长什么样先搞定格式再谈分析2.1 网站日志的通用格式与字段拆解这个项目里处理的日志是典型的 Web 服务器访问日志。不管你是用 Tomcat、Nginx 还是 Apache落盘的日志格式基本一致。项目里默认按 Apache 的 Combined Log Format 来模拟数据每一行长这样127.0.0.1 - - [10/Oct/2000:13:55:36 -0700] GET /apache_pb.gif HTTP/1.0 200 2326 http://www.example.com/start.html Mozilla/4.08 [en] (Win98; I ;Nav)用空格和引号做分隔整行能拆成 10 个字段客户端 IP、两个占位符通常为空指 RFC 1413 身份和认证用户、时间戳、请求行方法 资源路径 协议、状态码、返回字节数、Referer、User-Agent。在写解析代码之前有一点我得先强调别急着写正则先把分隔符弄清楚。这个格式里既有空格分隔又有[]和包裹的特殊字段普通split( )会把时间戳和请求行拆得乱七八糟。项目里采用的是分层解析——先用空格粗分再对特殊字段做二次处理这比一股脑上一个复杂的正则表达式要稳得多。2.2 定义日志 Bean让 Map 阶段处理更有条理解析完的日志字段建议封装成一个 Java Bean。项目里的LogBean类是这么设计的public class LogBean implements Writable { private String ip; // 客户端 IP private String time; // 访问时间 [10/Oct/2000:13:55:36 -0700] private String method; // 请求方法 GET/POST private String path; // 请求资源路径 private String status; // 状态码 200/404/500 private String bytes; // 返回字节数 private String referer; // 来源页 private String userAgent; // 用户代理 // 必须提供无参构造反序列化时反射调用 public LogBean() {} // 实现 write 和 readFields序列化顺序必须与字段声明顺序一致 Override public void write(DataOutput out) throws IOException { out.writeUTF(ip); out.writeUTF(time); // ... 其余字段 } Override public void readFields(DataInput in) throws IOException { this.ip in.readUTF(); this.time in.readUTF(); // ... 其余字段 } // getter/setter 省略 }这里有几个关键点需要说明。第一必须实现 Writable 接口因为 MapReduce 的 Shuffle 阶段要把对象序列化成字节流在节点间传输Java 自带的Serializable在这个场景下性能太差框架不认。第二write 和 readFields 里的字段顺序必须完全一致这俩方法本质上是同一份字节流的写入和读出顺序错了轻则字段错乱重则直接反序列化失败。第三一定要保留无参构造器框架通过反射创建对象时依赖它你要是只写了带参构造运行到一半会报NoSuchMethodException。2.3 手写解析器不用正则用分段切分解析这块我直接给出一段能在项目里落地的代码。它的设计思路是先把整行日志按空格拆成数组再利用特殊字段的边界符号做二次提取public static LogBean parse(String line) { LogBean bean new LogBean(); // 第一层按空格切分典型情况能分出 10 段 String[] parts line.split( ); if (parts.length 10) { return null; // 格式不完整直接丢弃 } bean.setIp(parts[0]); bean.setMethod(parts[5].substring(1)); // GET 去掉引号 bean.setPath(parts[6]); // 资源路径 bean.setStatus(parts[8]); bean.setBytes(parts[9]); // 第二层提取 [10/Oct/2000:13:55:36 -0700] 中的时间部分 String timePart parts[3].substring(1); // 去掉左方括号 bean.setTime(timePart); // 第三层Referer 和 User-Agent 被引号包裹从原始行截取 int refererStart line.indexOf(\, line.indexOf(\, line.indexOf(\) 1) 1) 1; int refererEnd line.indexOf(\, refererStart); bean.setReferer(line.substring(refererStart, refererEnd)); // User-Agent 是最后一段引号内容 int uaStart line.lastIndexOf(\) 1; int uaEnd line.lastIndexOf(\, uaStart - 1); bean.setUserAgent(line.substring(uaEnd 1, uaStart - 1)); return bean; }逻辑说明第一层空格拆分能覆盖大部分字段因为 IP、状态码、字节数这些都不会被特殊符号干扰第二层时间戳处理把[去掉保留可读时间串第三层用引号定位把 Referer 和 User-Agent 取出来因为这两个字段内部可能包含空格用split必然出错。整个方法的时间复杂度是 O(n)只扫描两遍字符串在大数据量下效率是可以接受的。参数说明如果日志格式不是 Combined 格式比如 Nginx 默认的 combined 格式里可能没有 Referer 或 User-Agent那就需要调整切分逻辑。最稳妥的做法是先跑一次这个 parse 方法打印出parts.length确认实际分段数再改代码。我在拆这个项目时第一次直接照搬别人的解析器跑了满屏 null后来学着先数据探查再写死下标就再没翻过车。3. 环境准备与项目导入把 Hadoop 跑起来才是第一步3.1 开发环境怎么选本机装还是 Docker 起这个项目对 Hadoop 环境的要求不算高Hadoop 2.x 或 3.x 都兼容。但在环境选择上我建议分情况处理情况一电脑内存 16G 以上直接装伪分布式伪分布式模式指的是所有 Hadoop 守护进程NameNode、DataNode、ResourceManager、NodeManager都跑到本机 JVM 上能完整模拟真实集群的文件存储和资源调度行为重点是这个模式下能看到全部日志输出排查问题比 Docker 方便不少。情况二内存只有 8G用 Docker 跑单节点集群伪分布式的守护进程加起来会吃 3G 左右内存8G 机器再开个 IDEA 会卡成 PPT。用 Docker 起一个 Hadoop 镜像宿主机只占一个轻量容器的开销。常见做法是拉一个单节点镜像启动时把 9870NameNode UI和 8088YARN UI端口映射出来。我自己的习惯是开发阶段用 Docker 起 HadoopIDEA 里跑代码两者通过网络连接等到要调优性能或看任务监控页面时再把代码打成 Jar 包丢进容器里执行。3.2 项目目录结构与导入要点这个资源解压后目录结构大致长这样web-log-analysis/ ├── src/main/java/ │ ├── bean/ # LogBean 等 Writable 自定义类型 │ ├── mapreduce/ # 各个分析任务的 Mapper、Reducer │ └── utils/ # 日志解析、时间工具 ├── input/ │ └── access.log # 模拟生成的日志文件 ├── output/ # 存放分析结果的目录 ├── pom.xml # Maven 依赖管理 └── README.md # 环境配置和运行说明直接 IDEA 导入 Maven 项目pom.xml 里的依赖主要是 hadoop-client我拆的这个版本是这样的dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version3.3.4/version /dependency注意这个依赖要把 Hadoop 的常用包全部带进来包括 HDFS 客户端、MapReduce 客户端和 YARN 相关依赖不用额外再引。如果你的机器本机装了 Hadoop 并且配置了HADOOP_HOMEIDE 里可以直接跑如果是连远程集群则要在工程里把core-site.xml和hdfs-site.xml放到 resources 目录下指定fs.defaultFS为远程 NameNode 地址configuration property namefs.defaultFS/name valuehdfs://192.168.1.100:9000/value /property /configuration这里踩坑概率最高的是hadoop-client的版本和集群版本不一致。版本不匹配会引出各种奇怪的 RPC 协议错误比如Failed to connect to /192.168.1.100:9000。第一次跑之前请先确认你连的 Hadoop 是什么版本让 pom 里的依赖版本尽量等于或小于集群版本。3.3 把本地日志上传到 HDFS三条命令搞定日志分析的数据源在 HDFS 里MT 项目里自然也会给你准备模拟日志生成器。数据准备好后上传命令是# 在 HDFS 根目录建一个 log 数据目录 hdfs dfs -mkdir -p /user/hadoop/log/input # 把本地生成的日志文件传到 HDFS hdfs dfs -put ./input/access.log /user/hadoop/log/input/ # 确认上传成功查看文件大小和块信息 hdfs dfs -ls /user/hadoop/log/input/ hdfs dfs -stat %b /user/hadoop/log/access.log参数说明-mkdir -p和 Linux 同名命令语义一致会自动创建缺失的父目录。-put适合一条条上传如果你的日志是按天生成的多个文件建议改用-put整个目录或使用-appendToFile增量追加。查看文件块信息用hdfs fsck /user/hadoop/log/access.log -files -blocks能看到这个日志被分成了几块、每块在哪台节点上对后续理解 MapReduce 的输入分片会有帮助。提示数据上传成功后建议先随手跑一下hdfs dfs -text /user/hadoop/log/input/access.log | head -5确认 HDFS 里的文件内容和本地一致防止上传时文件损坏或编码被改。4. 核心统计怎么算四个 MapReduce 任务的代码拆解4.1 PV 统计最简单的计数任务页面访问量PV是日志分析最基础的指标意思是每个页面被请求了多少次。这个任务的 Map 阶段把页面路径作为 keyReducer 阶段累加计数public class PVMapper extends MapperLongWritable, Text, Text, IntWritable { private final static IntWritable one new IntWritable(1); private Text path new Text(); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { LogBean bean LogParser.parse(value.toString()); if (bean ! null) { path.set(bean.getPath()); context.write(path, one); } } } public class PVReducer 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); } }逻辑说明Map 阶段每读入一行日志解析出路径字段直接输出中间的 Shuffle 阶段框架会把相同路径的记录自动集中到同一个 ReducerReducer 里对每个路径的IntWritable列表累加就得到该路径的总访问量。整个流程没有复杂的业务判断是 MapReduce 里最标准的“计数模型”。需要注意的是bean ! null这个判空它对应解析器里parts.length 10就返回 null 的逻辑。脏数据一定要在 Map 阶段过滤掉不能留到 Reducer否则一条格式错误的数据会导致整个任务失败。4.2 UV 统计用 Set 保持独立性独立访客数UV和 PV 的差别在于同一 IP 访问多次只算一次。MapReduce 里最直接的做法是利用 Reduce 阶段的去重特性——相同 key 会被分到同一个 Reducer在 Reducer 里用 Set 存 IP 再取 sizepublic class UVReducer extends ReducerText, Text, Text, IntWritable { private SetString ipSet new HashSet(); private IntWritable result new IntWritable(); Override protected void reduce(Text key, IterableText values, Context context) throws IOException, InterruptedException { ipSet.clear(); for (Text val : values) { ipSet.add(val.toString()); } result.set(ipSet.size()); context.write(key, result); } }逻辑说明这里的 key 是访问日期比如2024-06-01value 是每条日志里的 IP。Reducer 收到的 values 是去重前的全部 IP在 Reducer 内部用 Set 去重再取 size就能算出当日 UV。这个方案的缺点是内存占用和日期内 IP 数成正比但在单日百万级日志的场景下完全够用。如果日志量进一步上量更优的做法是“两阶段去重”先做一次Map 输出 IP 为 keyvalue 为空的 MapReduce 任务完成全局去重再去统计数量。但项目里用不到我就不展开说了。4.3 热门页面 Top N排序的两种实现方式Top N 在日志分析里特别常见比如看访问量最高的 10 个页面。实现方式有两种这个项目里用的是第二种方式一Reducer 内排序Map 阶段输出(path, 1)Reducer 里把所有数据先存进一个 Map最后统一排序取前 N 个。这个方法代码简单但 Reducer 内存压力大数据量超过千万级容易 OOM。方式二利用 Shuffle 的默认排序 自定义比较器MapReduce 框架在 Shuffle 阶段会对 key 做排序这一点是被很多人忽略的“隐藏功能”。利用这个机制我们让path做不了排序条件就把path和访问次数拼成一个组合 key实现一个自定义WritableComparable让框架先按访问次数降序排再取前 Npublic class CountKey implements WritableComparableCountKey { private String path; private long count; Override public int compareTo(CountKey o) { // 访问次数降序相同再按路径字典序 return Long.compare(o.count, this.count) ! 0 ? Long.compare(o.count, this.count) : this.path.compareTo(o.path); } Override public void write(DataOutput out) throws IOException { out.writeUTF(path); out.writeLong(count); } Override public void readFields(DataInput in) throws IOException { this.path in.readUTF(); this.count in.readLong(); } }逻辑说明这个类既是 key又包含排序规则。Map 阶段输出(CountKey, NullWritable)框架在分区和排序时会调用compareTo方法所有记录会先按访问次数降序排列。之后写一个 Reducer只取前 N 条输出就拿到了 Top N 页面。这个做法的核心价值是排序发生在框架内部而不是我们代码里它能把内存压力分散到多个 Reducer 节点上。但要注意只有指定了job.setNumReduceTasks(1)时全局排序才成立Reducer 数量大于 1 时数据是分区间排序的每个 Reducer 只能拿到局部的 Top N。4.4 独立 IP 数及其按小时分布独立 IP 和按小时统计 PV 这两个指标它们的共同点是都需要做时间维度处理。Map 阶段把时间戳里的小时提取出来作为 key 的一部分public class HourMapper extends MapperLongWritable, Text, Text, IntWritable { private Text hourKey new Text(); private final static IntWritable one new IntWritable(1); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { LogBean bean LogParser.parse(value.toString()); if (bean null) return; // 时间字段格式10/Oct/2000:13:55:36 String time bean.getTime(); String hour time.substring(time.indexOf(:) 1, time.indexOf(:) 3); hourKey.set(hour); context.write(hourKey, one); } }这里提取小时的逻辑是原始时间10/Oct/2000:13:55:36第一个冒号后面跟的就是小时13。indexOf(:)找到第一个冒号的位置加 1 跳到小时首位加 3 取到小时结束。这套写法避免了把字符串拆成数组再拼接的开销直接按固定偏移量截取。4.5 作业驱动的标准写法Job 配置参数详解有了 Mapper 和 Reducer还需要一个 Driver 类把所有东西串起来。项目里的 Job 配置是标准写法public static void main(String[] args) throws Exception { Configuration conf new Configuration(); Job job Job.getInstance(conf, pv-statistics); job.setJarByClass(PVJob.class); job.setMapperClass(PVMapper.class); job.setCombinerClass(PVReducer.class); // 能减则减减少 Shuffle 数据量 job.setReducerClass(PVReducer.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); }参数说明setCombinerClass是优化 Shuffle 的关键它的逻辑和 Reducer 完全相同——在本节点先做一次本地聚合把多条(path,1)合并成(path,n)这样网络传输的数据量能减少好几倍。但要注意 Combiner 的输入输出类型必须和 Mapper 的输出类型一致。在这个例子里 Mapper 输出Text, IntWritableCombiner 和 Reducer 输入输出也都是这两个类型所以能直接复用 Reducer 类。如果遇到求平均值的业务Combiner 就不能直接用 Reducer 的逻辑了因为平均值不可叠加需要单独写。job.waitForCompletion(true)的 boolean 参数表示是否打印任务进度。这个参数在生产环境我一般设成 true因为我们调试时特别需要看 Map 和 Reduce 的百分比进度但在上百个作业串联执行的场景建议设成 false 并把日志切到 WARN否则控制台会被刷爆。4.6 打包提交课程设计里最常见的翻车点IDEA 里跑通了不算完很多时候你需要把工程打成 Jar 包丢到服务器上运行。这个项目里 Maven 打包配置是这样的build plugins plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-shade-plugin/artifactId version3.2.4/version executions execution phasepackage/phase goals goalshade/goal /goals configuration transformers transformer implementationorg.apache.maven.plugins.shade.resource.ManifestResourceTransformer mainClasscom.example.PVJob/mainClass /transformer /transformers /configuration /execution /executions /plugin /plugins /build打包后提交到集群执行hadoop jar web-log-analysis-1.0.jar /user/hadoop/log/input /user/hadoop/log/output/pv这里maven-shade-plugin的作用是把依赖的 Hadoop 相关类一并打进 Jar 包防止提交时 classpath 缺失。如果你用普通的maven-jar-plugin经常会碰到ClassNotFoundException: org.apache.hadoop.conf.Configuration因为集群上的 Hadoop 环境和你本地的不一样。注意 shade 插件里必须配置 MainClass否则提交时只能写全类名hadoop jar xxx.jar cn.itcast.mapreduce.PVJob。5. 避坑笔记拆完这个项目后最想告诉你的五件事5.1 日志文件上传后没有换行符所有记录变成一行现象HDFS 里的文件看起来正常但跑 MapReduce 时发现 Map 阶段只处理了一条记录Reduce 输出结果数值全是 1。原因Windows 环境下生成的日志文件换行符是\r\n而 Linux 和 Hadoop 默认按\n切分。\r留在行尾导致部分行的解析失败或字段混入不可见字符。解决上传前统一转换格式。最简单的方式是用 IDEA 右下角把文件编码和换行符都切成 LF或者执行一条流式命令sed -i s/\r$// access.log另外写日志生成器时也要注意用\n拼接而不是System.lineSeparator()——后者在 Windows 下会产生\r\n。我在拆这个项目时用模拟生成器重新生成了一版日志才彻底绕开这个问题。5.2 提交到 YARN 后一直卡在 ACCEPTED 状态任务不启动现象IDEA 本地跑没问题用hadoop jar提交到集群后日志显示ACCEPTED但 ApplicationMaster 迟迟没有启动几分钟后报TimedOut waiting for AM to start。原因伪分布式集群的内存资源配置不足。YARN 默认给 ApplicationMaster 分配 1G 内存和 1 个虚拟核但你的yarn-site.xml里yarn.scheduler.maximum-allocation-mb设置的数值小于这个默认值资源调度器无法满足 AM 的申请。解决改yarn-site.xml把资源参数调大并重启 YARNproperty nameyarn.scheduler.maximum-allocation-mb/name value2048/value /property property nameyarn.scheduler.minimum-allocation-mb/name value256/value /property property nameyarn.nodemanager.resource.memory-mb/name value2048/value /property改完后执行mapred --daemon restart resourcemanager或直接重启整个集群。这个坑在 2G 内存的小集群上几乎是必现的建议不管资源够不够先把这两个参数按上面配好。5.3 Reduce 阶段大量 OOM堆内存设置失效现象任务在 Reduce 阶段报java.lang.OutOfMemoryError: Java heap space日志里能看到GC overhead limit exceeded。原因Hadoop 默认给 Map 和 Reduce 的子进程分配的堆内存是 1G 左右在mapred-site.xml里直接改mapreduce.reduce.java.opts有时不生效因为 Jar 包里自带的mapred-site.xml覆盖了集群的配置。解决在这个项目里最有效的做法是提交作业时通过-D参数显式指定hadoop jar web-log-analysis-1.0.jar \ -D mapreduce.reduce.memory.mb2048 \ -D mapreduce.reduce.java.opts-Xmx1536m \ /user/hadoop/log/input /user/hadoop/log/output/pv-D参数会被 JobConf 读取并覆盖集群配置优先级高于所有配置文件。reduce.memory.mb控制 YARN 分配给 Reduce 容器的总内存reduce.java.opts控制 JVM 实际使用的堆上限后者必须小于前者一般取前者的 75% 左右。如果你希望一劳永逸就把这两个参数写进集群的mapred-site.xml但提交时记得删掉 Jar 包内 resources 目录里的同名文件。5.4 结果文件里的中文变成了乱码现象Reduce 输出结果写入 HDFS 后用hdfs dfs -cat查看正常但下载到本地打开中文全部显示为??或方块。原因Hadoop 默认的输出编码是 UTF-8问题多半出在 Windows 本地查看工具不支持 UTF-8或者本地的core-site.xml里io.serializations配置了旧的编码器。解决先确认 HDFS 里的文件确实是 UTF-8。执行hdfs dfs -getmerge /user/hadoop/log/output /tmp/result.txt用 IDEA 打开文件右下角把文件编码显式切换成 UTF-8。如果确认 HDFS 端就是乱码则改mapred-site.xml里的mapreduce.output.fileoutputformat.compress.codec换成property namemapreduce.output.fileoutputformat.compress.codec/name valueorg.apache.hadoop.io.compress.GzipCodec/value /property property namemapreduce.output.fileoutputformat.compress/name valuetrue/value /property开启 Gzip 压缩能避免 HDFS 写入时的编码转换问题同时也顺带减小了结果文件体积。但这个方案会引入新坑——下游程序如果直接读结果文件需要支持 gzip 解压。所以如果没有读取兼容性要求直接用 UTF-8 原样输出通过本地 IDEA 打开就能看到正确中文。5.5 任务运行成功但输出目录为空或者只有_SUCCESS现象MapReduce 作业显示SUCCESS但输出目录里只有part-r-00000且内容是空的甚至只有_SUCCESS标记文件。原因这是最隐蔽的“伪成功”。最常见的原因是 Mapper 里的解析器把所有行都判成 null 了数据处理量等于 0Reducer 收到空列表自然没有输出。解决先在 Map 阶段加一个计数器验证输入行的解析率。// 在 Mapper 的 setup 或 map 里定义 enum ParserCounter { PARSED, SKIPPED } // map 方法里统计 if (bean ! null) { context.getCounter(ParserCounter.PARSED).increment(1); context.write(path, one); } else { context.getCounter(ParserCounter.SKIPPED).increment(1); }任务跑完看计数器hadoop jar web-log-analysis-1.0.jar /user/hadoop/log/input /user/hadoop/log/output/pv # 在控制台输出的最后几行会显示 # ParserCounter PARSED100000 SKIPPED5000如果SKIPPED数量很大直接拿一条原始日志到本地跑解析器打印parts.length根据实际分段数调整代码。这个调试思路在拆别人的项目时特别好用——先确认输入没问题再查业务逻辑别一上来就怀疑集群配置那是玄学调优浪费时间。6. 结果可视化验证从 HDFS 取数到前端图表一个不跑偏的闭环MapReduce 任务跑完只是第一步结果得让人看得懂才算闭环。这个项目的 output 目录里的结果文件是纯文本直接在控制台hdfs dfs -cat露一眼能看但看着不直观。我会把结果导出后用 Python 快速出图或者用 Hive 做二次关联分析验证数据的正确性。导出的标准流程是这样的# 把多个 part 文件合并成一个本地文件 hdfs dfs -getmerge /user/hadoop/log/output/pv /tmp/pv_result.txt # 确认合并结果看前 20 行 head -20 /tmp/pv_result.txt-getmerge会把part-r-00000、part-r-00001等文件按顺序拼成一个文件这个命令比-get一个个拉要省事得多。拿到合并后的文本用 Python 三行代码就能画个柱状图import matplotlib.pyplot as plt # 读取 MapReduce 的结果 path_counts {} with open(/tmp/pv_result.txt, r) as f: for line in f: path, count line.strip().split(\t) path_counts[path] int(count) # 按访问量降序取前 10 个页面 top10 sorted(path_counts.items(), keylambda x: x[1], reverseTrue)[:10] plt.figure(figsize(10, 6)) plt.bar([p[0] for p in top10], [p[1] for p in top10]) plt.xticks(rotation45) plt.title(Top 10 Popular Pages) plt.tight_layout() plt.savefig(/tmp/top_pages.png)这段脚本读的字段分隔符是 Tab——MapReduce 的默认输出格式就是key\tvalue。如果你的 Reducer 里用context.write(key, result)输出的不是这个格式需要回第 4 章检查 OutputValueClass 是不是设成了Text如果你用了NullWritable或自定义类型导出后的解析规则就要跟着变。验证维度上我习惯做几次“比对式验证”。拿 PV 结果来说原始日志有 100 万行所有页面的 PV 加起来应该等于 100 万。跑个简单的 Hive 查询SELECT SUM(cnt) FROM pv_result;这里pv_result是外部分区表对应的 HDFS 路径就是/user/hadoop/log/output/pv。如果 SUM 结果和日志总行数对不上说明 Map 阶段过滤了不该过滤的数据或者解析逻辑漏了字段。我拆这个项目的过程中就曾因为日志解析器里parts.length 10判断条件写成了结果丢了一整天的数据后来就是这个 SUM 比对把问题揪出来的。这个“总数核对”的习惯我一直到现在还在用。从那以后我每次跑完一个 MapReduce 任务都会强制走一遍这三件事先看 Counter 里的解析率和处理行数再 getmerge 导出结果最后用一个全量聚合的 SQL 交叉验证。这套流程才真的能发现隐性丢数据而不是只看任务是不是绿了。希望帮到你。本文还有配套的精品资源点击获取
网站建设高端定制企业官网