Agent工作流的可观测性设计:Trace追踪、断点重试与异常熔断机制

发布时间:2026/7/24 20:37:15
Agent工作流的可观测性设计:Trace追踪、断点重试与异常熔断机制 Agent工作流的可观测性设计Trace追踪、断点重试与异常熔断机制一、Agent的黑盒困境为什么它卡住了在哪一步多Agent工作流在生产环境面临的核心困境是可观测性缺失执行失败时不知道在哪个Agent卡住、哪个工具的调用超时、为什么循环了8次还没终止。在一个生活任务处理Agent的故障排查中一条帮我规划周末行程的请求因为调用了天气API但未设置超时导致整个工作流挂起12分钟——而日志只显示执行中。可观测性不是锦上添花是Agent工作流生产化的必要条件。以下设计围绕三个核心能力分布式Trace追踪定位瓶颈、断点重试从失败步骤恢复而非重头执行、熔断降级保护系统免受级联故障。二、Agent可观测性的三层指标Trace层使用OpenTelemetry在Agent、工具和LLM调用三个粒度上创建Span构建调用链的全局视图。Metric层实时统计每个Agent的延迟、成功率和Token消耗——这些都是可设定告警阈值的指标。Log层记录决策过程Agent选择了工具X参数为Y原因是Z用于事后故障复盘。三、断点重试与熔断的关键实现# agent_observability/circuit_breaker.py Agent级熔断器 设计意图 1. 半开状态的设计熔断后定期探测服务恢复情况 2. 每个Agent独立熔断不影响其他Agent的正常工作 3. 熔断期间的请求直接返回fallback结果而非报错 import time from enum import Enum from dataclasses import dataclass from functools import wraps from typing import Callable, Any class CircuitState(Enum): CLOSED closed # 正常 OPEN open # 熔断中 HALF_OPEN half_open # 探测恢复 dataclass class CircuitConfig: failure_threshold: int 5 # 连续失败N次后熔断 recovery_timeout: float 30.0 # 熔断30秒后尝试恢复 half_open_probes: int 2 # 半开状态允许2次探测 class AgentCircuitBreaker: Agent熔断器 def __init__(self, agent_name: str, config: CircuitConfig | None None): self.agent_name agent_name self.config config or CircuitConfig() self.state CircuitState.CLOSED self.failure_count 0 self.last_failure_time 0.0 self.probe_count 0 def call(self, func: Callable, *args, **kwargs) - Any: 执行Agent调用熔断期间返回None if self.state CircuitState.OPEN: if self._should_attempt_recovery(): self.state CircuitState.HALF_OPEN self.probe_count 0 else: # 熔断中直接返回fallback return self._fallback_response() try: result func(*args, **kwargs) self._on_success() return result except Exception as exc: self._on_failure() raise exc def _on_success(self): 调用成功重置熔断器 self.failure_count 0 if self.state CircuitState.HALF_OPEN: self.state CircuitState.CLOSED def _on_failure(self): 调用失败累加计数必要时触发熔断 self.failure_count 1 self.last_failure_time time.time() if self.failure_count self.config.failure_threshold: self.state CircuitState.OPEN def _should_attempt_recovery(self) - bool: 判断是否到了尝试恢复的时间 return (time.time() - self.last_failure_time self.config.recovery_timeout) def _fallback_response(self) - dict: 熔断期间的降级响应 return { status: degraded, message: fAgent {self.agent_name} 暂时不可用请稍后再试, } # agent_observability/retry.py 断点重试管理器 设计意图 1. 工作流执行到步骤N失败时从失败步骤重试而非从头开始 2. 每个步骤的执行结果保存为checkpoint 3. 重试只重新执行失败的步骤和后续步骤 from typing import Any import json import hashlib class CheckpointManager: 工作流断点管理器 def __init__(self, workflow_id: str): self.workflow_id workflow_id self.checkpoints: dict[int, dict[str, Any]] {} self.completed_steps: set[int] set() def save(self, step_index: int, input_data: dict, output_data: dict): 保存步骤检查点 checkpoint { step: step_index, input_hash: self._hash(input_data), output: output_data, } self.checkpoints[step_index] checkpoint self.completed_steps.add(step_index) def resume_from(self) - int: 返回应该从哪个步骤恢复执行已完成步骤的最大index1 if not self.completed_steps: return 0 return max(self.completed_steps) 1 def get_checkpoint(self, step_index: int) - dict | None: 获取某个步骤的输出结果用于后续步骤的输入 cp self.checkpoints.get(step_index) return cp[output] if cp else None def _hash(self, data: dict) - str: 对输入数据计算哈希用于检测输入是否变化 serialized json.dumps(data, sort_keysTrue, defaultstr) return hashlib.md5(serialized.encode()).hexdigest()[:8]CheckpointManager在工作流的每个步骤完成后保存检查点。重试时跳过已完成的步骤直接从失败位置继续执行——对于有5个步骤的工作流步骤4失败只需重试步骤4和5而非重新执行全部5个步骤。四、可观测性的开销控制OpenTelemetry的Span采样率在生产环境中需要控制。全量trace100%采样在高峰期会产生每分钟约12万条Span——Jaeger的后端存储压力巨大。通过Head-based Sampling设置30%的采样率同时使用Tail-based Sampling确保错误trace 100%保留。日志存储是另一个隐形成本。每个Agent每次调用产生的完整日志包括Prompt和响应内容约3KB日调用10000次就是30MB的日增日志量。通过设置日志保留策略Info级别保留7天Debug级别保留24小时平衡问题排查的需求和存储成本。五、总结本次Agent工作流可观测性设计的核心结论Trace/Metric/Log三层覆盖从监控到排查的完整链路Trace定位瓶颈Metric触发告警Log辅助复盘。断点重试节省75%的重试耗时5步骤工作流在第4步失败时只需重试2步而非5步。Agent级独立熔断防止级联故障一个Agent挂掉不影响其他Agent的正常运行降级响应替代报错。采样策略平衡监控覆盖和存储成本30%正常trace采样100%错误trace保留。日志保留策略按级别分级Info 7天、Debug 24小时遏制日志存储的线性增长。