You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Python实时监控日志:统计指定时段内关键词出现次数

实现流式日志的时间窗口统计

核心思路

要完成这个需求,需要覆盖三个核心环节:

  1. 实时读取日志:通过模拟tail -f的方式持续获取日志新增内容
  2. 时间窗口维护:用队列存储最近X秒内的日志,自动剔除过期条目
  3. 定时统计:每隔X秒对窗口内的日志按关键词计数

完整代码实现

import subprocess
from datetime import datetime, timedelta
from collections import deque
import threading
import time

# 配置参数
LOG_PATH = "/path/to/your/logfile.log"
WINDOW_SECONDS = 10  # 统计最近10秒的日志
CHECK_INTERVAL = 10  # 每隔10秒统计一次
TARGET_KEYWORDS = ["Error", "Normal"]  # 需要统计的关键词

# 存储最近WINDOW_SECONDS内的日志,每个元素是(时间戳datetime对象, 日志内容)
log_queue = deque()

def parse_timestamp(line):
    """解析日志行的时间戳,返回datetime对象"""
    try:
        time_str = line.split(" : ")[0]
        # 匹配日志时间格式:2022-11-15 14:00:00,000
        return datetime.strptime(time_str, "%Y-%m-%d %H:%M:%S,%f")
    except (IndexError, ValueError):
        # 无法解析时间的日志行直接忽略
        return None

def read_logs():
    """持续读取日志并维护时间窗口队列"""
    # 启动tail -f命令实时拉取日志
    with subprocess.Popen(["tail", "-f", LOG_PATH], stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True) as proc:
        for line in proc.stdout:
            line = line.strip()
            if not line:
                continue
            timestamp = parse_timestamp(line)
            if not timestamp:
                continue
            
            # 先清理队列中过期的日志(早于当前时间-WINDOW_SECONDS的条目)
            current_time = datetime.now()
            cutoff_time = current_time - timedelta(seconds=WINDOW_SECONDS)
            while log_queue and log_queue[0][0] < cutoff_time:
                log_queue.popleft()
            
            # 将当前日志加入队列
            log_queue.append((timestamp, line))

def count_keywords():
    """统计窗口内各关键词的出现次数"""
    count = {keyword: 0 for keyword in TARGET_KEYWORDS}
    for _, line in log_queue:
        for keyword in TARGET_KEYWORDS:
            if keyword in line:
                count[keyword] += 1
    
    # 打印统计结果
    print(f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] 最近{WINDOW_SECONDS}秒统计:")
    for keyword, cnt in count.items():
        print(f"  {keyword}: {cnt}")
    
    # 重启定时器,实现周期性统计
    threading.Timer(CHECK_INTERVAL, count_keywords).start()

if __name__ == "__main__":
    # 启动日志读取线程(后台运行)
    log_thread = threading.Thread(target=read_logs, daemon=True)
    log_thread.start()
    
    # 触发第一次统计,之后自动循环
    count_keywords()
    
    # 保持主线程存活
    while True:
        time.sleep(1)

关键细节说明

  • 日志读取:用subprocess.Popen调用系统tail -f命令,无需自己实现文件变更监控,逻辑更简洁
  • 时间窗口维护:每次新增日志前,先清理队列中超出时间窗口的旧条目,确保队列只保留有效数据
  • 定时统计:通过threading.Timer实现定时触发,每次统计完成后重启定时器,保证周期性执行
  • 异常处理:对无法解析时间的日志行直接忽略,避免程序因格式异常崩溃

注意事项

  • 若在Windows环境运行,需将tail -f替换为PowerShell命令Get-Content -Wait,对应修改subprocess.Popen的参数为["powershell", "-Command", f"Get-Content -Wait {LOG_PATH}"]
  • 若日志时间为UTC时区,需在解析时添加时区参数,避免本地时间与日志时间出现偏差
  • WINDOW_SECONDS和CHECK_INTERVAL可独立设置,比如每隔5秒统计最近10秒的日志

内容的提问来源于stack exchange,提问作者user370507

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.12 16:20:28