仰望星辰工作室

better-staridc-MNBT

better-staridc-MNBT/ bt_plugins/mnbt_connector/stats_collector.py 28.3 KB · 710 行 原始文件
Z zfhsh first commit 1 天前
1#!/usr/bin/env python3
2# coding: utf-8
3
4import logging
5import logging.handlers
6import os
7import re
8import sqlite3
9import time
10from datetime import datetime
11
12
13BASE_DIR = os.path.dirname(os.path.abspath(__file__))
14STATS_DB_PATH = os.path.join(BASE_DIR, "maxiaole.db")
15
16_stats_logger_inited = False
17
18
19def get_stats_logger():
20 global _stats_logger_inited
21 logger = logging.getLogger("mnbt_stats")
22 if _stats_logger_inited:
23 return logger
24 logger.setLevel(logging.DEBUG)
25 logger.propagate = True
26 root_logger = logging.getLogger()
27 if not root_logger.handlers and not logger.handlers:
28 try:
29 log_path = os.path.join(BASE_DIR, "worker.log")
30 handler = logging.handlers.RotatingFileHandler(
31 log_path,
32 maxBytes=10 * 1024 * 1024,
33 backupCount=5,
34 encoding="utf-8",
35 )
36 handler.setLevel(logging.INFO)
37 fmt = logging.Formatter(
38 "%(asctime)s [%(levelname)s] %(message)s",
39 datefmt="%Y-%m-%d %H:%M:%S",
40 )
41 handler.setFormatter(fmt)
42 root_logger.addHandler(handler)
43 root_logger.setLevel(logging.DEBUG)
44 except Exception:
45 pass
46 _stats_logger_inited = True
47 return logger
48
49
50def log_info(msg, *args):
51 get_stats_logger().info(msg, *args)
52
53
54def log_warn(msg, *args):
55 get_stats_logger().warning(msg, *args)
56
57
58def log_error(msg, *args):
59 get_stats_logger().error(msg, *args)
60LOG_DIR = os.path.join(BASE_DIR, "../wwwlogs") if os.path.isdir(os.path.join(BASE_DIR, "../wwwlogs")) else "/www/wwwlogs"
61
62SPIDER_AGENTS = {
63 "Baiduspider", "Googlebot", "bingbot", "YandexBot", "Sogou",
64 "360Spider", "Bytespider", "Amazonbot", "AhrefsBot", "SemrushBot",
65 "DotBot", "BLEXBot", "Exabot", "MJ12bot", "SeznamBot",
66}
67
68HOUR_DATE_RE = re.compile(r'(\d+)/(\w+)/(\d+):(\d+):(\d+):(\d+)')
69
70MONTH_MAP = {
71 "Jan":1,"Feb":2,"Mar":3,"Apr":4,"May":5,"Jun":6,
72 "Jul":7,"Aug":8,"Sep":9,"Oct":10,"Nov":11,"Dec":12,
73}
74
75
76def now_text():
77 return datetime.now().strftime("%Y-%m-%d %H:%M:%S")
78
79
80def init_stats_db():
81 conn = sqlite3.connect(STATS_DB_PATH)
82 conn.execute("PRAGMA journal_mode=WAL")
83 conn.execute("PRAGMA synchronous=NORMAL")
84 conn.execute("PRAGMA busy_timeout=5000")
85 cursor = conn.cursor()
86 cursor.execute("""
87 CREATE TABLE IF NOT EXISTS log_position (
88 site_name TEXT PRIMARY KEY NOT NULL,
89 inode REAL NOT NULL DEFAULT 0,
90 size INTEGER NOT NULL DEFAULT 0,
91 mtime REAL NOT NULL DEFAULT 0,
92 updated_at TEXT NOT NULL
93 )
94 """)
95 cursor.execute("""
96 CREATE TABLE IF NOT EXISTS site_hourly_stats (
97 id INTEGER PRIMARY KEY AUTOINCREMENT,
98 site_name TEXT NOT NULL,
99 hour TEXT NOT NULL,
100 pv INTEGER NOT NULL DEFAULT 0,
101 uv INTEGER NOT NULL DEFAULT 0,
102 total_bytes INTEGER NOT NULL DEFAULT 0,
103 UNIQUE(site_name, hour)
104 )
105 """)
106 cursor.execute("""
107 CREATE TABLE IF NOT EXISTS site_uri_stats (
108 id INTEGER PRIMARY KEY AUTOINCREMENT,
109 site_name TEXT NOT NULL,
110 date TEXT NOT NULL,
111 uri TEXT NOT NULL,
112 request_count INTEGER NOT NULL DEFAULT 0,
113 total_bytes INTEGER NOT NULL DEFAULT 0,
114 UNIQUE(site_name, date, uri)
115 )
116 """)
117 cursor.execute("""
118 CREATE TABLE IF NOT EXISTS site_ip_stats (
119 id INTEGER PRIMARY KEY AUTOINCREMENT,
120 site_name TEXT NOT NULL,
121 date TEXT NOT NULL,
122 ip TEXT NOT NULL,
123 request_count INTEGER NOT NULL DEFAULT 0,
124 total_bytes INTEGER NOT NULL DEFAULT 0,
125 UNIQUE(site_name, date, ip)
126 )
127 """)
128 cursor.execute("""
129 CREATE TABLE IF NOT EXISTS site_spider_stats (
130 id INTEGER PRIMARY KEY AUTOINCREMENT,
131 site_name TEXT NOT NULL,
132 date TEXT NOT NULL,
133 spider_name TEXT NOT NULL,
134 request_count INTEGER NOT NULL DEFAULT 0,
135 UNIQUE(site_name, date, spider_name)
136 )
137 """)
138 cursor.execute("""
139 CREATE TABLE IF NOT EXISTS site_client_stats (
140 id INTEGER PRIMARY KEY AUTOINCREMENT,
141 site_name TEXT NOT NULL,
142 date TEXT NOT NULL,
143 client_type TEXT NOT NULL,
144 client_name TEXT NOT NULL DEFAULT '',
145 request_count INTEGER NOT NULL DEFAULT 0,
146 UNIQUE(site_name, date, client_type, client_name)
147 )
148 """)
149 cursor.execute("""
150 CREATE TABLE IF NOT EXISTS site_status_stats (
151 id INTEGER PRIMARY KEY AUTOINCREMENT,
152 site_name TEXT NOT NULL,
153 date TEXT NOT NULL,
154 status_code INTEGER NOT NULL,
155 request_count INTEGER NOT NULL DEFAULT 0,
156 total_bytes INTEGER NOT NULL DEFAULT 0,
157 UNIQUE(site_name, date, status_code)
158 )
159 """)
160 cursor.execute("""
161 CREATE TABLE IF NOT EXISTS site_method_stats (
162 id INTEGER PRIMARY KEY AUTOINCREMENT,
163 site_name TEXT NOT NULL,
164 date TEXT NOT NULL,
165 method TEXT NOT NULL,
166 request_count INTEGER NOT NULL DEFAULT 0,
167 UNIQUE(site_name, date, method)
168 )
169 """)
170 cursor.execute("""
171 CREATE TABLE IF NOT EXISTS site_error_logs (
172 id INTEGER PRIMARY KEY AUTOINCREMENT,
173 site_name TEXT NOT NULL,
174 date TEXT NOT NULL DEFAULT '',
175 time_local TEXT NOT NULL,
176 ip TEXT NOT NULL,
177 method TEXT NOT NULL,
178 uri TEXT NOT NULL,
179 status INTEGER NOT NULL,
180 bytes INTEGER NOT NULL DEFAULT 0,
181 referer TEXT NOT NULL DEFAULT '',
182 ua TEXT NOT NULL DEFAULT ''
183 )
184 """)
185 try:
186 cursor.execute("ALTER TABLE site_error_logs ADD COLUMN date TEXT NOT NULL DEFAULT ''")
187 except sqlite3.OperationalError:
188 pass
189 cursor.execute("CREATE INDEX IF NOT EXISTS idx_hourly_site ON site_hourly_stats(site_name)")
190 cursor.execute("CREATE INDEX IF NOT EXISTS idx_uri_site ON site_uri_stats(site_name, date)")
191 cursor.execute("CREATE INDEX IF NOT EXISTS idx_ip_site ON site_ip_stats(site_name, date)")
192 cursor.execute("CREATE INDEX IF NOT EXISTS idx_spider_site ON site_spider_stats(site_name, date)")
193 cursor.execute("CREATE INDEX IF NOT EXISTS idx_client_site ON site_client_stats(site_name, date)")
194 cursor.execute("CREATE INDEX IF NOT EXISTS idx_status_site ON site_status_stats(site_name, date)")
195 cursor.execute("CREATE INDEX IF NOT EXISTS idx_method_site ON site_method_stats(site_name, date)")
196 cursor.execute("CREATE INDEX IF NOT EXISTS idx_error_site ON site_error_logs(site_name, status)")
197 cursor.execute("CREATE INDEX IF NOT EXISTS idx_error_date ON site_error_logs(site_name, date)")
198 # 清理 60 天前的统计数据
199 cutoff = time.strftime("%Y-%m-%d", time.localtime(time.time() - 60 * 86400))
200 for table in ("site_uri_stats", "site_ip_stats",
201 "site_spider_stats", "site_client_stats", "site_status_stats",
202 "site_method_stats", "site_error_logs"):
203 cursor.execute(f"DELETE FROM {table} WHERE date<?", (cutoff,))
204 # site_hourly_stats 用 hour 列(格式 YYYY-MM-DD HH)
205 hour_cutoff = cutoff + " 00"
206 cursor.execute("DELETE FROM site_hourly_stats WHERE hour<?", (hour_cutoff,))
207 conn.commit()
208 conn.close()
209
210
211def ts_to_hour_key(time_local):
212 m = HOUR_DATE_RE.match(time_local)
213 if not m:
214 return None, None
215 day, mon_str, year, hour = m.group(1), m.group(2), m.group(3), m.group(4)
216 mon = MONTH_MAP.get(mon_str)
217 if not mon:
218 return None, None
219 hour_key = f"{year}-{mon:02d}-{int(day):02d} {int(hour):02d}"
220 date_key = f"{year}-{mon:02d}-{int(day):02d}"
221 return hour_key, date_key
222
223
224def detect_spider(ua):
225 if not ua:
226 return None
227 for name in SPIDER_AGENTS:
228 if name.lower() in ua.lower():
229 return name
230 return None
231
232
233def detect_client(ua):
234 if not ua:
235 return "unknown", "unknown"
236 ua_lower = ua.lower()
237 mobile = any(k in ua_lower for k in ("mobile", "android", "iphone", "ipad", "ipod"))
238 client_type = "mobile" if mobile else "pc"
239 if "chrome" in ua_lower and "edge" not in ua_lower:
240 client_name = "Chrome"
241 elif "firefox" in ua_lower:
242 client_name = "Firefox"
243 elif "safari" in ua_lower and "chrome" not in ua_lower:
244 client_name = "Safari"
245 elif "edge" in ua_lower:
246 client_name = "Edge"
247 elif "msie" in ua_lower or "trident" in ua_lower:
248 client_name = "IE"
249 else:
250 client_name = "unknown"
251 return client_type, client_name
252
253
254# 宝塔 Nginx 默认日志格式:
255# log_format main '$remote_addr - $remote_user [$time_local] "$request" '
256# '$status $body_bytes_sent "$http_referer" '
257# '"$http_user_agent" "$http_x_forwarded_for"';
258# 示例: 192.168.1.1 - - [14/Jul/2026:10:30:00 +0800] "GET /index.php HTTP/1.1" 200 1234 "-" "Mozilla/5.0 ..." "-"
259# 兼容 $http_x_forwarded_for 可选、referer/UA 为 "-"、HTTP/2.0 等变体。
260
261_NGINX_RE = re.compile(
262 r'^(\S+)' # 1: IP
263 r'\s+\S+\s+\S+' # ident user (通常为 - -)
264 r'\s+\[([^\]]+)\]' # 2: time_local
265 r'\s+"([^"]*)"' # 3: request (method uri protocol)
266 r'\s+(\d{3})' # 4: status
267 r'\s+(\d+)' # 5: body_bytes_sent
268 r'(?:\s+"([^"]*)")?' # 6: referer (可选)
269 r'(?:\s+"([^"]*)")?' # 7: user_agent (可选)
270)
271
272def parse_nginx_line(line):
273 try:
274 m = _NGINX_RE.match(line.strip())
275 if not m:
276 return None
277 ip = m.group(1)
278 time_local = m.group(2)
279 request = m.group(3)
280 status = int(m.group(4))
281 body_bytes = int(m.group(5))
282 referer = m.group(6) or "-"
283 ua = m.group(7) or "-"
284 req_parts = request.split()
285 method = req_parts[0] if req_parts else "-"
286 uri = req_parts[1] if len(req_parts) > 1 else "-"
287 return {
288 "ip": ip,
289 "time_local": time_local,
290 "method": method,
291 "uri": uri,
292 "status": status,
293 "bytes": body_bytes,
294 "referer": referer,
295 "ua": ua,
296 }
297 except Exception:
298 return None
299
300
301def _get_stats_conn():
302 conn = sqlite3.connect(STATS_DB_PATH)
303 conn.execute("PRAGMA journal_mode=WAL")
304 conn.execute("PRAGMA synchronous=NORMAL")
305 conn.execute("PRAGMA busy_timeout=5000")
306 return conn
307
308def get_log_positions():
309 conn = _get_stats_conn()
310 cursor = conn.cursor()
311 cursor.execute("SELECT site_name, inode, size, mtime FROM log_position")
312 rows = cursor.fetchall()
313 conn.close()
314 return {r[0]: {"inode": r[1], "size": r[2], "mtime": r[3]} for r in rows}
315
316
317def update_log_position(site_name, inode_val, size_val, mtime_val):
318 conn = _get_stats_conn()
319 cursor = conn.cursor()
320 cursor.execute(
321 "INSERT OR REPLACE INTO log_position (site_name, inode, size, mtime, updated_at) VALUES (?, ?, ?, ?, ?)",
322 (site_name, inode_val, size_val, mtime_val, now_text())
323 )
324 conn.commit()
325 conn.close()
326
327
328def reset_site_stats(site_name):
329 """重置指定站点的所有统计数据和位置记录,下次运行时会重新回溯全部历史"""
330 conn = _get_stats_conn()
331 cursor = conn.cursor()
332 tables = [
333 "site_hourly_stats",
334 "site_ip_stats",
335 "site_uri_stats",
336 "site_spider_stats",
337 "site_client_stats",
338 "site_method_stats",
339 "site_error_logs",
340 ]
341 deleted = {}
342 for table in tables:
343 cursor.execute(f"DELETE FROM {table} WHERE site_name=?", (site_name,))
344 deleted[table] = cursor.rowcount
345 cursor.execute("DELETE FROM log_position WHERE site_name=?", (site_name,))
346 deleted["log_position"] = cursor.rowcount
347 conn.commit()
348 conn.close()
349 log_info("站点 [%s] 统计数据已重置,删除记录: %s", site_name, deleted)
350 return deleted
351
352
353def upsert_hourly_stats(site_name, hour_key, pv_inc, bytes_inc, unique_ips):
354 conn = _get_stats_conn()
355 cursor = conn.cursor()
356 cursor.execute(
357 "UPDATE site_hourly_stats SET pv=pv+?, uv=uv+?, total_bytes=total_bytes+? WHERE site_name=? AND hour=?",
358 (pv_inc, len(unique_ips), bytes_inc, site_name, hour_key))
359 if cursor.rowcount == 0:
360 cursor.execute(
361 "INSERT INTO site_hourly_stats (site_name, hour, pv, uv, total_bytes) VALUES (?, ?, ?, ?, ?)",
362 (site_name, hour_key, pv_inc, len(unique_ips), bytes_inc))
363 conn.commit()
364 conn.close()
365
366
367def upsert_uri_stats(site_name, date_key, uri, req_inc, bytes_inc):
368 conn = _get_stats_conn()
369 cursor = conn.cursor()
370 cursor.execute(
371 "UPDATE site_uri_stats SET request_count=request_count+?, total_bytes=total_bytes+? WHERE site_name=? AND date=? AND uri=?",
372 (req_inc, bytes_inc, site_name, date_key, uri))
373 if cursor.rowcount == 0:
374 cursor.execute(
375 "INSERT INTO site_uri_stats (site_name, date, uri, request_count, total_bytes) VALUES (?, ?, ?, ?, ?)",
376 (site_name, date_key, uri, req_inc, bytes_inc))
377 conn.commit()
378 conn.close()
379
380
381def upsert_ip_stats(site_name, date_key, ip, req_inc, bytes_inc):
382 conn = _get_stats_conn()
383 cursor = conn.cursor()
384 cursor.execute(
385 "UPDATE site_ip_stats SET request_count=request_count+?, total_bytes=total_bytes+? WHERE site_name=? AND date=? AND ip=?",
386 (req_inc, bytes_inc, site_name, date_key, ip))
387 if cursor.rowcount == 0:
388 cursor.execute(
389 "INSERT INTO site_ip_stats (site_name, date, ip, request_count, total_bytes) VALUES (?, ?, ?, ?, ?)",
390 (site_name, date_key, ip, req_inc, bytes_inc))
391 conn.commit()
392 conn.close()
393
394
395def upsert_spider_stats(site_name, date_key, spider_name):
396 conn = _get_stats_conn()
397 cursor = conn.cursor()
398 cursor.execute(
399 "UPDATE site_spider_stats SET request_count=request_count+1 WHERE site_name=? AND date=? AND spider_name=?",
400 (site_name, date_key, spider_name))
401 if cursor.rowcount == 0:
402 cursor.execute(
403 "INSERT INTO site_spider_stats (site_name, date, spider_name, request_count) VALUES (?, ?, ?, 1)",
404 (site_name, date_key, spider_name))
405 conn.commit()
406 conn.close()
407
408
409def upsert_client_stats(site_name, date_key, client_type, client_name):
410 conn = _get_stats_conn()
411 cursor = conn.cursor()
412 cursor.execute(
413 "UPDATE site_client_stats SET request_count=request_count+1 WHERE site_name=? AND date=? AND client_type=? AND client_name=?",
414 (site_name, date_key, client_type, client_name))
415 if cursor.rowcount == 0:
416 cursor.execute(
417 "INSERT INTO site_client_stats (site_name, date, client_type, client_name, request_count) VALUES (?, ?, ?, ?, 1)",
418 (site_name, date_key, client_type, client_name))
419 conn.commit()
420 conn.close()
421
422
423def upsert_status_stats(site_name, date_key, status_code, bytes_inc):
424 conn = _get_stats_conn()
425 cursor = conn.cursor()
426 cursor.execute(
427 "UPDATE site_status_stats SET request_count=request_count+1, total_bytes=total_bytes+? WHERE site_name=? AND date=? AND status_code=?",
428 (bytes_inc, site_name, date_key, status_code))
429 if cursor.rowcount == 0:
430 cursor.execute(
431 "INSERT INTO site_status_stats (site_name, date, status_code, request_count, total_bytes) VALUES (?, ?, ?, 1, ?)",
432 (site_name, date_key, status_code, bytes_inc))
433 conn.commit()
434 conn.close()
435
436
437def upsert_method_stats(site_name, date_key, method):
438 conn = _get_stats_conn()
439 cursor = conn.cursor()
440 cursor.execute(
441 "UPDATE site_method_stats SET request_count=request_count+1 WHERE site_name=? AND date=? AND method=?",
442 (site_name, date_key, method))
443 if cursor.rowcount == 0:
444 cursor.execute(
445 "INSERT INTO site_method_stats (site_name, date, method, request_count) VALUES (?, ?, ?, 1)",
446 (site_name, date_key, method))
447 conn.commit()
448 conn.close()
449
450
451def insert_error_log(site_name, date_key, time_local, ip, method, uri, status, bytes_val, referer, ua):
452 conn = _get_stats_conn()
453 cursor = conn.cursor()
454 cursor.execute(
455 "INSERT INTO site_error_logs (site_name, date, time_local, ip, method, uri, status, bytes, referer, ua) "
456 "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
457 (site_name, date_key, time_local, ip, method, uri, status, bytes_val, referer, ua)
458 )
459 conn.commit()
460 conn.close()
461
462
463def get_site_name_from_logpath(log_path):
464 basename = os.path.basename(log_path)
465 if basename.endswith(".log"):
466 return basename[:-4]
467 return basename
468
469
470def _find_site_log_files(log_dir, site_name):
471 """查找站点的所有日志文件(当前 + 轮转),返回 (path, mtime) 列表,按 mtime 倒序"""
472 prefix = f"{site_name}.log"
473 candidates = []
474 try:
475 for entry in os.listdir(log_dir):
476 if entry == prefix or (entry.startswith(prefix + ".") or entry.startswith(prefix + "-")):
477 fpath = os.path.join(log_dir, entry)
478 try:
479 candidates.append((fpath, os.path.getmtime(fpath)))
480 except OSError:
481 continue
482 except OSError:
483 pass
484 candidates.sort(key=lambda x: x[1], reverse=True)
485 return candidates
486
487
488def _process_log_lines(site_name, filename, lines):
489 total_lines = len(lines)
490 hour_batch = {}
491 ip_batch = {}
492 uri_batch = {}
493 spider_batch = {}
494 client_batch = {}
495 status_batch = {}
496 method_batch = {}
497 error_entries = []
498 for idx, line in enumerate(lines):
499 if idx % 10000 == 0 and idx > 0:
500 log_info("文件 [%s] 已处理 %d/%d 行", filename, idx, total_lines)
501 parsed = parse_nginx_line(line.strip())
502 if not parsed:
503 continue
504 hour_key, date_key = ts_to_hour_key(parsed["time_local"])
505 if not hour_key:
506 continue
507 hour_batch.setdefault((site_name, hour_key), {"pv": 0, "bytes": 0, "ips": set()})
508 hour_batch[(site_name, hour_key)]["pv"] += 1
509 hour_batch[(site_name, hour_key)]["bytes"] += parsed["bytes"]
510 hour_batch[(site_name, hour_key)]["ips"].add(parsed["ip"])
511 uri_key = (site_name, date_key, parsed["uri"])
512 uri_batch[uri_key] = uri_batch.get(uri_key, {"req": 0, "bytes": 0})
513 uri_batch[uri_key]["req"] += 1
514 uri_batch[uri_key]["bytes"] += parsed["bytes"]
515 ip_key = (site_name, date_key, parsed["ip"])
516 ip_batch[ip_key] = ip_batch.get(ip_key, {"req": 0, "bytes": 0})
517 ip_batch[ip_key]["req"] += 1
518 ip_batch[ip_key]["bytes"] += parsed["bytes"]
519 spider_name = detect_spider(parsed["ua"])
520 if spider_name:
521 sp_key = (site_name, date_key, spider_name)
522 spider_batch[sp_key] = spider_batch.get(sp_key, 0) + 1
523 client_type, client_name = detect_client(parsed["ua"])
524 cl_key = (site_name, date_key, client_type, client_name)
525 client_batch[cl_key] = client_batch.get(cl_key, 0) + 1
526 st_key = (site_name, date_key, parsed["status"])
527 status_batch[st_key] = status_batch.get(st_key, {"req": 0, "bytes": 0})
528 status_batch[st_key]["req"] += 1
529 status_batch[st_key]["bytes"] += parsed["bytes"]
530 md_key = (site_name, date_key, parsed["method"])
531 method_batch[md_key] = method_batch.get(md_key, 0) + 1
532 if parsed["status"] >= 400:
533 error_entries.append((
534 site_name, date_key, parsed["time_local"], parsed["ip"],
535 parsed["method"], parsed["uri"], parsed["status"],
536 parsed["bytes"], parsed["referer"], parsed["ua"]
537 ))
538 log_info("文件 [%s] 行处理完成,共解析 %d 行,begin write...", filename, total_lines)
539 wconn = _get_stats_conn()
540 wcur = wconn.cursor()
541 wcur.execute("BEGIN IMMEDIATE")
542 log_info("文件 [%s] 写入 hour stats (%d 条)...", filename, len(hour_batch))
543 for (sname, hk), data in hour_batch.items():
544 wcur.execute(
545 "UPDATE site_hourly_stats SET pv=pv+?, uv=uv+?, total_bytes=total_bytes+? WHERE site_name=? AND hour=?",
546 (data["pv"], len(data["ips"]), data["bytes"], sname, hk))
547 if wcur.rowcount == 0:
548 wcur.execute(
549 "INSERT INTO site_hourly_stats (site_name, hour, pv, uv, total_bytes) VALUES (?, ?, ?, ?, ?)",
550 (sname, hk, data["pv"], len(data["ips"]), data["bytes"]))
551 log_info("文件 [%s] 写入 uri stats (%d 条)...", filename, len(uri_batch))
552 for (sname, dk, uri), data in uri_batch.items():
553 wcur.execute(
554 "UPDATE site_uri_stats SET request_count=request_count+?, total_bytes=total_bytes+? WHERE site_name=? AND date=? AND uri=?",
555 (data["req"], data["bytes"], sname, dk, uri))
556 if wcur.rowcount == 0:
557 wcur.execute(
558 "INSERT INTO site_uri_stats (site_name, date, uri, request_count, total_bytes) VALUES (?, ?, ?, ?, ?)",
559 (sname, dk, uri, data["req"], data["bytes"]))
560 log_info("文件 [%s] 写入 ip stats (%d 条)...", filename, len(ip_batch))
561 for (sname, dk, ip), data in ip_batch.items():
562 wcur.execute(
563 "UPDATE site_ip_stats SET request_count=request_count+?, total_bytes=total_bytes+? WHERE site_name=? AND date=? AND ip=?",
564 (data["req"], data["bytes"], sname, dk, ip))
565 if wcur.rowcount == 0:
566 wcur.execute(
567 "INSERT INTO site_ip_stats (site_name, date, ip, request_count, total_bytes) VALUES (?, ?, ?, ?, ?)",
568 (sname, dk, ip, data["req"], data["bytes"]))
569 log_info("文件 [%s] 写入 spider stats (%d 种)...", filename, len(spider_batch))
570 for (sname, dk, sp), cnt in spider_batch.items():
571 wcur.execute(
572 "UPDATE site_spider_stats SET request_count=request_count+? WHERE site_name=? AND date=? AND spider_name=?",
573 (cnt, sname, dk, sp))
574 if wcur.rowcount == 0:
575 wcur.execute(
576 "INSERT INTO site_spider_stats (site_name, date, spider_name, request_count) VALUES (?, ?, ?, ?)",
577 (sname, dk, sp, cnt))
578 log_info("文件 [%s] 写入 client stats (%d 条)...", filename, len(client_batch))
579 for (sname, dk, ct, cn), cnt in client_batch.items():
580 wcur.execute(
581 "UPDATE site_client_stats SET request_count=request_count+? WHERE site_name=? AND date=? AND client_type=? AND client_name=?",
582 (cnt, sname, dk, ct, cn))
583 if wcur.rowcount == 0:
584 wcur.execute(
585 "INSERT INTO site_client_stats (site_name, date, client_type, client_name, request_count) VALUES (?, ?, ?, ?, ?)",
586 (sname, dk, ct, cn, cnt))
587 log_info("文件 [%s] 写入 status stats (%d 条)...", filename, len(status_batch))
588 for (sname, dk, sc), data in status_batch.items():
589 wcur.execute(
590 "UPDATE site_status_stats SET request_count=request_count+?, total_bytes=total_bytes+? WHERE site_name=? AND date=? AND status_code=?",
591 (data["req"], data["bytes"], sname, dk, sc))
592 if wcur.rowcount == 0:
593 wcur.execute(
594 "INSERT INTO site_status_stats (site_name, date, status_code, request_count, total_bytes) VALUES (?, ?, ?, ?, ?)",
595 (sname, dk, sc, data["req"], data["bytes"]))
596 log_info("文件 [%s] 写入 method stats (%d 条)...", filename, len(method_batch))
597 for (sname, dk, md), cnt in method_batch.items():
598 wcur.execute(
599 "UPDATE site_method_stats SET request_count=request_count+? WHERE site_name=? AND date=? AND method=?",
600 (cnt, sname, dk, md))
601 if wcur.rowcount == 0:
602 wcur.execute(
603 "INSERT INTO site_method_stats (site_name, date, method, request_count) VALUES (?, ?, ?, ?)",
604 (sname, dk, md, cnt))
605 log_info("文件 [%s] 写入 error logs (%d 条)...", filename, len(error_entries))
606 for err in error_entries:
607 wcur.execute(
608 "INSERT INTO site_error_logs (site_name, date, time_local, ip, method, uri, status, bytes, referer, ua) "
609 "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", err
610 )
611 wconn.commit()
612 wconn.close()
613 log_info("文件 [%s] 数据库写入完成", filename)
614
615
616def collect_site_stats(config):
617 log_dir = LOG_DIR
618 if not os.path.isdir(log_dir):
619 log_warn("日志目录 %s 不存在,跳过站点统计", log_dir)
620 return
621 positions = get_log_positions()
622 log_files = sorted(
623 [f for f in os.listdir(log_dir) if f.endswith(".log")],
624 key=lambda f: os.path.getmtime(os.path.join(log_dir, f)),
625 reverse=True
626 )
627 log_info("找到 %d 个日志文件", len(log_files))
628 processed = 0
629 for filename in log_files:
630 log_path = os.path.join(log_dir, filename)
631 site_name = get_site_name_from_logpath(filename)
632 try:
633 st = os.stat(log_path)
634 inode = st.st_ino
635 size = st.st_size
636 mtime = st.st_mtime
637 except OSError:
638 continue
639 pos = positions.get(site_name)
640 first_time = pos is None
641 log_info("处理文件 [%s] 大小=%d 字节 (first_time=%s)", filename, size, first_time)
642
643 if first_time:
644 log_info("首次发现站点 [%s],开始回溯全部历史日志", site_name)
645 all_site_files = _find_site_log_files(log_dir, site_name)
646 all_site_files.sort(key=lambda x: x[1])
647 total_all = len(all_site_files)
648 log_info("找到 %d 个历史日志文件,将按时间顺序全部处理", total_all)
649
650 for idx, (rpath, rmtime) in enumerate(all_site_files, 1):
651 rname = os.path.basename(rpath)
652 if not os.path.exists(rpath):
653 continue
654 try:
655 rst = os.stat(rpath)
656 rinode = rst.st_ino
657 rsize = rst.st_size
658 except OSError:
659 continue
660 if rsize == 0:
661 log_info("[%d/%d] 历史文件 [%s] 为空,跳过", idx, total_all, rname)
662 continue
663 log_info("[%d/%d] 处理历史文件 [%s] 大小=%.2fMB",
664 idx, total_all, rname, rsize / 1024 / 1024)
665 try:
666 with open(rpath, "r", encoding="utf-8", errors="ignore") as handle:
667 all_lines = handle.readlines()
668 except (OSError, UnicodeDecodeError) as exc:
669 log_error("读取历史文件 [%s] 失败: %s", rname, exc)
670 continue
671 if all_lines:
672 _process_log_lines(site_name, rname, all_lines)
673 log_info("[%d/%d] 历史文件 [%s] 处理完成 (%d 行)",
674 idx, total_all, rname, len(all_lines))
675 update_log_position(site_name, rinode, rsize, rmtime)
676 log_info("站点 [%s] 历史日志回溯完成", site_name)
677 processed += 1
678 continue
679
680 if pos and pos["inode"] == inode and pos["size"] == size:
681 log_info("文件 [%s] 无变化,跳过", filename)
682 continue
683
684 offset = pos["size"] if (pos and pos["inode"] == inode) else 0
685 if offset > size:
686 offset = 0
687 log_warn("文件 [%s] 大小回退(可能被截断/轮转),从头重读", filename)
688 if offset == size:
689 log_info("文件 [%s] offset=%d == size=%d,跳过", filename, offset, size)
690 continue
691 log_info("文件 [%s] 从 offset=%d 读取增量 (size=%d, 增量=%.2fKB)",
692 filename, offset, size, (size - offset) / 1024)
693 try:
694 with open(log_path, "r", encoding="utf-8", errors="ignore") as handle:
695 handle.seek(offset)
696 lines = handle.readlines()
697 except OSError as exc:
698 log_error("读取文件 [%s] 失败: %s", filename, exc)
699 continue
700 log_info("文件 [%s] 读取了 %d 行增量数据", filename, len(lines))
701 if not lines:
702 continue
703 _process_log_lines(site_name, filename, lines)
704 update_log_position(site_name, inode, size, mtime)
705 processed += 1
706 log_info("站点 [%s] 增量处理完成,新增 %d 行", site_name, len(lines))
707 if processed == 0:
708 log_info("站点统计:无新日志")
709 else:
710 log_info("站点统计完成,共处理 %d 个站点", processed)