简介:一个基于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 重启后找不到上次读取位置,它会把整个日志文件重新读一遍,导致重复数据。rollInterval=3600表示每小时滚动一次文件,不设置的话小文件会一直累积,对大集群来说很伤 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/dt=2024-12-01/access-00001.log /logs/dt=2024-12-01/access-00002.log /logs/dt=2024-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> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-client</artifactId> <version>3.3.4</version> <scope>provided</scope> </dependency> </dependencies> <build> <plugins> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <version>3.4.1</version> <executions> <execution> <phase>package</phase> <goals><goal>shade</goal></goals> </execution> </executions> </plugin> </plugins> </build>这里有两个容易翻车的地方。第一,如果依赖 scope 没写provided,打出来的 jar 可能有几十 MB,提交到 YARN 时还会出现类冲突。第二,shade 插件如果没有配置Main-Class,就得在运行时手动指定类名,体验很差。加上 manifest 配置才是完整的。
3.2 Mapper:解析一行日志输出日期与 1
PV 统计的核心就是写一个类似 WordCount 的 MapReduce。Mapper 的输入是 HDFS 上的一行日志,我们要从$time_local里拿到日期,再输出(date, 1)。Nginx 默认的日期格式是18/Dec/2024:14:23:45 +0800,我们需要用 SimpleDateFormat 解析,注意时区问题。
public class AccessLogMapper extends Mapper<LongWritable, 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 默认中文 locale,Dec可能解析失败。第二,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 Reducer<Text, LongWritable, Text, LongWritable> { @Override protected void reduce(Text key, Iterable<LongWritable> 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 统计独立 IP
PV 是简单的加法,但运营通常更关心独立访客 UV。精确 UV 需要按 IP + 用户标识去重,MapReduce 里最朴素的做法是在 Mapper 输出 IP 作为 key,Reducer 里只输出一次。但这种方法如果想直接拿到 UV 总量,还得再跑一轮计数。
更取巧的办法是使用 Hadoop 的 Counter。Counter 是全局的,每个 Mapper 都可以自增。我们可以在 Mapper 里用一个 HashSet 保存当前 map 分片内见过的 IP,最后在 cleanup 阶段把 set size 加到自定义 Counter 上。由于每个分片只处理一部分数据,不同分片之间会有重复 IP,所以这个值高估了 UV。要精确,还是得用 Distinct 的 Reduce。
下面给一个做精确 UV 的代码:
public class UVMapper extends Mapper<LongWritable, 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; // 仅输出 IP,NullWritable 表示不关心 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 要继承WritableComparable:
public class UrlCountKey implements WritableComparable<UrlCountKey> { 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 进同一个 reducer;GroupingComparator 保证一个 url 的所有组合 key 被视为一组。这套东西写起来啰嗦,但确实是面试和实际开发常考的硬功夫。
4.3 更省事的选择:能用 Hive 或 Spark 就不要硬写 MR
写完上面的代码,你可能已经觉得 MapReduce 很繁琐。实际上,如果只是算 PV、UV、TOP N,Hive 的 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 Mapper<LongWritable, 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 应该是 3,UV 应该是 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 日志分析上少走一段弯路。
本文还有配套的精品资源,点击获取