Agent 在凌晨 3 点崩了:我是 18 小时后才知道的

这篇讲的是真实事故复盘:一个每天处理上百条内容的自主 Agent,崩了整整 18 小时,期间没有任何告警,没有任何日志告诉我它挂了。事后我花了两周重建了整套监控体系,把血泪踩坑总结在这里。

去年 10 月,我在跑一个自动化内容分析 Agent——每隔几分钟拉一批数据、调 LLM 分析、写回数据库。它跑了好几周没出过事。

某天早上打开 dashboard,发现入库数量在凌晨 3 点 17 分突然归零,到第二天下午 9 点才有人(是我)意识到出了问题。中间 18 小时,Agent 完全沉默。没有报警、没有告警、没有任何迹象。

检查日志,发现是一次数据库连接超时触发了未捕获的异常,进程 crash 了。而我的"监控"实际上只有一行:

# 当时的"监控"代码
import logging
logging.basicConfig(level=logging.INFO)
logger.info("任务完成")

进程不在了,日志自然也不会再写了。这是个非常经典的 盲区——你只能监控到活着的进程,死了的进程什么都不会说。

一、Agent 监控为什么和普通服务不一样

先说清楚这个问题。普通 Web 服务挂了很容易感知:请求超时、状态码 5xx、负载均衡健康检查失败——外部有人在"问"它,它不应答,你就知道了。

自主 Agent 不一样。它是主动型的——自己按计划做事,没有外部触发。你没有办法通过"发请求看响应"来判断它还活着,因为它根本不接请求。

所以传统的 uptime 监控对 Agent 几乎失效。你需要一套反过来的机制:让 Agent 主动证明自己还在工作

这就是心跳检测(Heartbeat)的核心思路,也是我后来花精力最多的一块。

另一个坑是状态孤岛。Agent 在内存里维护的工作进度,进程一死就全没了。如果没有持久化,重启后要么从头跑(可能重复处理)、要么完全不知道从哪里继续。

三类典型 Agent 故障

我整理了自己踩过的故障模式:

故障类型 触发原因 常见误判 真正需要的检测方式
进程 crash 未捕获异常、OOM、信号终止 以为 Agent 在"休眠" 进程存活检查 + 心跳超时
静默挂起 网络等待死锁、LLM 请求无限期 pending 日志显示"正常"但什么都没干 任务进度时间戳监控
逻辑死循环 重试逻辑 bug、状态机卡死 CPU 高,看起来"在运行" 任务完成计数器 + 速率监控
数据污染 LLM 返回不符合预期格式、写库 silent fail 表面正常,数据悄悄出错 输出数据质量校验

这四类我全踩过,前三类都发生在 Agent 层,最后一类发生在数据层。

二、心跳机制的工程实现

心跳的思路很简单:Agent 定期"签到",如果一段时间没签到,就报警。

难点在于实现细节。我走了不少弯路。

2.1 第一版:写文件时间戳(失败)

最初我让 Agent 每分钟 touch 一个文件,用文件 mtime 判断存活:

# 第一版:写文件心跳(不推荐)
import time
from pathlib import Path

HEARTBEAT_FILE = Path("/tmp/agent-heartbeat")

def heartbeat():
    HEARTBEAT_FILE.touch()

# 检查脚本(cron 每 5 分钟跑)
def check_alive():
    if not HEARTBEAT_FILE.exists():
        alert("Agent 心跳文件不存在")
        return
    
    age = time.time() - HEARTBEAT_FILE.stat().st_mtime
    if age > 300:  # 5 分钟没更新
        alert(f"Agent 心跳超时 {age:.0f}s")

看起来没问题,实际上坑很多:

  1. 文件系统不是原子的:NFS / tmpfs / 容器 overlay 各种场景下 mtime 更新不可靠
  2. 检查进程本身可能挂:cron 任务如果因为什么原因跳过,就会误报
  3. 无法区分"忙"和"挂":Agent 在处理大任务时可能几分钟不回来写文件,导致误报

2.2 第二版:HTTP 心跳端点(改进但有开销)

让 Agent 内置一个小 HTTP server,外部定期 poll:

# Agent 内嵌心跳 server
import threading
import time
from http.server import HTTPServer, BaseHTTPRequestHandler

class HeartbeatHandler(BaseHTTPRequestHandler):
    def do_GET(self):
        if self.path == "/healthz":
            status = agent.get_status()  # 拿 Agent 自己的状态
            payload = {
                "alive": True,
                "last_task_at": status["last_task_ts"],
                "tasks_completed": status["total_tasks"],
                "current_task": status["current_task_id"],
                "queue_depth": status["queue_len"],
            }
            self.send_response(200)
            self.send_header("Content-Type", "application/json")
            self.end_headers()
            self.wfile.write(json.dumps(payload).encode())
        else:
            self.send_response(404)
            self.end_headers()
    
    def log_message(self, format, *args):
        pass  # 静默,不要污染 Agent 日志

def start_heartbeat_server(port=18888):
    server = HTTPServer(("0.0.0.0", port), HeartbeatHandler)
    thread = threading.Thread(target=server.serve_forever, daemon=True)
    thread.start()
    return server

这版比文件心跳好,但有个根本问题:HTTP server 线程活着不代表主逻辑在工作。我遇到过主循环卡死(等 LLM 响应),但 heartbeat server 还在愉快地响应 200 OK 的情况。

2.3 第三版:进度时间戳心跳(目前在用)

关键洞察:要监控的不是进程存活,而是任务进度

import time
import json
import redis

class AgentHeartbeat:
    """
    基于任务进度的心跳:只有真正完成了工作才算"活着"
    """
    
    def __init__(self, agent_id: str, redis_client: redis.Redis):
        self.agent_id = agent_id
        self.redis = redis_client
        self.key = f"agent:heartbeat:{agent_id}"
        self.ttl = 300  # 5 分钟没更新 = 死亡
    
    def beat(self, task_id: str, tasks_done: int, queue_depth: int):
        """完成一个任务后调用,不是按时间调用"""
        payload = {
            "ts": time.time(),
            "task_id": task_id,
            "tasks_done": tasks_done,
            "queue_depth": queue_depth,
        }
        # 原子操作:set + expire
        self.redis.setex(self.key, self.ttl, json.dumps(payload))
    
    def is_alive(self) -> tuple[bool, dict | None]:
        """外部检查调用"""
        raw = self.redis.get(self.key)
        if raw is None:
            return False, None
        data = json.loads(raw)
        age = time.time() - data["ts"]
        return age < self.ttl, data


# Agent 主循环使用方式
heartbeat = AgentHeartbeat("content-analyzer-prod", redis_client)

for task in task_queue:
    result = process_task(task)  # 真正的工作
    save_result(result)
    
    # 只有任务完成后才打心跳——如果卡在 process_task 里不出来,心跳就超时
    heartbeat.beat(
        task_id=task.id,
        tasks_done=counter.increment(),
        queue_depth=task_queue.qsize(),
    )

这个设计有个微妙但重要的点:心跳是在任务完成之后打,而不是按固定时间间隔打。这样如果 Agent 卡在某个任务里(比如 LLM 请求超时 60 秒),心跳就会自然超时,触发告警。

不过这也带来一个新问题:任务处理时间如果本来就很长(超过 TTL),会误报。解法是把 TTL 设为"正常最大处理时间的 3 倍",根据实际任务耗时 p99 来定。

三、状态持久化:进程死了也不怕

进程 crash 后,怎么让 Agent 重启后继续工作,而不是从头开始或者乱做一通?

这是"状态持久化"要解决的问题。

3.1 最小状态机设计

先搞清楚 Agent 到底有哪些状态需要保存。我用的是一个最小状态集:

from dataclasses import dataclass, asdict
from typing import Literal
import time

@dataclass
class AgentCheckpoint:
    """Agent 运行检查点——进程重启后从这里恢复"""
    
    agent_id: str
    
    # 进度信息
    last_processed_id: str       # 上次处理到哪里了
    tasks_completed: int         # 总计完成任务数
    
    # 当前任务状态(用于幂等重试)
    current_task_id: str | None  # 正在处理的任务 ID
    current_task_started_at: float | None  # 任务开始时间
    
    # 元信息
    checkpoint_at: float = 0.0
    version: int = 1             # schema 版本,方便升级
    
    def save(self, redis_client):
        """原子写入 Redis,同时备份到文件"""
        data = asdict(self)
        data["checkpoint_at"] = time.time()
        
        key = f"agent:checkpoint:{self.agent_id}"
        redis_client.set(key, json.dumps(data))
        
        # 双写文件备份(Redis 挂了还能从文件恢复)
        backup_path = Path(f"/var/agent-state/{self.agent_id}.json")
        backup_path.parent.mkdir(parents=True, exist_ok=True)
        backup_path.write_text(json.dumps(data, indent=2))
    
    @classmethod
    def load(cls, agent_id: str, redis_client) -> "AgentCheckpoint | None":
        """启动时加载检查点"""
        key = f"agent:checkpoint:{agent_id}"
        raw = redis_client.get(key)
        
        if raw:
            return cls(**json.loads(raw))
        
        # Redis 里没有,尝试文件备份
        backup_path = Path(f"/var/agent-state/{agent_id}.json")
        if backup_path.exists():
            return cls(**json.loads(backup_path.read_text()))
        
        return None  # 全新启动

3.2 幂等任务处理

有了检查点,还需要确保任务可以安全重复执行。Crash 时可能任务已经执行到一半——重启后要能重试而不产生副作用。

class IdempotentTaskRunner:
    """
    保证每个任务 ID 最多被处理一次
    用 Redis SET NX 做分布式互斥
    """
    
    def __init__(self, redis_client: redis.Redis):
        self.redis = redis_client
    
    def run_once(self, task_id: str, fn, *args, **kwargs):
        """只运行一次——幂等保证"""
        done_key = f"task:done:{task_id}"
        lock_key = f"task:lock:{task_id}"
        
        # 如果已经完成,跳过
        if self.redis.exists(done_key):
            return {"skipped": True, "task_id": task_id}
        
        # 抢锁(60 秒超时,防止死锁)
        acquired = self.redis.set(lock_key, "1", nx=True, ex=60)
        if not acquired:
            raise RuntimeError(f"Task {task_id} 正在被其他进程处理")
        
        try:
            result = fn(*args, **kwargs)
            # 成功后标记完成,保留 24 小时记录
            self.redis.setex(done_key, 86400, json.dumps({"done_at": time.time()}))
            return result
        finally:
            self.redis.delete(lock_key)

3.3 重启恢复逻辑

把上面两块拼在一起的主循环:

def run_agent(agent_id: str):
    redis_client = redis.Redis(...)
    heartbeat = AgentHeartbeat(agent_id, redis_client)
    runner = IdempotentTaskRunner(redis_client)
    
    # 1. 加载检查点(如果有)
    checkpoint = AgentCheckpoint.load(agent_id, redis_client)
    if checkpoint:
        start_from = checkpoint.last_processed_id
        logger.info(f"从检查点恢复,上次处理到 {start_from},已完成 {checkpoint.tasks_completed} 个任务")
    else:
        start_from = None
        logger.info("全新启动")
    
    # 2. 主循环
    counter = checkpoint.tasks_completed if checkpoint else 0
    
    for task in fetch_tasks(after_id=start_from):
        # 幂等处理
        result = runner.run_once(task.id, process_task, task)
        
        if not result.get("skipped"):
            save_result(result)
            counter += 1
        
        # 更新检查点
        cp = AgentCheckpoint(
            agent_id=agent_id,
            last_processed_id=task.id,
            tasks_completed=counter,
            current_task_id=None,
            current_task_started_at=None,
        )
        cp.save(redis_client)
        
        # 打心跳
        heartbeat.beat(task.id, counter, task_queue.qsize())

四、告警:知道挂了还不够,要知道挂在哪

心跳超时只是告诉你"出问题了",但不告诉你出了什么问题。我最后搭的告警体系分三层:

第一层:存活告警(心跳超时)

  • 触发条件:Redis key 过期
  • 响应 SLA:5 分钟
  • 告警渠道:飞书机器人(高优先级)

第二层:进度告警(任务完成速率异常)

  • 触发条件:1 小时内任务完成数 < 历史均值的 20%
  • 响应 SLA:30 分钟
  • 这个能发现"活着但没干活"的情况

第三层:数据质量告警(输出异常)

  • 触发条件:入库数据字段缺失率 > 5%,或字段值分布异常
  • 响应 SLA:1 小时

实现用的是一个简单的 cron 脚本,每 2 分钟跑:

# monitor.py — 放进 cron 每 2 分钟跑
def check_agents():
    agents = get_registered_agents()  # 从配置读取期望运行的 agent 列表
    
    for agent_id in agents:
        hb = AgentHeartbeat(agent_id, redis_client)
        alive, data = hb.is_alive()
        
        if not alive:
            send_alert(
                level="critical",
                title=f"[Agent Down] {agent_id}",
                body=f"心跳超时,最后心跳:{data['ts'] if data else '无记录'}"
            )
            continue
        
        # 检查进度速率
        rate = get_completion_rate(agent_id, window_minutes=60)
        baseline = get_baseline_rate(agent_id)
        
        if rate < baseline * 0.2:
            send_alert(
                level="warning",
                title=f"[Agent Slow] {agent_id}",
                body=f"任务速率 {rate:.1f}/h,基线 {baseline:.1f}/h,仅 {rate/baseline:.0%}"
            )

五、实战效果

上面这套体系跑了大约 4 个月,收集到了一些数据(从运维日志统计):

指标 引入前(3 个月平均) 引入后(4 个月平均)
平均故障发现时间 MTTD 约 11 小时 约 6 分钟
平均恢复时间 MTTR 约 4 小时(手动重建状态) 约 8 分钟(自动重启 + 检查点恢复)
因故障导致的任务丢失率 约 12%(估算) 0.3%(可追踪)
误报率 约 4%(主要是任务本来就耗时长导致的)

说一下 4% 误报率怎么来的——有几类任务本身处理时间就可能超过 10 分钟(调 LLM 做长文分析),这些任务会触发心跳超时误报。解法是给不同任务类型配不同的 TTL:

# 不同任务类型用不同 TTL
TASK_TTLS = {
    "quick_classify": 60,      # 快速分类:1 分钟超时
    "long_analysis": 900,      # 长文分析:15 分钟超时
    "batch_embed": 300,        # 批量 embedding:5 分钟超时
}

自动重启是用 systemd 做的:

# /etc/systemd/system/content-agent.service
[Service]
Restart=on-failure
RestartSec=10
StartLimitInterval=60
StartLimitBurst=3

加上状态持久化之后,重启后 Agent 会从检查点继续——用户几乎感知不到有过 crash。

六、一个容易忽视的坑:检查点本身可能被污染

踩了这个坑才加的:如果 Agent 在处理过程中遇到格式错误的数据,可能把 current_task_id 更新了但任务实际上没有正确完成。重启后从这个检查点恢复,会跳过这个坏任务,然后一路跑下去——表面上恢复了,实际数据有洞。

解法:给检查点加版本和完整性校验:

import hashlib

@dataclass
class AgentCheckpoint:
    # ...之前的字段...
    checksum: str = ""
    
    def compute_checksum(self) -> str:
        """计算关键字段的校验和"""
        payload = f"{self.agent_id}:{self.last_processed_id}:{self.tasks_completed}"
        return hashlib.sha256(payload.encode()).hexdigest()[:16]
    
    def save(self, redis_client):
        self.checksum = self.compute_checksum()
        # ...之前的保存逻辑...
    
    @classmethod
    def load(cls, agent_id, redis_client) -> "AgentCheckpoint | None":
        # ...加载逻辑...
        cp = cls(**json.loads(raw))
        
        # 校验完整性
        expected = cp.compute_checksum()
        if cp.checksum != expected:
            logger.error(f"检查点校验失败!expected={expected}, got={cp.checksum}")
            logger.error("丢弃损坏的检查点,从头开始")
            return None
        
        return cp

常见问题

Q:Agent 的心跳和普通 Web 服务的健康检查有什么本质区别?

A:Web 服务健康检查是被动响应——你问它,它回答,超时就是挂了。Agent 心跳是主动上报——它自己定期汇报进度,超时才说明出问题了。更关键的区别在于:好的 Agent 心跳应该反映"任务进度"而不是"进程存活"——进程活着但卡在某个 LLM 调用上无法前进,对用户来说和挂了没区别。

Q:Redis 本身挂了怎么办?检查点和心跳都依赖 Redis?

A:所以要双写。检查点同时写文件系统(本文第 3.1 节代码里有);心跳可以加一个 fallback——Redis 不可达时退化为写本地文件,监控脚本同时检查两处。Redis 单点故障是个真实风险,如果 Agent 是生产级的,用 Redis Sentinel 或 Redis Cluster。

Q:这套方案适用于多 Agent 并行场景吗?

A:基本适用,但有两处要改:一是每个 Agent 实例需要唯一 ID(加 hostname 或 pod 名后缀);二是幂等锁(第 3.2 节的 run_once)可以天然支持多实例抢锁,只需要确保 task_id 唯一。真正复杂的多 Agent 协调(比如 task 分发、结果聚合)需要额外的编排层,不在本文范围内。


从第一次 18 小时无感知故障,到现在 MTTD 压到 6 分钟以内,核心不是什么高深的技术,就是把三件事做扎实:让 Agent 主动说话(心跳)、让状态能活过重启(检查点)、让告警真的告警(分层监控)

监控 Agent 这件事,工程味道比 AI 味道重得多。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐