简介:面向大数据离线分析学习者的完整项目案例文档,以某技术学习论坛的访问日志为数据源,系统讲解网站日志分析从数据采集到指标输出的完整流程。文档先介绍项目来源与数据情况,说明历史数据追加写入、自某日起每日生成数据文件的变化,让读者理解真实日志的两种组织形式;随后围绕浏览量、注册用户数、独立IP数、跳出率、板块热度排行榜五项关键指标,逐一梳理定义和计算公式,并给出上传日志至HDFS、用MapReduce清洗、用Hive统计、用Sqoop导入MySQL、用HBase存明细的整体开发链路,同时讨论了不同数据规模下日志上传方案的选择,以及统计结果表与明细日志表的结构设计。资源为1个docx文档,压缩包大小约1013KB,文字精炼、结构清晰,便于按章节查阅。目前已有1401人学习下载,适合正在学习或实践大数据离线分析、希望从综合案例中掌握日志处理全流程的读者。
1. 网站日志分析到底在解决什么问题:三个角色盯同一份 access.log
做过网站运维的人都有过这种体验:线上出问题、流量莫名波动、运营追问转化率,最后大家不约而同打开服务器上的 access.log,却发现几 GB 的文本根本没法直接读。这个大数据综合案例做的是把 nginx 这类 Web 服务器每天产生的访问日志,从零搭建一条完整链路:采集、落盘、离线统计、实时计算、可视化。日志里的 PV、UV、独立访客、热门 URL、在线人数和异常流量,最终变成一张张能直接给业务看的报表。
文章面向两类读者:一类是刚学完 Hadoop、Hive、Spark 基础,想找一个能串起全流程的落地项目;另一类是负责网站运维或数据分析,想把手里的原始日志真正用起来的从业者。我不会只讲概念,每一段都会落到配置、参数和踩坑经验上,你照着能复现,跑通了能改出自己需要的那套指标。
2. 从 nginx 日志到 HDFS:Flume 采集链路的配置与参数
2.1 采集选型:Flume 相比自写脚本,赢在断点续传与批处理
做网站日志分析,第一个要解决的问题不是「怎么分析」,而是「怎么把日志稳定地搬进大数据平台」。很多人第一反应是写个 crontab 脚本,定时把日志 cp 到 Hadoop 的 HDFS 上,简单直接。但实际跑两周就会撞上三个问题:一是日志文件正在被 nginx 写,cp 出来的文件可能缺尾部数据;二是凌晨重启服务器或脚本异常中断后,不知道哪一批文件已经传过,重传一次指标就重复一遍;三是每天几十个节点、每小时上百个小文件,靠脚本一个个 put 效率太低。
我一般会选 Flume 而不是自写脚本,核心原因有三条。第一,Flume 的 taildir source 支持断点续传,会把读取位置记录在一个 positionFile 里,进程重启后从上次的位置继续读,不需要自己维护「哪个文件传过」的状态。第二,Flume 的 HDFS sink 自带批量写入、按大小或时间滚动文件、压缩格式选择,这些是脚本里要写很久才能稳定的功能。第三,Flume 本身就是 Java 进程,异常退出后重启即可,不需要额外处理进程守护逻辑。
如果只是单机日志量每天不超过几百 MB,自写脚本也能扛住;但一旦涉及多台 Web 服务器、每天几 GB 甚至几十 GB 的日志,Flume 的可靠性和吞吐优势就很明显了。选型时还有一个常见选择:用 Kafka 直接做采集端。我的建议是,日志量大且下游有实时计算需求时,用 Flume 先落 HDFS、同步推一份到 Kafka;日志量小、只做离线分析,Flume 直接落 HDFS 就够,少一条链路少一组故障点。这个案例先解决离线,所以重点讲 Flume 落 HDFS 的完整配置。
2.2 一份能直接落地的 Flume 配置:taildir + HDFS Sink
下面这份配置是我多次实际使用后保留的最小可运行版本,采集目录、HDFS 路径、滚动参数都直接标好,替换成你自己的路径就能跑:
# flume-nginx.conf agent.sources = r1 agent.channels = c1 agent.sinks = k1 # 1. taildir source: 断点续传, 读取 nginx 日志 agent.sources.r1.type = TAILDIR agent.sources.r1.filegroups = g1 agent.sources.r1.filegroups.g1 = /var/log/nginx/access.*.log agent.sources.r1.positionFile = /data/flume/position/taildir-position.json agent.sources.r1.batchSize = 1000 agent.sources.r1.backoffSleepIncrement = 1000 agent.sources.r1.maxBackoffSleep = 5000 # 2. channel: 用 memory channel, 速度优先 agent.channels.c1.type = memory agent.channels.c1.capacity = 10000 agent.channels.c1.transactionCapacity = 1500 # 3. HDFS sink: 按小时滚动, 输出为 DataStream agent.sinks.k1.type = hdfs agent.sinks.k1.hdfs.path = /data/nginx/logs/%Y%m%d/%H agent.sinks.k1.hdfs.fileType = DataStream agent.sinks.k1.hdfs.writeFormat = Text agent.sinks.k1.hdfs.batchSize = 1000 agent.sinks.k1.hdfs.rollSize = 134217728 agent.sinks.k1.hdfs.rollCount = 0 agent.sinks.k1.hdfs.rollInterval = 3600 agent.sinks.k1.hdfs.idleTimeout = 60 agent.sinks.k1.hdfs.filePrefix = access agent.sinks.k1.hdfs.useLocalTimeStamp = true agent.sources.r1.channels = c1 agent.sinks.k1.channel = c1这份配置里最值得注意的有三个地方。第一,taildir 的positionFile要放在持久化磁盘上,不要放临时目录,否则重启丢位置,就会重新读一遍近期日志,产生重复数据。第二,hdfs.fileType = DataStream表示按普通文本方式写文件,不是 SequenceFile。很多新手抄网上的配置抄成了 SequenceFile,后续 Hive 建表就要多处理一层,没必要。第三,useLocalTimeStamp = true表示用 Flume 所在服务器的时间生成 HDFS 路径,而不是读日志内容里的时间。这样路径里的小时分区基本准确,但要注意服务器时区和日志时区一致,否则分区偏移 8 小时的问题会在后面出现。
滚动参数rollSize = 134217728表示文件达到 128 MB 就滚动,rollInterval = 3600表示最多一小时也滚动一次,idleTimeout = 60表示文件空闲 60 秒关闭。这三个参数直接决定 HDFS 上的小文件数量,建议同时保留大小和时间两个维度,避免某些小时流量低,文件永远滚不动。
2.3 采集后的第一道工序:分区规范与对账校验
日志进了 HDFS,不代表就可以放心跑 SQL。我习惯在采集当天就把分区规范定死,并在第二天做一个简单的对账。分区路径采用/data/nginx/logs/日期/小时这种粒度,对应配置里的%Y%m%d/%H。日期分区的好处是,Hive 查某一天的数据只需要读对应目录,查询速度差一个量级。
对账方法很朴素:每天凌晨对比「nginx 自己记录的 access 日志条数」和「HDFS 上对应分区文件的条数」。nginx 的 access.log 每行一条记录,直接数行数即可。我常用这样的命令检查:
# 统计原始日志行数, 作为当日基准 wc -l /var/log/nginx/access.2024-06-01.log # 统计 HDFS 上对应分区的文件行数 hdfs dfs -cat /data/nginx/logs/20240601/*/* | wc -l基准行数和 HDFS 行数相差超过 1%,就要查是采集丢失还是重复。常见原因集中在三处:taildir 的 positionFile 权限不对导致反复重读;HDFS 滚动瞬间文件还没 close 就被读走;日志轮转时 nginx 的符号链接指向了新文件而旧文件还有残余行。这些排查点后面避坑章节还会展开,这里先记住一个原则:分区要当天校验,别等一周后跑报表才发现数据不对,到那时候原始日志可能已经轮转过,想补救都难。
3. Hive 离线统计:PV、UV、热门 URL 的可复用 SQL
3.1 建表不能拍脑袋:分区、文件格式与日志切割方式
日志落到 HDFS 后,下一步是在 Hive 里建表。建表这一步看着简单,其实有两个前置决策很关键:用什么文件格式、日志行怎么切分。
文件格式我建议直接用 ORC,原因很务实:ORC 带列式存储和压缩,同样一份日志,ORC 占的空间大约是纯文本的四分之一到三分之一,查询扫描数据量也能明显下降。但注意,Flume 落地时写的是 DataStream 文本格式,所以需要先建一个指向 HDFS 文本目录的临时表,再从临时表 INSERT 到 ORC 的结果表,不能直接在文本文件上建 ORC 表。
日志切分方式取决于你的 nginx 日志格式。假设日志以空格分隔,或自定义成|分隔,我一般用默认空格切分,避免用正则解析这种性能陷阱。下面是一份可用的建表 SQL:
-- 原始文本表: 指向 Flume 写入的 HDFS 目录 CREATE EXTERNAL TABLE ods_nginx_access_text ( ip STRING, time_local STRING, request STRING, status STRING, body_bytes_sent STRING, http_referer STRING, http_user_agent STRING ) PARTITIONED BY (dt STRING, hr STRING) ROW FORMAT DELIMITED FIELDS TERMINATED BY ' ' STORED AS TEXTFILE LOCATION '/data/nginx/logs'; -- ORC 结果表: 列式存储 + 压缩, 供后续查询 CREATE TABLE dwd_nginx_access ( ip STRING, time_local STRING, request STRING, status STRING, body_bytes_sent STRING, http_referer STRING, http_user_agent STRING ) PARTITIONED BY (dt STRING, hr STRING) STORED AS ORC; -- 每日分区数据从文本表转换到 ORC 表 INSERT OVERWRITE TABLE dwd_nginx_access PARTITION (dt='2024-06-01', hr='00') SELECT ip, time_local, request, status, body_bytes_sent, http_referer, http_user_agent FROM ods_nginx_access_text WHERE dt='2024-06-01' AND hr='00';这段 SQL 里有三个细节新手容易踩。第一,FIELDS TERMINATED BY ' '只处理空格分隔的情况,如果日志里有 URL 参数带空格,会导致字段错位,这时建议在采集端用统一分隔符重写,而不是在这里改。第二,ODS 表是外部表,删除表不会删 HDFS 文件,适合指向采集链路写出的原始目录,避免误删原始数据。第三,INSERT OVERWRITE按分区写入,重复跑不会产生重复分区,这也是一种天然的幂等手段。
3.2 三个必写的统计指标 SQL:PV、UV、Top URL
建好表之后,最核心的离线统计就是 PV、UV 和热门 URL。PV 是最简单的,count 一下分区内所有行就行;UV 要复杂一些,得先明确「一个用户」怎么定义。业界常见两种口径:按用户 Cookie ID 去重,或者按 IP 去重。Cookie 口径更接近真实访客数,但需要前端埋点把 Cookie 写进日志,很多日志格式里没有;IP 口径容易把同一个 NAT 出口的多个人算成一个用户。这个案例里我们以日志中的 ip 字段为例做 UV,并把口径写清楚,避免报表对不上。三个指标的 SQL 如下:
-- PV: 统计一天内所有访问次数 SELECT COUNT(*) AS pv FROM dwd_nginx_access WHERE dt = '2024-06-01'; -- UV: 按 IP 去重, 只保留首次出现 SELECT COUNT(*) AS uv FROM ( SELECT ip FROM dwd_nginx_access WHERE dt = '2024-06-01' GROUP BY ip ) t; -- 热门 URL Top10: 按请求路径聚合 SELECT request, COUNT(*) AS visit_cnt FROM dwd_nginx_access WHERE dt = '2024-06-01' AND request LIKE 'GET /%' GROUP BY request ORDER BY visit_cnt DESC LIMIT 10;这里特别说明一下 UV 的写法。很多人第一反应是COUNT(DISTINCT ip),这个写法在小数据量时没问题,但数据量一大,Hive 里COUNT(DISTINCT)会把所有不同 ip 拉到一个 reducer 上计算,内存容易爆,速度也慢。我一般先GROUP BY ip再用COUNT(*),让聚合分散到多个 reducer,速度提升非常明显。热门 URL 那一段,LIKE 'GET /%'是为了过滤掉非请求行,有些日志里混合了错误记录或健康检查,带这个过滤条件能减少脏数据。
3.3 数据量上来之后:Spark SQL 跑同套逻辑的参数调整
当每天日志量涨到几亿行,Hive 的 MapReduce 引擎就有点力不从心了。这时我会把同样的 SQL 交给 Spark SQL 跑,不改 SQL 逻辑,只调几个关键参数。Spark SQL 的优势是内存计算,中间结果不需要落盘,同样的聚合任务往往比 Hive 快几倍。下面是一份常用的提交参数:
spark-submit \ --master yarn \ --deploy-mode client \ --driver-memory 4g \ --executor-memory 8g \ --executor-cores 4 \ --num-executors 20 \ --conf spark.sql.shuffle.partitions=200 \ --conf spark.sql.adaptive.enabled=true \ --conf spark.sql.adaptive.coalescePartitions.enabled=true \ --conf spark.sql.adaptive.skewJoin.enabled=true \ --class org.apache.spark.sql.hive.thriftserver.HiveThriftServer2 \ spark-sql参数背后的考量是:--num-executors 20+--executor-cores 4表示同时用 80 个核处理,量级按你的集群规模调整;spark.sql.shuffle.partitions决定 shuffle 时数据分成多少个分区,200 是一个起步值,如果发现某个 reducer 处理时间特别长,说明分区数相对数据量偏少,可以往上调;spark.sql.adaptive.enabled开启后,Spark 会运行时自动合并小分区,尤其是skewJoin能自动识别数据倾斜的 join 键,把大 key 单独拆分处理,这对日志里某些 IP 访问量极高的情况很有用。
注意一点:Spark SQL 跑 Hive 表,需要在提交时带上 Hive 的元数据连接配置。更省事的做法是直接用spark-sql命令行,它默认读取 Hive 的 metastore,只要把 Hive 的hive-site.xml放在 Spark 的 conf 目录下即可。这个点上我翻过车:漏放配置文件,连上 Spark 后表全部看不到,一度以为集群权限问题,排查了很久才发现是 metastore 没连上。
4. 分钟级实时链路:Kafka + Flink 计算在线人数与异常流量
4.1 为什么要在离线之外再加一套实时计算
离线链路解决了「昨天」的问题,但运营经常要的是「现在」的数字。举两个真实场景:活动页面上线十分钟,运营要立刻知道当前在线人数有没有冲上来;某个接口突然被大量请求打满,运维希望第一时间看到异常流量的来源 IP。离线报表 T+1 的时效根本覆盖不了这些需求,所以需要一条实时链路。
实时链路常见做法是:Flume 在把日志写 HDFS 的同时,再通过 Kafka Channel 把数据发一份到 Kafka,Flink 从 Kafka 消费日志,按分钟窗口计算在线人数、按 IP 聚合检测异常流量,结果写入 Redis 或 Elasticsearch 供前端展示。整体链路比离线多两个组件,但逻辑并不复杂,核心是 Flink 的窗口计算。
4.2 Flink 消费 Kafka 计算窗口 UV 的最小代码骨架
下面这段代码是一个简化版的 Flink 作业骨架,功能是每分钟统计一个窗口内的唯一访客数:
import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.api.common.serialization.SimpleStringSchema; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.windowing.time.Time; import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer; import org.apache.flink.streaming.connectors.elasticsearch.util.RetryRequestConfig; import java.time.Duration; import java.util.Properties; public class MinuteUvJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(4); env.enableCheckpointing(60000); // 每分钟做一次 checkpoint Properties kafkaProps = new Properties(); kafkaProps.setProperty("bootstrap.servers", "kafka1:9092,kafka2:9092"); kafkaProps.setProperty("group.id", "nginx-log-min-uv"); FlinkKafkaConsumer<String> consumer = new FlinkKafkaConsumer<>( "nginx-access-log", // Kafka topic new SimpleStringSchema(), kafkaProps); // 从日志行中解析出 ip 和时间戳 DataStream<String> rawStream = env.addSource(consumer) .assignTimestampsAndWatermarks( WatermarkStrategy .<String>forBoundedOutOfOrderness(Duration.ofSeconds(10)) .withTimestampAssigner((event, timestamp) -> extractTimestamp(event))); rawStream .map(line -> parseIp(line)) .keyBy(ip -> ip) .window(TumblingProcessingTimeWindows.of(Time.minutes(1))) .apply(new MinuteWindowFunction()) .print(); env.execute("nginx-log-minute-uv"); } }这段代码的核心逻辑是三段:assignTimestampsAndWatermarks负责告诉 Flink 日志里的时间字段是什么,以及最多允许乱序多少秒;keyBy(ip)把同一 IP 的访问分到同一个处理槽,窗口内自动去重;window按 1 分钟翻滚窗口聚合,输出每分钟的独立访客数。注意这里用了 Processing Time 而非 Event Time,是为了代码可读性。实际生产中建议用 Event Time,配合日志里的时间字段,这样即使到达顺序乱了,统计结果也不会偏。
4.3 乱序与迟到的取舍:Watermark 和 allowedLateness 怎么配
实时计算里最常被忽略的是乱序问题。网络抖动、Flume 批量发送、Kafka 分区之间的不均衡,都会导致日志到达 Flink 的顺序和实际发生顺序不一致。我在代码里写了forBoundedOutOfOrderness(Duration.ofSeconds(10)),表示允许日志乱序最多 10 秒,超过这个范围的迟到数据会被丢弃。这个值不能拍脑袋定,要结合日志采集链路的总延迟来设:Flume 从产生到进 Kafka 一般 1-3 秒,所以 10 秒的阈值相对保守,既能覆盖大部分抖动,又不会让窗口一直等数据。
另外一个参数是allowedLateness。Flink 窗口默认在 Watermark 越过窗口末端时触发计算,但还可以再等一段时间,允许迟到的数据补算进前一个窗口。配置方式是在窗口后加.allowedLateness(Time.seconds(5)),这样窗口关闭后 5 秒内到达的数据,如果时间戳属于上一个窗口,会触发重新计算。它的代价是同一窗口可能输出多次结果,下游 Redis 更新时要按窗口 ID 做幂等覆盖,否则同一个窗口的 UV 会被累加多次。这个点在实时报表里特别容易出问题,数值看着跳来跳去,就是因为下游没有做覆盖写。
5. 日志分析避坑:数据倾斜、时区偏移、小文件、重复采集
5.1 小文件问题:NameNode 告警,查询越来越慢
现象:跑了一段时间后,NameNode 告警说文件数量超阈值,HDFS 页面卡顿,Hive 查询 1 亿行数据比刚开始还慢。
原因:Flume 滚动参数配置不当。比如rollInterval设为 60 秒,每小时每个 source 产出 60 个文件,一天下来几千个小文件。每个文件在 NameNode 里是一条元数据记录,文件数量一多,NameNode 内存吃紧,所有文件操作都变慢。
解决:把滚动策略从时间驱动改成大小驱动为主。我的建议是rollInterval提到 3600 秒,rollSize保持在 128MB,idleTimeout设为 60 秒;另外,每天凌晨加一个合并小文件的定时任务,把前一天零碎分区合并成少量大文件。合并常用 Hive 的INSERT OVERWRITE,从 ODS 文本表读取,写入 ORC 表时用DISTRIBUTE BY dt让同一分区的数据落进同一个 reducer,减少输出文件数量。
5.2 时间全部差 8 小时:日志分区错位排查
现象:某天看报表,发现昨晚 23:00 的访问量几乎为零,但今天 07:00 的数据量异常偏高,肉眼可见的时间错位。
原因:服务器日志里记录的是北京时间,集群所在区域是 UTC 时区。Flume 配置里useLocalTimeStamp=true用的是 Flume 进程所在系统的时区,如果集群是 UTC,HDFS 上的小时目录就会偏 8 小时。凌晨附近的数据都被写进了错误的日期分区。
解决:统一时间口径,最省事的是在 Flume 配置里把useLocalTimeStamp相关的时区对齐,让所有机器都用同一时区;同时 Hive 查询里不要直接拿time_local字符串和dt分区比较,而是定义from_utc_timestamp(to_utc_timestamp(time_local,'GMT+8'),'GMT+8')这样的转换规则。核心原则是:原始日志留原始时间,分区字段用规范时间,查询时再做显式转换,避免在采集端反复改时间格式。
5.3 Reduce 卡在 99%:数据倾斜的两种解法
现象:Hive 跑 UV 统计,进度条卡在 99% 一两个小时不动,日志里看到某个 Reduce 任务的输入数据量是其他任务的几十倍。
原因:典型的数据倾斜。某个热门 IP 或某个爬虫 IP 的访问量占了全天日志的 30% 以上,按 IP 分组去重时,这个 IP 对应的数据全涌入同一个 Reduce,其他 Reduce 早干完了,就它一个还在慢慢算。
解决:分两步处理。第一步在 SQL 层面做「加盐」两阶段聚合:先给 IP 加一个随机后缀拆成多份,分别聚合去重,再对结果做最终聚合。第二步在 Spark 引擎层面,开启spark.sql.adaptive.skewJoin.enabled,让引擎自动检测倾斜 key 并拆分执行。还有一条额外建议:如果日志里明确知道某几个固定 IP 是公司健康检查或爬虫,可以在 ETL 阶段提前过滤掉,既避免倾斜,也净化了指标。
5.4 UV 突然虚高:重复采集与去重方案的取舍
现象:某天 UV 突然比前一天翻了一倍,PV 正常,用户访问行为看上去没有异常。
原因:采集链路出了重复。最典型的是 Flume 重启时 positionFile 没写成功,taildir 从较早位置重新读了一遍同一批日志;或者 Flume 在写 HDFS 时文件还没完全关闭,下游读走了部分数据,文件滚动后 Flume 又写了一次。
解决:先从根上减少重复,确保 positionFile 目录可写且持久化,Flume 重启后立即看启动日志确认读到了正确位置。再从计算侧兜底,离线 UV 改用 ROW_NUMBER 按 ip + request + time_local 去重后再计数;实时 UV 在 Redis 里用 SETNX 按窗口存已见 IP,窗口内重复 IP 不计数。两条线配合才能把重复率压到可接受范围。这也是为什么我坚持采集后必须做行数对账,重复和丢失都在 1% 以外能及时发现。
6. 把指标变成看板:三种可视化落地方案与结果验证
离线结果和实时指标都有了,最后一步是把它们变成能看的看板。这里有三条落地路径,按成本从低到高排列。
第一种,查询结果导出 MySQL,接开源 BI 工具。Hive 的统计结果用INSERT OVERWRITE写到 MySQL 表,BI 工具直连 MySQL 出折线图和柱状图。优点是部署成本最低、业务同事可以直接拖拽看数,缺点是 MySQL 只适合存放聚合结果,明细数据量太大容易拖垮查询。
第二种,明细进 Elasticsearch,用 Kibana 做探索式分析。适合需要频繁下钻的场景,比如分析某一天某个 URL 的状态码分布、某个 IP 的访问序列。Elasticsearch 对日志明细的全文检索友好,但写入吞吐有限,需要控制索引分片数和副本数。
第三种,自建报表 API + 前端图表库。实时链路的结果写 Redis,API 读取 Redis 直接返回,前端用图表库渲染。适合在线人数这类需要秒级刷新的指标,灵活性最高,代价是前后端都要投入开发量。
做完全链路之后,我强烈建议做一个「指标正确性验证」,不然报表上线了也心虚。常用做法是抽样对比:手工数一个 10 分钟窗口内的原始日志 IP 数,和实时 UV 对比,误差通常在 2% 以内;再拿今天的 PV 和昨天的同期 PV 做对比,波动超过 30% 就要查是活动流量还是采集异常。我早期的血泪经验是,全链路跑通不等于结果可信,很多指标口径不一致的问题,就是在这种对比里暴露出来的。这个案例的完整落地路径就是这样:采集定分区、离线算指标、实时补时效、最后用对账和验证把数据链路的每一环都看住。希望帮到你。
本文还有配套的精品资源,点击获取