Nginx访问日志是业务系统最基础的数据资产之一,每一次HTTP请求都会在access log中留下完整记录。当日志量达到每天几个GB甚至几十GB时,单机脚本已经无法胜任清洗和统计任务,这时候引入Spark做离线批处理就成了自然的选型。本文将围绕Nginx日志的采集、解析、分析和结果输出四个环节,完整讲解一套可落地的离线日志处理方案。

一、定制Nginx日志格式,为解析打好基础
默认的combined日志格式信息量有限,缺少上游响应时间、请求体大小等关键字段。在做离线分析之前,第一步应该调整log_format,把需要的字段显式打出来,并且用不可见于URL的字符做分隔,方便后续Spark解析时减少出错概率。
推荐使用JSON格式输出日志,虽然体积会增大约百分之二十,但解析的可靠性远高于正则匹配,尤其是URL中可能携带空格、引号等特殊字符的场景下,JSON的优势非常明显。配置示例如下:
http {
log_format json_log escape=json '{'
'"remote_addr":"$remote_addr",'
'"time_local":"$time_iso8601",'
'"request_method":"$request_method",'
'"request_uri":"$request_uri",'
'"status":"$status",'
'"body_bytes_sent":"$body_bytes_sent",'
'"http_referer":"$http_referer",'
'"http_user_agent":"$http_user_agent",'
'"upstream_response_time":"$upstream_response_time",'
'"request_time":"$request_time"'
'}';
access_log /var/log/nginx/access_json.log json_log buffer=32k flush=5s;
}
两个细节值得注意。第一,escape=json参数会对变量值中的特殊字符做转义,保证输出的JSON合法;第二,buffer=32k flush=5s开启了日志缓冲写入,减少磁盘IO压力,代价是日志最多延迟5秒落盘,对离线分析场景完全可以接受。
二、日志采集与HDFS存储规划
日志分散在每台Nginx机器本地,必须先汇聚到HDFS才能交给Spark处理。常见做法有两种:Flume和Filebeat。Flume的优势是和HDFS生态集成成熟,天然支持按时间滚动目录;Filebeat更轻量,适合机器数量多、不想额外部署Java进程的场景,通常搭配Kafka作为中转。
如果采用Flume方案,配置一个TAILDIR source监听日志文件,sink直接写HDFS,按小时切分目录:
a1.sources = r1 a1.sinks = k1 a1.channels = c1 a1.sources.r1.type = TAILDIR a1.sources.r1.filegroups = f1 a1.sources.r1.filegroups.f1 = /var/log/nginx/access_json.log a1.sources.r1.positionFile = /opt/flume/position.json a1.sinks.k1.type = hdfs a1.sinks.k1.hdfs.path = hdfs://namenode:9000/nginx/logs/dt=%Y%m%d/hour=%H a1.sinks.k1.hdfs.rollInterval = 300 a1.sinks.k1.hdfs.rollSize = 134217728 a1.sinks.k1.hdfs.fileType = DataStream a1.channels.c1.type = file a1.sources.r1.channels = c1 a1.sinks.k1.channel = c1
HDFS目录建议带上日期和小时分区,比如/nginx/logs/dt=20250101/hour=13,这样Spark做按天调度时可以直接用分区裁剪,只扫描需要的目录,避免全量读取。rollSize设置为128MB左右与HDFS块大小对齐,能避免产生大量小文件,这是很多团队踩过的坑:小文件过多会导致NameNode内存吃紧,同时Spark任务读取时产生海量Task,整体耗时成倍增加。
三、Spark清洗与解析实现
拿到JSON格式日志后,用Spark解析非常直接。Spark内置了JSON数据源,可以把每行日志读成StructType结构。清洗环节要做的事情包括:过滤掉非法行、剔除健康检查请求、把状态码和响应时间转成数值类型、对空字段填充默认值。
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.functions;
SparkSession spark = SparkSession.builder()
.appName("NginxLogClean")
.enableHiveSupport()
.getOrCreate();
Dataset<Row> logs = spark.read()
.json("hdfs://namenode:9000/nginx/logs/dt=20250101/*");
// 过滤健康检查与非法记录
Dataset<Row> cleaned = logs
.filter(functions.col("request_uri").isNotNull())
.filter(!functions.col("request_uri").startsWith("/health"))
.withColumn("status", functions.col("status").cast("int"))
.withColumn("request_time", functions.col("request_time").cast("double"))
.withColumn("hour", functions.substring(functions.col("time_local"), 11, 2));
cleaned.createOrReplaceTempView("nginx_logs");
清洗后的数据可以注册为临时视图,直接用SQL做统计分析,这对数据分析人员更加友好。常见的分析维度包括PV、UV、状态码分布、慢接口统计等:
-- 每小时PV和UV统计
SELECT hour, COUNT(*) AS pv, COUNT(DISTINCT remote_addr) AS uv
FROM nginx_logs
GROUP BY hour ORDER BY hour;
-- 慢接口Top 10,平均耗时超过1秒
SELECT split_part(request_uri, '?', 1) AS api,
COUNT(*) AS cnt,
ROUND(AVG(request_time), 3) AS avg_time,
MAX(request_time) AS max_time
FROM nginx_logs
GROUP BY split_part(request_uri, '?', 1)
HAVING AVG(request_time) > 1
ORDER BY avg_time DESC
LIMIT 10;
-- 状态码分布
SELECT status, COUNT(*) AS cnt
FROM nginx_logs
GROUP BY status ORDER BY cnt DESC;
注意split_part(request_uri, '?', 1)这一步的去参数化处理。如果不把查询串去掉,同一个接口会因为参数不同被统计成无数条记录,结果完全没有参考价值。另外COUNT(DISTINCT)在大数据量下是性能杀手,如果UV统计量级很大,可以考虑改用近似去重函数approx_count_distinct,误差率在百分之一以内,速度却能快一个数量级。
四、结果落库与调度运维
分析结果一般要写入MySQL或ClickHouse供报表系统查询。数据量不大时直接用JDBC写入MySQL即可;如果下游是实时看板,推荐写入ClickHouse,它的列式存储对聚合查询做了深度优化,亿级数据的分组统计也能秒级返回。
// 结果写入MySQL
hourStats.write()
.format("jdbc")
.option("url", "jdbc:mysql://dbhost:3306/log_report")
.option("dbtable", "nginx_hour_stats")
.option("user", "report")
.option("password", "******")
.option("batchsize", 2000)
.mode("append")
.save();
调度层面建议使用Airflow或DolphinScheduler,每天凌晨两点拉起前一天的日志处理任务,任务之间用依赖串联:采集完成、清洗完成、分析完成、落库完成。同时要考虑失败重跑机制,因为整个链路是幂等的(按分区覆盖写入),重跑不会产生重复数据。
最后提几个性能调优要点。读取小文件多时开启repartition合并分区,避免Task过多;shuffle类操作适当调大spark.sql.shuffle.partitions,默认200在大集群往往偏小;缓存复用频繁访问的中间结果集时用cache()并选择MEMORY_AND_DISK级别,防止内存不足直接报错。这套方案经过合理配置后,单日百GB级别的日志处理通常可以在十分钟内完成,完全满足离线分析的需求。