Agent 故障隔离与自愈:单个 Agent 异常不影响其他会话

一个 Agent 崩了不可怕,可怕的是一个崩了把二十个正常的也拉下水。

一、场景痛点

你们的 Agent 平台上线了 500 个并发会话,每个会话对应一个独立的 Agent 实例。一切运行良好,直到一个用户的会话触发了一个边缘 case——Agent 调用的工具返回了畸形的 JSON,导致 Agent 的推理循环陷入了死循环,CPU 吃到 100%,内存瞬间涨到 32GB。

如果没有隔离机制,这个"发疯"的 Agent 会影响同 Pod 内其他 20 个正常的 Agent——CPU 被抢占、内存被挤占、OOM Killer 随机杀掉进程。最终一个用户的错误输入,导致了 21 个用户的服务中断。

这就是故障传播(Failure Propagation)。在 Agent 系统中,故障传播比传统微服务更危险——因为每个 Agent 实例都承载着用户的实时交互会话,不可恢复的状态(对话历史、工具调用栈)一旦丢失就是一次糟糕的 UX 体验。

二、底层机制与原理剖析

2.1 Agent 故障隔离的分层架构

2.2 故障隔离的核心手段

隔离层 技术方案 解决的问题 隔离效果
进程隔离 每个 Agent 独立 fork 内存泄漏、segfault 不传播 强隔离
cgroup 限资源 cgroup v2 memory.max / cpu.max 单个 Agent 不能吃光资源 硬限制
命名空间 PID/IPC/Net namespace 网络故障不扩散 强隔离
断路器 连续失败计数 + 状态机 问题 Agent 自动熔断停止调用 软限制
会话恢复 Redis + Checkpoint + Replay 熔断后用户可无缝恢复 业务连续性
健康检查 心跳 + 超时检测 及时发现僵死 Agent 快速发现

2.3 断路器(Circuit Breaker)的状态机

          ┌──────────────────────────────────┐
          │              CLOSED               │ ← 正常状态,请求通过
          │  失败计数器: 0                    │
          │  超时/异常 → 计数器 +1            │
          └──────────┬───────────────────────┘
                     │ 失败次数 >= 阈值 (3)
                     ▼
          ┌──────────────────────────────────┐
          │              OPEN                 │ ← 熔断状态,直接拒绝
          │  启动 30s 冷却计时器              │
          │  所有请求立即返回 CircuitBreaker  │
          └──────────┬───────────────────────┘
                     │ 冷却时间到
                     ▼
          ┌──────────────────────────────────┐
          │            HALF_OPEN              │ ← 半开状态,探测恢复
          │  允许 1 个请求通过                │
          │  成功 → 回到 CLOSED               │
          │  失败 → 回到 OPEN (重置计时器)    │
          └──────────────────────────────────┘

三、生产级代码实现

3.1 进程级隔离 + cgroup 资源限制

"""
Agent 沙箱容器:为每个 Agent 会话创建独立的进程 + cgroup 资源组
实现进程级故障隔离,确保单个 Agent 异常不影响其他会话

架构:
  AgentManager (主进程)
    ├── Agent Worker 1 (forked, cgroup: /agent_pool/agent_1)
    ├── Agent Worker 2 (forked, cgroup: /agent_pool/agent_2)
    ├── ...
    └── Agent Worker N (forked, cgroup: /agent_pool/agent_N)
"""
import os
import sys
import time
import json
import signal
import multiprocessing as mp
from multiprocessing import Process, Queue, Event
from dataclasses import dataclass, field
from typing import Optional
from enum import Enum
import threading

class AgentState(Enum):
    IDLE = "idle"           # 空闲,等待分配会话
    RUNNING = "running"     # 正在处理用户请求
    BLOCKED = "blocked"     # 等待工具调用返回
    DEGRADED = "degraded"   # 性能下降(高延迟)
    CIRCUIT_OPEN = "open"   # 断路器打开(被熔断)
    DEAD = "dead"           # 进程已死亡

@dataclass
class AgentSandboxConfig:
    """Agent 沙箱资源配置"""
    agent_id: str
    session_id: str
    cpu_limit_cores: float = 1.0        # CPU 核数上限
    memory_limit_mb: int = 4096          # 内存上限 (MB)
    max_consecutive_failures: int = 3    # 连续失败阈值
    cooldown_seconds: int = 30           # 断路器冷却时间
    heartbeat_interval: float = 1.0      # 心跳间隔 (秒)
    heartbeat_timeout: float = 5.0       # 心跳超时
    max_runtime_seconds: int = 300       # 单次请求最大执行时间

class CgroupManager:
    """
    cgroup v2 资源限制管理器
    为每个 Agent 进程组创建独立的 cgroup,实现 CPU 和内存的硬限制
    
    cgroup v2 路径结构:
      /sys/fs/cgroup/agent_pool/
        └── agent_{id}/
            ├── cpu.max       # "100000 100000" (1 核 = 100%)
            ├── memory.max    # "4294967296" (4GB)
            └── cgroup.procs  # Agent 主进程 + 子进程 PID
    """
    
    CGROUP_ROOT = "/sys/fs/cgroup/agent_pool"
    
    @classmethod
    def init_root_group(cls):
        """初始化 cgroup 根目录"""
        os.makedirs(cls.CGROUP_ROOT, exist_ok=True)
        # 设置根组的默认限制(总资源池)
        with open(f"{cls.CGROUP_ROOT}/memory.max", 'w') as f:
            f.write("max")  # 根组不限制
    
    @classmethod
    def create_agent_group(cls, agent_id: str, config: AgentSandboxConfig):
        """
        为指定 Agent 创建 cgroup 资源组
        
        注意: 需要 root 权限或适当的 cgroup 委派
        生产环境通过 systemd 的 Delegate=yes 授权
        """
        cgroup_path = f"{cls.CGROUP_ROOT}/agent_{agent_id}"
        os.makedirs(cgroup_path, exist_ok=True)
        
        # CPU 限制: 格式 "MAX PERIOD" (单位: 微秒)
        cpu_quota = int(config.cpu_limit_cores * 100_000)  # 1 核 = 100000 µs
        with open(f"{cgroup_path}/cpu.max", 'w') as f:
            f.write(f"{cpu_quota} 100000")
        
        # 内存限制
        memory_bytes = config.memory_limit_mb * 1024 * 1024
        with open(f"{cgroup_path}/memory.max", 'w') as f:
            f.write(str(memory_bytes))
        
        # 启用内存 OOM 通知(超限时内核发通知而不是直接 kill)
        try:
            with open(f"{cgroup_path}/memory.oom.group", 'w') as f:
                f.write("1")
        except OSError:
            pass  # 某些内核版本可能不支持
        
        return cgroup_path
    
    @classmethod
    def add_process_to_group(cls, agent_id: str, pid: int):
        """将进程 PID 加入 Agent 的 cgroup"""
        cgroup_path = f"{cls.CGROUP_ROOT}/agent_{agent_id}"
        procs_file = f"{cgroup_path}/cgroup.procs"
        with open(procs_file, 'w') as f:
            f.write(str(pid))
    
    @classmethod
    def remove_agent_group(cls, agent_id: str):
        """清理 Agent 的 cgroup 资源组(需要先移出所有进程)"""
        cgroup_path = f"{cls.CGROUP_ROOT}/agent_{agent_id}"
        try:
            os.rmdir(cgroup_path)
        except OSError:
            pass  # 目录非空说明还有子 cgroup,保留

class CircuitBreaker:
    """
    Agent 断路器
    连续失败检测 + 自动熔断 + 半开探测恢复
    
    状态转换:
    CLOSED --(failures >= threshold)--> OPEN
    OPEN   --(cooldown elapsed)-------> HALF_OPEN
    HALF_OPEN --(probe success)-------> CLOSED
    HALF_OPEN --(probe failure)-------> OPEN
    """
    
    def __init__(self, max_failures: int = 3, cooldown_sec: int = 30):
        self.max_failures = max_failures
        self.cooldown_sec = cooldown_sec
        self.failure_count = 0
        self.last_failure_time = 0.0
        self.state = "CLOSED"
        self._lock = threading.Lock()
    
    def record_success(self):
        """记录一次成功调用"""
        with self._lock:
            self.failure_count = 0
            if self.state == "HALF_OPEN":
                self.state = "CLOSED"
    
    def record_failure(self):
        """记录一次失败调用,检查是否需要触发熔断"""
        with self._lock:
            self.failure_count += 1
            self.last_failure_time = time.time()
            
            if self.failure_count >= self.max_failures and self.state == "CLOSED":
                self.state = "OPEN"
                return True  # 触发熔断
        return False
    
    def allow_request(self) -> bool:
        """
        判断当前是否允许请求通过
        
        CLOSED → True (允许所有请求)
        OPEN   → 检查冷却时间是否到期
        HALF_OPEN → True (只允许 1 个探测请求)
        """
        with self._lock:
            if self.state == "CLOSED":
                return True
            
            if self.state == "OPEN":
                elapsed = time.time() - self.last_failure_time
                if elapsed >= self.cooldown_sec:
                    self.state = "HALF_OPEN"
                    self.failure_count = 0  # 重置计数器
                    return True
                return False
            
            if self.state == "HALF_OPEN":
                return True  # 允许探测请求
        return False

class AgentSandbox:
    """
    Agent 沙箱实例
    封装一个独立的 Agent 进程 + 资源限制 + 断路器 + 心跳
    
    生命周期:
    1. fork 子进程 → 2. 绑定 cgroup → 3. 初始化 Agent → 
    4. 心跳监控循环 → 5. 接收到 exit 信号 → 6. 清理资源
    """
    
    def __init__(self, config: AgentSandboxConfig):
        self.config = config
        self.agent_id = config.agent_id
        self.session_id = config.session_id
        self.state = AgentState.IDLE
        
        # 进程间通信
        self.command_queue: Queue = mp.Queue()  # 父 → 子: 发送命令
        self.response_queue: Queue = mp.Queue() # 子 → 父: 返回结果
        self.heartbeat_queue: Queue = mp.Queue()# 子 → 父: 心跳
        
        self.process: Optional[Process] = None
        self.circuit_breaker = CircuitBreaker(
            config.max_consecutive_failures,
            config.cooldown_seconds
        )
        
        # 最后一次心跳的时间
        self.last_heartbeat = time.time()
    
    def start(self):
        """启动 Agent 子进程并绑定资源限制"""
        # 创建 cgroup
        cgroup_path = CgroupManager.create_agent_group(
            self.agent_id, self.config
        )
        
        # fork 子进程
        self.process = mp.get_context('spawn').Process(
            target=self._agent_main,
            name=f"agent-{self.agent_id}",
            args=(self.command_queue, self.response_queue, self.heartbeat_queue)
        )
        self.process.start()
        
        # 将子进程加入 cgroup
        CgroupManager.add_process_to_group(self.agent_id, self.process.pid)
        
        self.state = AgentState.IDLE
        return cgroup_path
    
    def execute(self, user_message: str, timeout: float = 300) -> str:
        """
        执行用户请求,带断路器和超时保护
        
        流程:
        1. 检查断路器状态 → 不通过则直接返回错误
        2. 发送消息到子进程
        3. 等待子进程响应(超时自动终止)
        4. 更新断路器状态
        5. 返回结果
        """
        # Step 1: 断路器检查
        if not self.circuit_breaker.allow_request():
            raise CircuitBreakerOpenError(
                f"Agent {self.agent_id} 已熔断 (连续失败 {self.config.max_consecutive_failures} 次)"
            )
        
        # Step 2: 检查子进程存活
        if not self.process or not self.process.is_alive():
            self.state = AgentState.DEAD
            raise AgentDeadError(f"Agent {self.agent_id} 进程已死亡")
        
        try:
            # Step 3: 发送执行命令
            self.command_queue.put({
                "type": "execute",
                "message": user_message,
                "timestamp": time.time()
            })
            self.state = AgentState.RUNNING
            
            # Step 4: 等待响应(带整体超时)
            response = self.response_queue.get(timeout=timeout)
            
            # Step 5: 处理响应
            if response.get("status") == "success":
                self.circuit_breaker.record_success()
                self.state = AgentState.IDLE
                return response.get("result", "")
            else:
                # 子进程返回错误(业务逻辑错误,非宕机)
                self.circuit_breaker.record_failure()
                self.state = AgentState.IDLE
                raise AgentExecutionError(response.get("error", "unknown error"))
        
        except Exception as e:
            # 超时或通信错误 → 记录失败
            self.circuit_breaker.record_failure()
            self.state = AgentState.IDLE
            
            # 如果连续失败达标,更新 CIrcuit 状态
            if self.circuit_breaker.state == "OPEN":
                self.state = AgentState.CIRCUIT_OPEN
            
            raise
    
    def _agent_main(self, cmd_queue: Queue, resp_queue: Queue, hb_queue: Queue):
        """
        子进程主循环
        独立进程中运行,通过队列与父进程通信
        崩溃不会影响父进程
        """
        # 设置进程标题(便于 ps 查看)
        try:
            import setproctitle
            setproctitle.setproctitle(f"agent-worker-{self.session_id}")
        except ImportError:
            pass
        
        # 初始化 Agent(在子进程中加载模型等)
        agent = self._init_agent()
        
        # 心跳线程
        heartbeat_stop = threading.Event()
        heartbeat_thread = threading.Thread(
            target=self._heartbeat_loop,
            args=(hb_queue, heartbeat_stop),
            daemon=True
        )
        heartbeat_thread.start()
        
        try:
            while True:
                # 阻塞等待父进程命令
                cmd = cmd_queue.get()
                
                if cmd["type"] == "execute":
                    try:
                        # 执行 Agent 推理
                        result = agent.invoke(cmd["message"])
                        resp_queue.put({
                            "status": "success",
                            "result": result,
                            "timestamp": time.time()
                        })
                    except Exception as e:
                        # 错误信息回传,但不崩溃
                        resp_queue.put({
                            "status": "error",
                            "error": str(e),
                            "error_type": type(e).__name__,
                            "timestamp": time.time()
                        })
                
                elif cmd["type"] == "shutdown":
                    break
                
                elif cmd["type"] == "health_check":
                    # 健康检查:返回当前状态
                    resp_queue.put({
                        "status": "healthy",
                        "memory_mb": self._get_memory_usage(),
                        "timestamp": time.time()
                    })
        
        finally:
            heartbeat_stop.set()
            heartbeat_thread.join(timeout=2)
    
    def _heartbeat_loop(self, hb_queue: Queue, stop_event: threading.Event):
        """心跳发送循环(子进程内)"""
        interval = self.config.heartbeat_interval
        while not stop_event.wait(interval):
            hb_queue.put({
                "agent_id": self.agent_id,
                "pid": os.getpid(),
                "timestamp": time.time(),
                "memory_mb": self._get_memory_usage()
            })
    
    def _init_agent(self):
        """初始化 Agent 实例(在子进程中调用)"""
        # 实际实现会初始化 LLM client、加载 prompt 模板等
        class SimpleAgent:
            def invoke(self, message: str) -> str:
                return f"Response to: {message}"
        return SimpleAgent()
    
    def _get_memory_usage(self) -> float:
        """获取当前进程的内存使用量 (MB)"""
        try:
            with open(f"/proc/{os.getpid()}/status") as f:
                for line in f:
                    if line.startswith("VmRSS:"):
                        # VmRSS: 实际占用的物理内存 (kB)
                        kb = int(line.split()[1])
                        return kb / 1024.0
        except Exception:
            pass
        return 0.0
    
    def stop(self, timeout: float = 10):
        """优雅关闭 Agent 进程"""
        if not self.process or not self.process.is_alive():
            return
        
        # 发送关闭命令
        try:
            self.command_queue.put({"type": "shutdown"}, timeout=2)
        except Exception:
            pass
        
        # 等待进程退出
        self.process.join(timeout=timeout)
        
        # 超时强制 kill
        if self.process.is_alive():
            self.process.terminate()
            self.process.join(timeout=2)
            
            if self.process.is_alive():
                self.process.kill()
                self.process.join(timeout=1)
        
        # 清理 cgroup
        CgroupManager.remove_agent_group(self.agent_id)
        self.state = AgentState.DEAD

# ========== 错误类型定义 ==========
class CircuitBreakerOpenError(Exception):
    """断路器打开错误"""
    pass

class AgentDeadError(Exception):
    """Agent 进程死亡错误"""
    pass

class AgentExecutionError(Exception):
    """Agent 执行错误(业务逻辑)"""
    pass

# ========== Agent 管理器:协调多个沙箱实例 ==========
class AgentPoolManager:
    """
    Agent 池管理器
    管理多个 AgentSandbox 实例的生命周期
    
    核心职责:
    1. 会话到 Agent 的映射 (session_id → sandbox)
    2. 健康监控:定期检查心跳、内存使用、CPU 使用
    3. 故障恢复:死掉的 Agent 自动创建新实例 + 恢复会话状态
    4. 资源回收:空闲过久的 Agent 自动销毁
    """
    
    def __init__(self, max_agents: int = 500):
        self.max_agents = max_agents
        self.sandboxes: dict[str, AgentSandbox] = {}  # session_id → sandbox
        
        # 启动健康监控后台线程
        self._health_monitor_stop = threading.Event()
        self._health_monitor = threading.Thread(
            target=self._health_monitor_loop,
            daemon=True,
            name="agent-health-monitor"
        )
        self._health_monitor.start()
    
    def create_agent(self, session_id: str) -> AgentSandbox:
        """为会话创建新的 Agent 沙箱"""
        if len(self.sandboxes) >= self.max_agents:
            raise RuntimeError(f"Agent 池已满 ({self.max_agents})")
        
        config = AgentSandboxConfig(
            agent_id=f"agent_{session_id}",
            session_id=session_id
        )
        
        sandbox = AgentSandbox(config)
        sandbox.start()
        self.sandboxes[session_id] = sandbox
        
        return sandbox
    
    def get_agent(self, session_id: str) -> Optional[AgentSandbox]:
        """获取会话对应的 Agent 沙箱"""
        return self.sandboxes.get(session_id)
    
    def _health_monitor_loop(self):
        """
        健康监控循环
        每秒检查一次所有 Agent 的健康状态
        """
        while not self._health_monitor_stop.wait(1.0):
            dead_sessions = []
            
            for session_id, sandbox in list(self.sandboxes.items()):
                # 检查 1: 进程存活
                if sandbox.process and not sandbox.process.is_alive():
                    dead_sessions.append(session_id)
                    continue
                
                # 检查 2: 心跳超时
                if sandbox.process and sandbox.process.is_alive():
                    elapsed = time.time() - sandbox.last_heartbeat
                    if elapsed > sandbox.config.heartbeat_timeout:
                        print(f"[MONITOR] Agent {session_id} 心跳超时 ({elapsed:.1f}s)")
                        dead_sessions.append(session_id)
                
                # 检查 3: 断路器状态
                if sandbox.circuit_breaker.state == "OPEN":
                    print(f"[MONITOR] Agent {session_id} 断路器打开,等待恢复")
            
            # 处理死掉的 Agent
            for session_id in dead_sessions:
                self._recover_agent(session_id)
    
    def _recover_agent(self, session_id: str):
        """
        恢复死掉的 Agent
        1. 从 Redis 读取 checkpoint
        2. 创建新 Agent 实例
        3. 回放对话历史
        4. 更新会话映射
        """
        print(f"[RECOVERY] 开始恢复 Agent {session_id}")
        
        # 清理旧实例
        if session_id in self.sandboxes:
            old_sandbox = self.sandboxes[session_id]
            try:
                old_sandbox.stop(timeout=5)
            except Exception:
                pass
            del self.sandboxes[session_id]
        
        # 恢复会话状态
        checkpoint = self._load_checkpoint(session_id)
        
        if checkpoint:
            # 创建新 Agent 并回放历史
            new_sandbox = self.create_agent(session_id)
            # 回放对话历史...
            print(f"[RECOVERY] Agent {session_id} 恢复成功")
        else:
            print(f"[RECOVERY] Agent {session_id} 无可恢复的状态,创建新实例")
            new_sandbox = self.create_agent(session_id)
        
        self.sandboxes[session_id] = new_sandbox
    
    def _load_checkpoint(self, session_id: str) -> Optional[dict]:
        """从 Redis 加载会话 checkpoint"""
        # 实际实现: redis.get(f"agent:checkpoint:{session_id}")
        return None
    
    def shutdown(self):
        """关闭所有 Agent"""
        self._health_monitor_stop.set()
        self._health_monitor.join(timeout=5)
        
        for sandbox in self.sandboxes.values():
            try:
                sandbox.stop(timeout=5)
            except Exception:
                pass
        
        self.sandboxes.clear()

if __name__ == "__main__":
    # 初始化 cgroup 根组
    CgroupManager.init_root_group()
    
    # 创建 Agent 池
    pool = AgentPoolManager(max_agents=100)
    
    # 模拟创建 5 个 Agent
    for i in range(5):
        session_id = f"session_{i}"
        sandbox = pool.create_agent(session_id)
        print(f"Agent {session_id} 已启动, PID: {sandbox.process.pid}")
    
    # 模拟执行
    try:
        result = pool.get_agent("session_0").execute("Hello, agent!")
        print(f"Result: {result}")
    except CircuitBreakerOpenError as e:
        print(f"断路器拦截: {e}")
    
    # 清理
    pool.shutdown()

四、边界分析与架构权衡

4.1 进程隔离 vs 线程隔离

大多数 Agent 框架(LangChain、AutoGPT)默认使用线程模型——一个进程内多个线程分别服务不同的会话。好处是内存共享(模型权重只加载一次),坏处是一个线程的 segfault 或死循环会 kill 整个进程。

策略:一个 GPU 节点跑 2-4 个主进程,每个主进程管理 10-20 个线程。进程间是强隔离(GPU 通过 MPS 共享),进程内是轻量隔离 + cgroup 限制。

4.2 cgroup 的内存超限行为

cgroup v2 的 memory.max 超过后不会直接 OOM Kill,而是先暂停 cgroup 内所有进程并触发 OOM 通知。应用层可以:

  1. 通过 memory.events 文件检测 OOM 事件
  2. 主动清理缓存或减少负载
  3. 如果 5 秒内未恢复,由内核 OOM Killer 介入

但很多 Agent 引擎对"被暂停"没有防御能力——它们在 invoke() 调用中同步等待 LLM 响应,被暂停后恢复可能导致状态不一致。

建议:将 cgroup 内存限制设为"软限制"(memory.high),超过后内核温和地回收内存,而不是硬停。

4.3 断路器误触发问题

连续 3 次失败就熔断听起来合理,但如果这 3 次失败都是因为同一个外部工具暂时不可用(比如 Google Search API 限流),熔断实际上是惩罚了正常的 Agent。

优化策略:区分"内部故障"和"外部故障":

  • 内部故障(segfault、OOM、死循环)→ 立即熔断 + 重建进程
  • 外部故障(工具超时、API 限流)→ 只降级不熔断,返回"工具暂不可用"

4.4 会话恢复(Session Recovery)的代价

Checkpoint 恢复不是免费的。每生成一个 token 就写一次 Redis(P99 延迟 +2ms),对于流式对话体验有明显影响。

异步 Checkpoint:在后台线程中批量写入,每 10 个 token 或每 5 秒写一次。丢失的最坏情况是用户需要重复最后 10 个 token——对于故障恢复来说是可以接受的代价。

五、总结

Agent 故障隔离的核心思路和微服务治理一脉相承,但更严苛——因为 Agent 的状态(对话历史)一旦丢失就是用户体验的断裂。

四条铁律:

  1. 进程级隔离是底线。不能接受"一个用户把整个服务搞崩"。
  2. cgroup 资源限制是守护线。CPU 和内存都要硬限制,超了就熔断而不是扩散。
  3. 断路器和心跳是双保险。断路器防止"坏 Agent"继续浪费资源,心跳确保能及时发现"死 Agent"。
  4. 会话恢复是业务连续性。故障不可避免,但恢复应该是无缝的——用户看到的应该是"刚刚在加载,现在好了"而不是"聊天记录全丢了"。

最后一句:Agent 平台的可靠性不取决于你的 Agent 多聪明,而取决于一个疯掉的 Agent 能不能被优雅地关进笼子里。

Logo

这里是“一人公司”的成长家园。我们提供从产品曝光、技术变现到法律财税的全栈内容,并连接云服务、办公空间等稀缺资源,助你专注创造,无忧运营。

更多推荐