Agent 心跳与健康检查:长连接场景下的会话状态监控
Agent 连着连着就没了反应——你不知道它是真的在思考,还是已经悄悄挂了。
一、场景痛点
你的 Agent 系统用 WebSocket 维持长连接,用户发一条消息后 Agent 需要调用多个工具,耗时可能 10-60 秒。问题来了:这 60 秒里,用户不知道 Agent 是在处理还是已经挂了。如果 Agent 进程崩溃,WebSocket 连接不会自动断开——客户端一直等着,等到超时才报错,但此时用户已经以为系统卡死了。
你加了一个进度条,每 5 秒推送一条"还在处理中"的消息。但 Agent 真的挂了的时候,进度条还在转——因为推送线程和 Agent 执行线程是独立的,Agent 挂了推送线程还在跑。
更棘手的是会话恢复。Agent 处理到第 3 步时崩溃了,用户重连后从第 1 步重新开始——前两步的结果全部丢失。如果是付费场景(每次调用消耗 token),重复执行的成本直接翻倍。
核心矛盾:长连接场景下,Agent 的存活状态和处理进度必须被持续监控,不能靠"连接还在就以为还活着"。
二、底层机制与原理剖析
2.1 心跳机制的层次
2.2 心跳数据模型
进程心跳不只是"我还活着",它包含处理状态:
| 字段 | 含义 | 用途 |
|---|---|---|
agent_id | Agent 实例标识 | 关联会话与实例 |
session_id | 会话标识 | 恢复会话时使用 |
status | idle/processing/error | 区分状态 |
current_step | 当前执行步骤 | 进度追踪 |
total_steps | 总步骤数 | 进度百分比计算 |
cpu_percent | CPU 使用率 | 资源监控 |
memory_mb | 内存占用 | 资源监控 |
last_tool_call | 最近一次工具调用信息 | 卡住定位 |
timestamp | 心跳时间戳 | 判断是否过期 |
2.3 健康检查的判定逻辑
健康判定不是简单的"心跳在就健康"。需要根据心跳间隔和业务状态综合判断:
- 健康:心跳间隔 < 预期间隔 × 2,status 不是 error
- 亚健康:心跳间隔 > 预期间隔 × 2 但 < 预期间隔 × 5,status 是 processing 但 current_step 长时间不变
- 不健康:心跳间隔 > 预期间隔 × 5,或 status 是 error,或连续 3 次心跳缺失
三、生产级代码实现
3.1 Agent 心跳上报器
// agent-heartbeat.ts —— Agent 进程心跳上报器 import { EventEmitter } from 'events'; export enum AgentStatus { IDLE = 'idle', // 等待用户输入 PROCESSING = 'processing', // 正在处理用户请求 ERROR = 'error', // 处理出错,等待恢复 TERMINATING = 'terminating', // 正在优雅关闭 } export interface HeartbeatPayload { agent_id: string; session_id: string; status: AgentStatus; current_step: number; total_steps: number; cpu_percent: number; memory_mb: number; last_tool_call: string | null; timestamp: number; // Unix 时间戳(毫秒) } export class AgentHeartbeat extends EventEmitter { private agentId: string; private sessionId: string; private status: AgentStatus = AgentStatus.IDLE; private currentStep: number = 0; private totalSteps: number = 0; private lastToolCall: string | null = null; private intervalMs: number; // 心跳间隔 private maxMissedHeartbeats: number; // 允许缺失的最大心跳数 private heartbeatTimer: NodeJS.Timeout | null = null; // 心跳发送通道:WebSocket / HTTP / 消息队列 private sender: (payload: HeartbeatPayload) => Promise<void>; constructor( agentId: string, sessionId: string, intervalMs: number = 5000, // 默认 5 秒心跳间隔 maxMissedHeartbeats: number = 3, sender: (payload: HeartbeatPayload) => Promise<void> ) { super(); this.agentId = agentId; this.sessionId = sessionId; this.intervalMs = intervalMs; this.maxMissedHeartbeats = maxMissedHeartbeats; this.sender = sender; } /** 启动心跳循环:定时上报状态 */ start(): void { if (this.heartbeatTimer) return; // 已启动则不重复启动 // 定时发送心跳:intervalMs 间隔 // 心跳是"推"模式,不是"拉"模式——服务端不需要轮询检查 Agent 状态 this.heartbeatTimer = setInterval(() => { this.sendHeartbeat(); }, this.intervalMs); // 立即发送一次心跳:启动时让服务端知道 Agent 已上线 this.sendHeartbeat(); } /** 停止心跳循环:优雅关闭前调用 */ stop(): void { if (this.heartbeatTimer) { clearInterval(this.heartbeatTimer); this.heartbeatTimer = null; } // 发送最终心跳:标记为 terminating,服务端知道 Agent 正在关闭 this.status = AgentStatus.TERMINATING; this.sendHeartbeat(); } /** 更新处理状态:Agent 每完成一步调用此方法 */ updateProgress(currentStep: number, totalSteps: number, toolCall: string): void { this.currentStep = currentStep; this.totalSteps = totalSteps; this.lastToolCall = toolCall; this.status = AgentStatus.PROCESSING; // 状态变化时立即发送一次心跳(不等定时器) // 用户在等待结果,状态变化应该第一时间告知服务端 this.sendHeartbeat(); } /** 标记错误状态 */ markError(): void { this.status = AgentStatus.ERROR; this.sendHeartbeat(); } /** 标记空闲状态 */ markIdle(): void { this.status = AgentStatus.IDLE; this.currentStep = 0; this.totalSteps = 0; this.lastToolCall = null; this.sendHeartbeat(); } /** 发送心跳:组装 payload 并调用 sender */ private async sendHeartbeat(): void { const payload: HeartbeatPayload = { agent_id: this.agentId, session_id: this.sessionId, status: this.status, current_step: this.currentStep, total_steps: this.totalSteps, cpu_percent: this.getCpuUsage(), memory_mb: this.getMemoryUsage(), last_tool_call: this.lastToolCall, timestamp: Date.now(), }; try { await this.sender(payload); this.emit('heartbeat:sent', payload); } catch (err) { // 心跳发送失败:不中断 Agent 处理流程 // 心跳是辅助功能,核心业务不能因为心跳通道故障而停止 this.emit('heartbeat:failed', { error: err, payload }); } } /** 获取 CPU 使用率:简化实现,生产环境用 process.cpuUsage() */ private getCpuUsage(): number { // Node.js 的 process.cpuUsage() 返回微秒级的 CPU 时间 const usage = process.cpuUsage(); const totalUsec = usage.user + usage.system; // 转换为百分比(近似值,需要采样间隔才能精确计算) return Math.min(totalUsec / 1000 / this.intervalMs, 100); } /** 获取内存使用量 */ private getMemoryUsage(): number { return process.memoryUsage().heapUsed / 1024 / 1024; // MB } }3.2 服务端健康检查监控器
// health-monitor.ts —— 服务端 Agent 健康检查监控器 // 监控所有 Agent 实例的心跳,判断健康状态,触发告警和会话恢复 export enum HealthStatus { HEALTHY = 'healthy', DEGRADED = 'degraded', UNHEALTHY = 'unhealthy', DEAD = 'dead', } interface AgentHealthRecord { agentId: string; sessionId: string; lastHeartbeat: HeartbeatPayload; lastHeartbeatTime: number; missedHeartbeats: number; healthStatus: HealthStatus; // 停滞检测:current_step 连续 N 次心跳未变化 stagnantCount: number; } export class AgentHealthMonitor extends EventEmitter { private agents: Map<string, AgentHealthRecord> = new Map(); private intervalMs: number; private maxMissed: number; private stagnantThreshold: number; // 心跳停滞阈值:step 不变的次数 private checkTimer: NodeJS.Timeout | null = null; constructor( intervalMs: number = 10000, // 每 10 秒检查一次所有 Agent maxMissed: number = 3, stagnantThreshold: number = 5 // 5 次心跳 step 不变判定为停滞 ) { super(); this.intervalMs = intervalMs; this.maxMissed = maxMissed; this.stagnantThreshold = stagnantThreshold; } /** 接收 Agent 心跳:更新健康记录 */ receiveHeartbeat(payload: HeartbeatPayload): void { const existing = this.agents.get(payload.agent_id); if (existing) { // 检查 current_step 是否变化:停滞检测 if (payload.current_step === existing.lastHeartbeat.current_step && payload.status === AgentStatus.PROCESSING) { existing.stagnantCount++; } else { existing.stagnantCount = 0; } // 更新记录 existing.lastHeartbeat = payload; existing.lastHeartbeatTime = payload.timestamp; existing.missedHeartbeats = 0; // 重新评估健康状态 this.evaluateHealth(existing); } else { // 新 Agent 上线:初始化健康记录 this.agents.set(payload.agent_id, { agentId: payload.agent_id, sessionId: payload.session_id, lastHeartbeat: payload, lastHeartbeatTime: payload.timestamp, missedHeartbeats: 0, healthStatus: HealthStatus.HEALTHY, stagnantCount: 0, }); this.emit('agent:registered', { agentId: payload.agent_id }); } } /** 启动健康检查循环 */ start(): void { this.checkTimer = setInterval(() => { this.checkAllAgents(); }, this.intervalMs); } /** 检查所有 Agent 的健康状态 */ private checkAllAgents(): void { const now = Date.now(); const expectedInterval = 5000; // Agent 心跳间隔 for (const [agentId, record] of this.agents) { const elapsed = now - record.lastHeartbeatTime; // 心跳缺失检测:超过预期间隔则计数 +1 if (elapsed > expectedInterval * 2) { record.missedHeartbeats++; } // 停滞检测:step 不变的次数超过阈值 if (record.stagnantCount >= this.stagnantThreshold) { // 工具调用卡住:Agent 还活着但处理停滞 this.emit('agent:stagnant', { agentId, sessionId: record.sessionId, currentStep: record.lastHeartbeat.current_step, lastToolCall: record.lastHeartbeat.last_tool_call, }); } // 重新评估健康状态 this.evaluateHealth(record); // 不健康或死亡:触发告警 if (record.healthStatus === HealthStatus.UNHEALTHY) { this.emit('agent:unhealthy', { agentId, sessionId: record.sessionId, missedHeartbeats: record.missedHeartbeats, }); } if (record.healthStatus === HealthStatus.DEAD) { this.emit('agent:dead', { agentId, sessionId: record.sessionId, }); // 死亡 Agent 从监控列表移除:不再等待心跳 // 但会话状态保留,用于后续恢复 this.agents.delete(agentId); } } } /** 评估单个 Agent 的健康状态 */ private evaluateHealth(record: AgentHealthRecord): void { const previousStatus = record.healthStatus; if (record.missedHeartbeats >= this.maxMissed * 2) { // 连续缺失超过 2 倍阈值:判定死亡 record.healthStatus = HealthStatus.DEAD; } else if (record.missedHeartbeats >= this.maxMissed) { // 连续缺失超过阈值:判定不健康 record.healthStatus = HealthStatus.UNHEALTHY; } else if (record.missedHeartbeats > 0 || record.stagnantCount >= this.stagnantThreshold) { // 有缺失但未超阈值,或处理停滞:亚健康 record.healthStatus = HealthStatus.DEGRADED; } else { // 正常心跳且处理推进中:健康 record.healthStatus = HealthStatus.HEALTHY; } // 状态变化时发出事件:外部系统可以订阅做自动恢复 if (previousStatus !== record.healthStatus) { this.emit('health:changed', { agentId: record.agentId, from: previousStatus, to: record.healthStatus, }); } } }3.3 会话恢复与断点续传
# session_recovery.py —— Agent 崩溃后的会话恢复与断点续传 import json import logging import time from datetime import datetime logger = logging.getLogger('session-recovery') class SessionRecoveryManager: """会话恢复管理器:Agent 崩溃后从断点继续处理""" def __init__(self, storage_client, heartbeat_monitor): self.storage = storage_client self.monitor = heartbeat_monitor def save_checkpoint(self, session_id: str, step_index: int, step_results: dict): """保存检查点:每完成一步就保存,崩溃后从检查点恢复""" checkpoint = { 'session_id': session_id, 'step_index': step_index, 'step_results': step_results, 'timestamp': datetime.utcnow().isoformat(), } # 检查点存到对象存储:比数据库更快,且不影响业务表 key = f"checkpoints/{session_id}/step_{step_index}.json" self.storage.put(key, json.dumps(checkpoint)) def recover_session(self, session_id: str) -> dict: """从最新检查点恢复会话""" # 查找该会话的所有检查点,取最新的 pattern = f"checkpoints/{session_id}/step_*.json" checkpoints = self.storage.list(pattern) if not checkpoints: logger.warning(f"No checkpoints found for session {session_id}") return {'step_index': 0, 'step_results': {}} # 取最新的检查点(step_index 最大的) latest = max(checkpoints, key=lambda k: int(k.split('step_')[1].split('.')[0])) checkpoint_data = self.storage.get(latest) checkpoint = json.loads(checkpoint_data) logger.info( f"Recovered session {session_id} from step {checkpoint['step_index']}" ) return checkpoint def handle_dead_agent(self, agent_id: str, session_id: str): """处理死亡 Agent:恢复会话并分配新 Agent""" # 1. 从检查点恢复会话状态 checkpoint = self.recover_session(session_id) # 2. 创建新 Agent 实例,传入恢复的检查点 # 新 Agent 从断点继续,不从第 0 步重新开始 new_agent = self.create_agent_with_checkpoint( session_id, checkpoint ) # 3. 通知用户:会话恢复,从第 N 步继续 logger.info( f"Session {session_id} recovered: " f"new agent {new_agent.agent_id}, " f"resuming from step {checkpoint['step_index']}" ) # 4. 清理旧 Agent 的残留资源(内存中的会话数据等) self.cleanup_agent_resources(agent_id) return new_agent def create_agent_with_checkpoint(self, session_id: str, checkpoint: dict): """创建新 Agent 并注入检查点数据""" # 新 Agent 启动时接收检查点, # 从 checkpoint['step_index'] + 1 开始执行 # 前面步骤的结果从 checkpoint['step_results'] 中读取 agent_config = { 'session_id': session_id, 'resume_from_step': checkpoint['step_index'] + 1, 'previous_results': checkpoint['step_results'], } # 调用 Agent 启动接口 return self.start_new_agent(agent_config) def cleanup_agent_resources(self, agent_id: str): """清理死亡 Agent 的残留资源""" # 释放内存中的会话缓存、关闭未完成的工具调用连接等 logger.info(f"Cleaning up resources for dead agent {agent_id}")四、边界分析与架构权衡
4.1 心跳间隔的权衡
心跳间隔太短(1 秒):网络开销大,服务端处理压力大。心跳间隔太长(30 秒):Agent 挂了 30 秒你才知道,用户已经等了 30 秒才发现系统没响应。
折中:基础心跳 5 秒(覆盖大部分场景),状态变化时立即发送一次即时心跳。这样正常情况下每 5 秒一次心跳,状态变化时秒级感知。
4.2 检查点的存储频率
每完成一步就保存检查点,意味着每步都有一次存储写入。如果步骤执行很快(每步 1 秒),写入频率就是每秒一次。对象存储的写入延迟约 50-100ms,不影响步骤执行。
但如果步骤执行很慢(每步 10 秒),保存频率是每 10 秒一次,崩溃后最多丢失 10 秒的工作量。
4.3 适用边界与禁用场景
- 适用:WebSocket/SSE 长连接的 Agent 系统、多步骤工具调用链路、需要会话恢复的付费场景
- 禁用:单次请求-响应的简单 Agent(不需要心跳)、短连接 HTTP API(不需要长连接监控)、离线批处理 Agent(不需要实时状态)
4.4 心跳通道与业务通道的隔离
心跳消息和业务消息走同一个 WebSocket 连接时,如果业务消息阻塞(比如大结果传输),心跳也会延迟。解决方案:心跳走独立连接或独立的消息类型(WebSocket 的 ping 帧与数据帧是独立的)。
五、总结
Agent 长连接场景的健康检查需要三层心跳:连接层检测网络可达、进程层检测 Agent 存活与资源状态、业务层检测处理进度是否推进。单靠连接层心跳无法区分"在思考"和"已挂掉"。核心设计:5 秒基础心跳 + 状态变化即时心跳、停滞检测(step 不变的次数超阈值判定卡住)、缺失检测(连续 N 次心跳缺失判定死亡)。检查点机制保证崩溃后断点续传:每完成一步保存结果,恢复时从最新检查点继续,不从头重跑。心跳通道与业务通道隔离,避免业务阻塞影响心跳延迟。