
简介一个基于Hadoop的网站日志分析项目面向正在学习大数据处理与人工智能数据预处理的学生和开发者尤其适合数据科学、人工智能方向的初学者。它以网站日志为处理对象演示了从日志解析、数据清洗、访问量统计到热门页面排序、用户行为分析的完整流程并覆盖了无效日志过滤、异常访问识别等常见分析任务适合作为课程设计、毕业设计或Hadoop入门练习素材。压缩包共十四个文件主要包含七份Java源码和对应的七份字节码代码精简整体仅十六KB便于快速查看和直接运行。目前已有170人学习下载。通过该项目可以掌握Mapper与Reducer的实现方式、作业配置与提交方法了解如何将统计结果用于用户画像、推荐系统等人工智能场景同时学习HDFS文件读写与日志格式解析等基础操作为独立开发日志分析工具打下基础。1. 一个 zip 能装下的日志分析Hadoop 到底解决了什么问题网站日志分析这个需求一旦日均 PV 过了百万、日志文件每天几个 GB 起步传统的 awk 加 grep 加 Excel 流水线就彻底转不动了。基于Hadoop的网站日志分析程序.zip这个标题代表了一类非常典型的课程设计/生产级小工程把分散在 Nginx 或 Apache 服务器上的 access.log 统一收进 HDFS再用 MapReduce 或 Hive 跑出 PV、UV、独立 IP、TOP 页面这些指标。这类 zip 里的程序往往不大但踩过的坑一点不少从 Hadoop 伪分布式搭建、环境配置到集群跑通每一步都可能有反直觉的报错。这篇笔记就把我做过的最常见方案讲清楚适合刚搭好 Hadoop 伪分布式或准备做课程设计的同学也适合想把日志分析落到 YARN 上的从业者照着复现。2. 把日志送进 HDFS采集与落地的三种选择2.1 日志源头的格式陷阱先统一 Nginx 的 log_format做日志分析第一步不是写代码而是确定日志长什么样。Nginx 默认的 combined 格式已经够用但很多线上环境会自定义字段比如加上$request_time、$upstream_status或者$http_x_forwarded_for。如果直接拿默认解析器去切很容易错位。我一般会让运维把 log_format 固定成一套并在日志里预留 JSON 或者统一的字段分隔符。常见的 Nginx 配置是这样log_format main $remote_addr - $remote_user [$time_local] $request $status $body_bytes_sent $http_referer $http_user_agent $http_x_forwarded_for; access_log /var/log/nginx/access.log main;这里的关键点是日志里本身没有转义处理如果用户请求的 URL 里带引号或空格整行解析就会偏。实际处理时我的做法是先用nginx -t验证配置然后写一小段脚本读取样例行用空格做切分时注意$request里面是带引号的整体。对 MapReduce 程序而言日志一行就是一个 record解析责任全部落在 Mapper所以格式越规整越省事。2.2 上传 HDFS 的三条路Flume、定时脚本还是 Web 直传拿到日志文件后往 HDFS 传数据的方式会直接影响后续处理效率。常见做法有三种各有利弊。第一条路是 Flume。Flume 的 TailDir Source 可以实时监听本地文件新增行落到 HDFS 上时还能按天或者按小时滚动目录。配置大概长这样agent.sources tail agent.sources.tail.type TAILDIR agent.sources.tail.positionFile /opt/flume/position/taildir.json agent.sources.tail.filegroups f1 agent.sources.tail.filegroups.f1 /var/log/nginx/.*access.*\.log agent.sinks hdfsSink agent.sinks.hdfsSink.type hdfs agent.sinks.hdfsSink.hdfs.path hdfs://namenode:9000/logs/%Y%m%d/ agent.sinks.hdfsSink.hdfs.filePrefix access agent.sinks.hdfsSink.hdfs.rollInterval 3600这段配置里最容易被忽略的是positionFile。如果 Flume 重启后找不到上次读取位置它会把整个日志文件重新读一遍导致重复数据。rollInterval3600表示每小时滚动一次文件不设置的话小文件会一直累积对大集群来说很伤 NameNode 内存。第二条路是直接用hadoop fs -put配合 crontab。适合日志量不大、实时性要求不高的场景。命令很简单#!/bin/bash date_str$(date %Y%m%d) hadoop fs -mkdir -p /logs/$date_str hadoop fs -put /var/log/nginx/access.log /logs/$date_str/access_$(date %H%M).log第三条路是用 Web 服务接收日志再写 HDFS适合跨机房或者日志源特别分散的场景。这条路要自己做幂等控制否则重复写会很频繁。2.3 分区目录设计按天分区是所有后续查询的命根子HDFS 上的日志目录结构直接决定 Hive 分区的效率。我见过很多工程把日志一股脑丢进/logs/结果后面查某一天的数据要扫描全量文件。正确姿势是建多层分区至少年/月/日或者直接按%Y%m%d扁平存放。推荐目录结构/logs/dt2024-12-01/access-00001.log /logs/dt2024-12-01/access-00002.log /logs/dt2024-12-02/access-00001.log如果将来要做小时级分析就加一层hour。目录层级不要超过三级否则路径解析本身就变成负担。用 Hive 建表时分区字段直接对应目录名比如CREATE EXTERNAL TABLE access_log ( ip STRING, time_local STRING, request STRING, status INT ) PARTITIONED BY (dt STRING) ROW FORMAT DELIMITED FIELDS TERMINATED BY LOCATION /logs;注意这里的分隔符只能处理简单空格分割的字段。如果字段内部有空格比如$request里的 URL需要用 CSV 或者正则解析或者干脆在采集时把日志转成 TSV 格式再入库。3. 写第一个 MapReduce 日志分析程序从 WordCount 到 PV 统计3.1 Maven 工程结构一个能提交到集群的 jar很多人卡在写不出能跑的 jar问题多半出在依赖范围和打包插件。MapReduce 程序的依赖只需要hadoop-client而且 scope 必须是provided这样打包时不会把 Hadoop 自带类混进 jar。用 Maven 的maven-shade-plugin做可执行包主类指定好。dependencies dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version3.3.4/version scopeprovided/scope /dependency /dependencies build plugins plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-shade-plugin/artifactId version3.4.1/version executions execution phasepackage/phase goalsgoalshade/goal/goals /execution /executions /plugin /plugins /build这里有两个容易翻车的地方。第一如果依赖 scope 没写provided打出来的 jar 可能有几十 MB提交到 YARN 时还会出现类冲突。第二shade 插件如果没有配置Main-Class就得在运行时手动指定类名体验很差。加上 manifest 配置才是完整的。3.2 Mapper解析一行日志输出日期与 1PV 统计的核心就是写一个类似 WordCount 的 MapReduce。Mapper 的输入是 HDFS 上的一行日志我们要从$time_local里拿到日期再输出(date, 1)。Nginx 默认的日期格式是18/Dec/2024:14:23:45 0800我们需要用 SimpleDateFormat 解析注意时区问题。public class AccessLogMapper extends MapperLongWritable, Text, Text, LongWritable { private static final LongWritable ONE new LongWritable(1); private Text outKey new Text(); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line value.toString(); if (line null || line.length() 0) return; // 按空格切分取 [1] 是 remote_addr取 [3] 是时间戳部分 // 实际日志格式: $remote_addr - - [$time_local] ... String[] parts line.split( ); // 防数组越界日志行过短直接丢弃 if (parts.length 4) return; String timeStr parts[3].replace([, ); // 日期解析简单做截取前 11 个字符如 18/Dec/2024 String date parseDate(timeStr); if (date ! null) { outKey.set(date); context.write(outKey, ONE); } } private String parseDate(String timeStr) { try { java.text.SimpleDateFormat sdf new java.text.SimpleDateFormat(dd/MMM/yyyy, java.util.Locale.US); java.util.Date d sdf.parse(timeStr); java.text.SimpleDateFormat out new java.text.SimpleDateFormat(yyyy-MM-dd); return out.format(d); } catch (Exception e) { return null; } } }这段代码有几个设计上的取舍。第一parseDate里用了Locale.US因为 Nginx 的月份缩写是英文如果 JVM 默认中文 localeDec可能解析失败。第二timeStr.replace([, )是因为日志里时间戳带中括号切分后$time_local字段实际是[18/Dec/2024:14:23:45需要剥掉左括号。第三parts.length 4的过滤能在源头丢弃脏数据避免后续序列化问题。3.3 Reducer 与 Combiner本地汇总降低网络压力Reducer 端做累加没有技术含量但设置 Combiner 很关键。Combiner 是在每个 Map 节点本地先做一轮 reduce能显著减少 shuffle 阶段的数据传输量。如果 reduce 函数满足交换律和结合律就可以直接用同一个类做 Combiner。public class AccessLogReducer extends ReducerText, LongWritable, Text, LongWritable { Override protected void reduce(Text key, IterableLongWritable values, Context context) throws IOException, InterruptedException { long sum 0; for (LongWritable val : values) { sum val.get(); } context.write(key, new LongWritable(sum)); } }Combiner 的用法是在 Job 里设置job.setMapperClass(AccessLogMapper.class); job.setCombinerClass(AccessLogReducer.class); job.setReducerClass(AccessLogReducer.class);有一点要提醒Combiner 并不保证被调用而且如果 reducer 逻辑不是简单的聚合就不能硬套。比如后面要算 UV 去重在 Mapper 端做顺序文件合并就必须小心。对 PV 这种纯计数场景Combiner 能带来 30% 以上的性能提升。3.4 提交到 YARN 运行命令与参数集群跑起来之前先确认 HDFS 上有输入数据。没有的话用这条命令把日志放进去hadoop fs -mkdir -p /input/logs hadoop fs -put /opt/data/access.log /input/logs/然后打包并提交mvn clean package -DskipTests hadoop jar target/log-analysis-1.0.jar com.example.PVJob /input/logs /output/pv运行期间一定要看 YARN 的日志yarn application -list yarn logs -applicationId application_xxx这里有一个非常常见的坑运行时指定/output/pv如果目录已存在Hadoop 会直接抛FileAlreadyExistsException。所以每次重跑前都要手动删掉输出目录或者在代码里用FileSystem.delete()先清理。我自己的习惯是在 Job 里写一段自动清理逻辑但要注意如果路径写错了可能把不该删的数据删掉。4. 更实用的分析维度跳出 PV看 UV、独立 IP 与 TOP 页面4.1 用 MapReduce 自带 Counter 统计独立 IPPV 是简单的加法但运营通常更关心独立访客 UV。精确 UV 需要按 IP 用户标识去重MapReduce 里最朴素的做法是在 Mapper 输出 IP 作为 keyReducer 里只输出一次。但这种方法如果想直接拿到 UV 总量还得再跑一轮计数。更取巧的办法是使用 Hadoop 的 Counter。Counter 是全局的每个 Mapper 都可以自增。我们可以在 Mapper 里用一个 HashSet 保存当前 map 分片内见过的 IP最后在 cleanup 阶段把 set size 加到自定义 Counter 上。由于每个分片只处理一部分数据不同分片之间会有重复 IP所以这个值高估了 UV。要精确还是得用 Distinct 的 Reduce。下面给一个做精确 UV 的代码public class UVMapper extends MapperLongWritable, Text, Text, NullWritable { private Text outKey new Text(); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String[] parts value.toString().split( ); if (parts.length 1) return; // 仅输出 IPNullWritable 表示不关心 value // 做法以 IP 为 key在 reducer 端去重 outKey.set(parts[0]); context.write(outKey, NullWritable.get()); } }Reducer 直接原样写出即可。这里的性能瓶颈在于 IP 重复度高的场景shuffle 会传很多NullWritable数据量还是很大。更好的做法是在 Mapper 内先做 local dedup用 TreeSet 或 HashSet 缓冲再定时 flush。这个优化在日志量大时效果很明显能减少 90% 的无效 shuffle。4.2 二次排序按访问量取 TOP 页面统计每个 URL 的访问量是常见需求。用 MR 默认排序只能按 key 排序如果我们想按照访问量倒序取 TOP N就需要二次排序。二次排序的核心思路是构造一个组合 key(url, count)在分区和分组时只按 url在排序时按 count 降序。这样每个 url 的 reducer 收到的第一条记录就是 count 最大的一条。组合 key 要继承WritableComparablepublic class UrlCountKey implements WritableComparableUrlCountKey { private String url; private long count; Override public int compareTo(UrlCountKey o) { int cmp this.url.compareTo(o.url); if (cmp ! 0) return cmp; // count 大的排在前面用负号翻转 return Long.compare(o.count, this.count); } }Job 里还要自定义Partitioner和GroupingComparator。Partitioner 只按 url 分区保证同一个 url 进同一个 reducerGroupingComparator 保证一个 url 的所有组合 key 被视为一组。这套东西写起来啰嗦但确实是面试和实际开发常考的硬功夫。4.3 更省事的选择能用 Hive 或 Spark 就不要硬写 MR写完上面的代码你可能已经觉得 MapReduce 很繁琐。实际上如果只是算 PV、UV、TOP NHive 的 SQL 几分钟就能写完执行引擎走 Tez 或 Spark 比原生 MR 快得多。很多 zip 里的程序只是为了满足课程设计要求才用 Java 写 MR。用 Hive 算 PV 的 SQL 堪称模板SELECT dt, COUNT(*) AS pv FROM access_log GROUP BY dt;算独立 IP 就是SELECT dt, COUNT(DISTINCT ip) AS uv FROM access_log GROUP BY dt;取 TOP 10 页面SELECT url, COUNT(*) AS cnt FROM access_log WHERE dt 2024-12-01 GROUP BY url ORDER BY cnt DESC LIMIT 10;Hive 的优势是写起来快但缺点也很明显如果在伪分布式环境下跑启动 HiveServer2 和 Tez 会话消耗的资源可能比作业本身还大。生产环境建议直接上 Spark SQL流式处理用 Flink。MapReduce 的价值在于理解分布式计算的底层模型一旦你搞清楚 Mapper、Reducer、Shuffle 和 Sort之后学 Spark 或 Flink 都是降维打击。5. Hadoop 日志分析避坑指南5 个常见的翻车现场5.1 时间戳时区与解析失败日志错位一天现象统计出来的日期比实际晚 8 小时或者某些小时的数据为 0。原因Nginx 日志记录的是本地时间并且$time_local里自带时区偏移例如[18/Dec/2024:14:23:45 0800]。你的SimpleDateFormat如果不解析时区会按服务器默认时区处理导致边界数据分到错误的日期。另外如果日志是 GMT 时间Hadoop 集群的user.timezone设置也会影响最终输出。解决写日志到 HDFS 之前统一在采集层把时间转成 UTC 或北京时间并在 Hive 表里把分区字段明确为dt STRING不用 TIMESTAMP 类型避免隐式转换。MapReduce 解析日期时直接用SimpleDateFormat(dd/MMM/yyyy:HH:mm:ss Z, Locale.US)把完整的偏移量也吃掉。5.2 小文件过多NameNode 内存被打爆现象HDFS 上每天有几百个几 KB 的日志文件集群跑着跑着 NameNode 开始频繁 GC最后 Active/Standby 切换。原因Flume 默认每个文件小于 1024 字节就滚动或者你用了tail -F手动 put导致每个日志块都变成独立小文件。每个文件在 NameNode 上对应一份元数据约 150 字节文件数量一多内存就扛不住。解决调整 Flume 的参数让文件至少 128 MB 或每小时滚动一次并把hdfs.minBlockReplicas保持在默认值。对已经存在的小文件可以用hadoop archive -archiveName logs.har -p /logs /archive做 HAR 归档或者用 Spark 批量做一次小文件合并。我自己常用的方案是每天凌晨跑一个定时任务把前一天的小文件用hadoop fs -getmerge合并成几个大文件再回传。5.3 数据倾斜某一个 IP 访问量远超其他 IP卡住整个 Reducer现象同一个 MR 任务其他 Reducer 几十秒跑完但有一个 Reducer 跑了几个小时还没结束。原因日志分析里高热度 IP 或热门 URL 是天然的倾斜源。默认 HashPartitioner 按 key 哈希取模热门 key 全部进同一个分区导致长尾。解决对 key 加盐。比如以 URL 为 key 做统计时可以把 URL 后面拼一个随机数让数据分散到多个 Reducer最后再做一次全局聚合。加盐的具体做法是public class SkewMapper extends MapperLongWritable, Text, Text, LongWritable { Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String url extractUrl(value.toString()); // 加盐随机数范围取决于 reducer 数量一般取 10~100 String saltedKey url _ ThreadLocalRandom.current().nextInt(50); context.write(new Text(saltedKey), ONE); } }注意加盐后统计结果需要再跑一次去掉盐的作业或者直接使用 Hive 的GROUP BY配合skewed table处理。还有一个更轻量的办法在 Mapper 侧设置job.setCombinerClass让热门 key 先在本地聚合减少倾斜总量。5.4 YARN 内存参数与容器溢出现象作业刚提交就报Container killed on request. Exit code is 143或者GC overhead limit exceeded。原因伪分布式模式下默认的yarn.nodemanager.resource.memory-mb只有 8 GB而 MapReduce 作业的mapreduce.map.memory.mb默认 1 GB如果日志解析逻辑复杂或者有内存泄漏容器会被 NodeManager 强制杀掉。解决给 YARN 和 MapReduce 设置合理的内存关系。比如单机 16 GB 内存可以这样配yarn.nodemanager.resource.memory-mb 12288 yarn.scheduler.maximum-allocation-mb 12288 mapreduce.map.memory.mb 2048 mapreduce.reduce.memory.mb 4096 mapreduce.map.java.opts -Xmx1638m mapreduce.reduce.java.opts -Xmx3276mjava.opts一般设为 memory.mb 的 80% 左右因为 JVM 除了堆还要留一些 Metaspace 和线程栈。如果你的程序用到了很大堆外内存比如解析超长日志还要相应调低比例。5.5 二次排序时比较器与分区器不一致导致数据错乱现象Reducer 里拿到的同一个 key 的数据被分成了多组reduce方法被调用多次最后统计结果偏大。原因二次排序要求Partitioner、SortComparator、GroupingComparator三个类对 key 的比较逻辑保持一致。你只重写了组合 key 的compareTo但 Job 默认的HashPartitioner使用整个组合 key 的 hashCode 去分区导致相同 url 的 key 分到不同 Reducer而默认GroupingComparator又按完整 key 比较同一 url 的不同 count 不能组成一组。解决显式指定分区器和分组比较器。分区器只取组合 key 中的 url 字段做 hashCode分组比较器只用 url 字段做相等判断。这三者的代码量确实不少建议使用 Hive 的ROW_NUMBER()窗口函数代替手写二次排序性能差距对于大多数日志分析场景可以忽略。6. 结果验证与调参技巧如何确认你的分析没算错6.1 用已知数据集做回归验证每次写完新的分析任务我都会准备一个只有几条日志的小测试文件手工算出期望结果再跑程序。比如只有三条日志两个来自同一个 IP另一个来自不同 IP那 PV 应该是 3UV 应该是 2。流程是echo 1.1.1.1 - - [18/Dec/2024:10:00:01 0800] GET /a HTTP/1.1 200 123 - curl/8.0 sample.log echo 1.1.1.1 - - [18/Dec/2024:10:00:02 0800] GET /b HTTP/1.1 200 456 - curl/8.0 sample.log echo 2.2.2.2 - - [18/Dec/2024:10:05:00 0800] GET /c HTTP/1.1 200 789 - curl/8.0 sample.log hadoop fs -mkdir -p /test/in hadoop fs -put sample.log /test/in/ hadoop jar log-analysis-1.0.jar com.example.PVJob /test/in /test/out hadoop fs -cat /test/out/part-r-00000输出为2024-12-01 3则正确。这一步能拦住大部分解析和逻辑错误。不要直接拿线上几 GB 日志跑否则一个小 bug 会让你在 YARN 日志里翻半天。6.2 调优并行度Map 和 Reduce 的数量不是越大越好Map 数量的默认逻辑是每个 HDFS block 一个 Map所以文件的大小和块大小决定了 Map 数。如果你觉得 Map 数太少可以把输入文件多分几个或者使用FileInputFormat.setMaxInputSplitSize控制分片大小。但 Map 数太多也会导致启动开销占比高。Reduce 数量相对简单我一般设置成集群可用内存的 1/2 左右或者直接根据业务指标个数来定。比如统计 PV 和 UV 可以共用一个 Job输出两个结果Reduce 设 2 就行。设置过大的 Reduce 数会让每个 Reduce 只处理很少的数据反而增加 shuffle 开销。还有一个常被忽略的参数是mapreduce.task.io.sort.mb它控制 Map 端排序缓冲区大小。日志解析场景下如果每条记录较短可以把这个值调大到 256 MB减少 spill 次数。6.3 从日志分析这个方向往后走把批处理升级为实时链路如果只是课程设计跑通 MR 已经合格。但如果你在生产环境做网站日志分析我更推荐用一个能随时重算的清洗层。把原始日志清洗成 Parquet 格式的标准表之后所有 PV、UV、留存、转化指标都基于这一层构建。这样面向变化的指标需求时你只需要写新的 SQL不用再重跑一堆 MR。这也是我从手动写 MR 的教训里学到的经验——每次需求变化都改 Java 代码、重新打包、重新跑作业非常耗时。另外一个值得养成的习惯是给所有统计结果写一个完整性校验脚本。比如对比 HDFS 输入日志的总行数和统计出来的 PV 总和偏差超过 0.1% 就报警。因为不管你用什么计算引擎脏数据、重复数据、解析失败都会静默吞掉只有对账才能发现。我现在的常规操作是每个月底跑一次全量重算确认和每日增量结果一致防止某个补数任务产生的脏数据长期潜伏。希望这篇笔记能帮你在 Hadoop 日志分析上少走一段弯路。本文还有配套的精品资源点击获取