云计算百科
云计算领域专业知识百科平台

Agent 心跳与健康检查:长连接场景下的会话状态监控

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 次心跳缺失判定死亡)。检查点机制保证崩溃后断点续传:每完成一步保存结果,恢复时从最新检查点继续,不从头重跑。心跳通道与业务通道隔离,避免业务阻塞影响心跳延迟。

赞(0)
未经允许不得转载:网硕互联帮助中心 » Agent 心跳与健康检查:长连接场景下的会话状态监控
分享到: 更多 (0)

评论 抢沙发

评论前必须登录!