AI Agent 从 Chains 到 Systems:Tool Calling、可观测性与生产稳定性架构实战
一、从 Chain 到 System:复杂度爆炸的根因
传统 RAG 链路是典型的 Chain 模式:检索 → 拼接上下文 → 一次 LLM 调用 → 返回结果。它的状态空间小、失败模式单一(检索为空或超时),容易工程化。
但真正的 Agent 是 System:模型自主决定调用哪些工具、调用几次、是否回退重试。这意味着:
- 状态空间指数级膨胀:一个允许 8 次工具调用的 Agent,若每次有 5 种工具可选,理论路径数为 5⁸ = 390,625 条。
- 失败模式不可枚举:工具超时、返回格式不符合 schema、API 限流、模型幻觉输出虚假参数、工具结果互相矛盾。
- 失败成本不可控:每多一次 Token 调用都在计费,一个失控的 Agent 循环可能几分钟内烧掉数百美元。
生产中的典型故障链路如下:
用户请求 → Agent 调用 SearchAPI → 返回空结果 → Agent 再次调用 SearchAPI(修改参数)
→ API 限流 429 → Agent 误判为永久失败 → 切换到 FallbackDB
→ FallbackDB 返回过期数据 → Agent 基于过期数据生成错误结论
→ 用户投诉 → 运维发现时已产生 143 次工具调用、$47 费用
这个链路涉及了 5 种失败模式:重试风暴、限流感知缺失、幻觉切换、数据过期无感知、成本失控。每一类都需要独立的基础设施保障。
二、Tool Calling 的可靠性保障架构
2.1 分层重试与指数退避
工具调用失败分两类:瞬时故障(网络抖动、503)应重试;永久故障(400 参数错误、404 资源不存在)必须立即终止并上报模型。
import asyncio
import time
from dataclasses import dataclass, field
from enum import Enum
from typing import Callable, Any, Optional
class FailureMode(Enum):
TRANSIENT = "transient" # 可重试
PERMANENT = "permanent" # 永久失败
RATE_LIMIT = "rate_limit" # 限流,需等待后重试
AMBIGUOUS = "ambiguous" # 不确定,需模型重新判断
@dataclass
class ToolResult:
success: bool
data: Any = None
error_detail: str = ""
failure_mode: FailureMode = FailureMode.TRANSIENT
retry_after_seconds: float = 0.0
attempts_used: int = 0
@dataclass
class ToolCallConfig:
max_retries: int = 3
base_backoff_ms: int = 500 # 起始退避
max_backoff_ms: int = 30_000 # 最大 30s
jitter_factor: float = 0.3 # 抖动系数,避免重试风暴
circuit_breaker_threshold: int = 5 # 连续失败 5 次触发熔断
circuit_breaker_timeout_s: int = 60 # 熔断恢复时间
cost_budget_usd: float = 1.0 # 单次工具调用链的成本上限
timeout_per_call_ms: int = 15_000 # 单次超时
class CircuitBreaker:
"""三态熔断器:CLOSED → OPEN → HALF_OPEN"""
CLOSED = "closed"
OPEN = "open"
HALF_OPEN = "half_open"
def __init__(self, threshold: int, recovery_s: int):
self.threshold = threshold
self.recovery_s = recovery_s
self.state = self.CLOSED
self.consecutive_failures = 0
self.last_failure_time: float = 0
def record_success(self):
self.consecutive_failures = 0
self.state = self.CLOSED
def record_failure(self) -> bool:
self.consecutive_failures += 1
self.last_failure_time = time.time()
if self.consecutive_failures >= self.threshold:
self.state = self.OPEN
return True # 刚触发熔断
return False
def can_execute(self) -> bool:
if self.state == self.CLOSED:
return True
if self.state == self.OPEN:
if time.time() - self.last_failure_time > self.recovery_s:
self.state = self.HALF_OPEN
return True
return False
return True # HALF_OPEN: 允许试探请求
class ResilientToolWrapper:
"""带熔断、退避、重试的工具调用包装器"""
def __init__(self, name: str, config: ToolCallConfig):
self.name = name
self.config = config
self.cb = CircuitBreaker(
threshold=config.circuit_breaker_threshold,
recovery_s=config.circuit_breaker_timeout_s
)
async def execute_with_resilience(self, fn: Callable, *args, **kwargs) -> ToolResult:
if not self.cb.can_execute():
return ToolResult(
success=False,
failure_mode=FailureMode.PERMANENT,
error_detail=f"Circuit breaker OPEN for {self.name}"
)
last_error = None
for attempt in range(self.config.max_retries + 1):
try:
result = await asyncio.wait_for(
fn(*args, **kwargs),
timeout=self.config.timeout_per_call_ms / 1000
)
self.cb.record_success()
return ToolResult(success=True, data=result, attempts_used=attempt + 1)
except asyncio.TimeoutError:
last_error = "timeout"
mode = FailureMode.TRANSIENT
except RateLimitError as e:
last_error = str(e)
mode = FailureMode.RATE_LIMIT
# 使用服务端返回的 retry-after,或默认退避
wait = e.retry_after or self._calc_backoff(attempt)
await asyncio.sleep(wait)
continue
except PermanentAPIError as e:
# 永久错误:不重试,直接上报
self.cb.record_failure()
return ToolResult(
success=False, failure_mode=FailureMode.PERMANENT,
error_detail=str(e), attempts_used=attempt + 1
)
except Exception as e:
last_error = str(e)
mode = FailureMode.AMBIGUOUS
# 瞬时/模糊错误:退避后重试
if attempt < self.config.max_retries:
wait = self._calc_backoff(attempt)
await asyncio.sleep(wait)
# 重试耗尽
triggered = self.cb.record_failure()
detail = f"Retries exhausted. Last error: {last_error}"
if triggered:
detail += " | Circuit breaker OPENED"
return ToolResult(
success=False, failure_mode=mode,
error_detail=detail, attempts_used=self.config.max_retries + 1
)
def _calc_backoff(self, attempt: int) -> float:
"""指数退避 + 全抖动"""
import random
exp = min(
self.config.max_backoff_ms,
self.config.base_backoff_ms * (2 ** attempt)
)
jitter = random.uniform(0, exp * self.config.jitter_factor)
return (exp + jitter) / 1000.0
关键设计决策:
- 熔断器使用三态模型(HALF_OPEN 允许一个试探请求验证恢复),避免在故障未完全恢复时直接恢复全量。
- 抖动因子(jitter)对重试风暴的防护至关重要:当 1000 个 Agent 实例同时遇到限流,无抖动必然导致 Thundering Herd。
- 永久失败(参数 400 类错误)必须跳过重试直接上报模型,因为重试只会浪费成本,必须让模型自己修正参数。
2.2 成本预算与紧急停止
Agent 在生产中最危险的故障不是"答错了",而是死循环调用导致成本失控。硬性限制必须由基础设施层强制实施,而非依赖模型自身判断。
@dataclass
class AgentBudget:
max_tool_calls: int = 50 # 单次任务最多工具调用次数
max_total_tokens: int = 100_000 # 单次任务 Token 总量
max_wall_time_s: int = 120 # 单次任务最大运行时间
max_tool_cost_usd: float = 5.0 # 单个工具调用链费用上限
cost_per_1k_tokens: float = 0.01 # 按 gpt-4o-mini input 计费
# 预算消耗 80% 时触发降级(切换为更便宜的模型或简化输出)
degradation_threshold: float = 0.8
class BudgetGuard:
"""与 Agent runtime 嵌入的预算守卫"""
def __init__(self, budget: AgentBudget):
self.budget = budget
self.tool_calls_made = 0
self.tokens_consumed = 0
self.start_time = time.time()
self._emergency_stop = False
def check_tool_call(self, estimated_tokens: int = 500) -> dict:
"""每次工具调用前检查预算,返回决策指令"""
self.tool_calls_made += 1
elapsed = time.time() - self.start_time
decisions = {"allow": True, "action": "continue", "reason": ""}
if self.tool_calls_made > self.budget.max_tool_calls:
decisions = {
"allow": False,
"action": "terminate",
"reason": f"tool_calls limit ({self.budget.max_tool_calls}) exceeded"
}
elif elapsed > self.budget.max_wall_time_s:
decisions = {
"allow": False,
"action": "terminate",
"reason": f"wall time limit ({self.budget.max_wall_time_s}s) exceeded"
}
elif self._degradation_triggered(estimated_tokens):
decisions = {
"allow": True,
"action": "degrade",
"reason": "budget 80% consumed, switching to lighter model"
}
if not decisions["allow"]:
self._emergency_stop = True
return decisions
def _degradation_triggered(self, estimated_tokens: int) -> bool:
token_ratio = (self.tokens_consumed + estimated_tokens) / self.budget.max_total_tokens
time_ratio = (time.time() - self.start_time) / self.budget.max_wall_time_s
return max(token_ratio, time_ratio) >= self.budget.degradation_threshold
def record_tokens(self, input_tokens: int, output_tokens: int):
cost = (input_tokens + output_tokens) / 1000 * self.budget.cost_per_1k_tokens
self.tokens_consumed += input_tokens + output_tokens
return cost
这个守卫不是"建议",而是硬拦截。即使模型认为必须继续执行,runtime 也必须在预算耗尽时强制终止并返回降级响应。
三、Agent 可观测性:超越传统 APM 的追踪体系
3.1 Agent Trace 的数据模型
传统 APM 追踪的是 HTTP 请求的生命周期,而 Agent 的追踪核心是"思维链的可观测性"——模型在每一步做决策的原因、工具选择逻辑、废弃的推理路径。
Trace (session_id=abc123)
├── Span: "用户请求 session_start" [ROOT]
├── Span: "Agent 思考: 分析需求类型" [LLM_CALL, tokens=1200/800]
├── Span: "Tool: web_search(query='...')" [TOOL_CALL, status=success, 230ms]
├── Span: "Tool: read_file(path='...')" [TOOL_CALL, status=error, 50ms]
├── Span: "Agent 思考: 文件不存在, 尝试替代方案" [LLM_CALL, tokens=900/600]
├── Span: "Tool: list_files(dir='...')" [TOOL_CALL, status=success, 180ms]
├── Span: "Agent 思考: 找到目标文件" [LLM_CALL, tokens=400/350]
├── Span: "Tool: read_file(path='final.md')"[TOOL_CALL, status=success, 45ms]
└── Span: "Agent 思考: 生成回答" [LLM_CALL, tokens=300/1500]
每个 Span 必须携带的关键属性:
{
"trace_id": "abc123",
"span_id": "span_007",
"kind": "TOOL_CALL",
"name": "read_file",
"status": "ERROR",
"attributes": {
"tool.name": "read_file",
"tool.call_id": "call_20261004_0830",
"tool.arguments": {"path": "/tmp/missing.txt"},
"tool.result_summary": "FileNotFoundError",
"tool.latency_ms": 52,
"tool.cost_usd": 0.0,
"llm.model": "gpt-4o-mini",
"llm.input_tokens": 1200,
"llm.output_tokens": 800,
"agent.decision_reason": "\"文件不存在, 尝试 list_files 找替代路径\"",
"agent.retry_count": 0,
"budget.tool_calls_remaining": 43,
"budget.tokens_remaining": 85000,
"error.type": "FileNotFoundError",
"error.recoverable": false
}
}
注意 agent.decision_reason 和 budget.tool_calls_remaining——这是 Agent 追踪区别于传统 APM 的关键字段。前者记录了模型在失败后的推理逻辑(对事后调试至关重要),后者记录了预算消耗曲线(对成本归因至关重要)。
3.2 基于 Agent Trace 的异常检测
Agent 系统的异常检测不能套用基于 P99 的阈值告警。我们需要识别的是行为序列异常——例如"模型连续 3 次调用同一工具但每次都失败,却没有尝试替代方案"。
from collections import deque
from dataclasses import dataclass
from typing import List, Callable
import re
@dataclass
class AnomalyRule:
name: str
description: str
severity: "info" | "warning" | "critical"
detector: Callable[[List] -> bool] # 接收 Trace spans,返回是否命中
action: str # 触发后的操作
def detect_stuck_loop(spans: List) -> bool:
"""检测模型陷入工具调用死循环的特征"""
if len(spans) < 6:
return False
recent = spans[-6:]
# 同一工具连续调用 ≥ 4 次,且参数相似度 > 80%
same_tool_count = sum(1 for s in recent[-4:] if s.name == recent[-1].name)
if same_tool_count < 4:
return False
# 检查参数是否几乎相同(去掉时间戳后对比)
params = [str(s.attributes.get("tool.arguments", "")) for s in recent[-4:]]
simplified = [re.sub(r'"timestamp"\s*:\s*\d+', '"timestamp":0', p) for p in params]
return len(set(simplified)) == 1 # 4 次参数完全一致
def detect_cost_anomaly(spans: List) -> bool:
"""检测 Token 消耗异常加速"""
if len(spans) < 3:
return False
# 最近 3 轮 LLM 调用的 output_tokens 是否逐次翻倍
llm_spans = [s for s in spans if s.kind == "LLM_CALL"][-3:]
if len(llm_spans) < 3:
return False
return (llm_spans[1].attributes > 2 * llm_spans[0].attributes and
llm_spans[2].attributes > 2 * llm_spans[1].attributes)
def detect_irrelevant_tool_selection(spans: List) -> bool:
"""检测模型选择了与用户意图无关的工具"""
if not spans:
return False
root_intent = spans[0].attributes.get("user.intent", "")
recent_tools = [s.name for s in spans[-4:] if s.kind == "TOOL_CALL"]
if len(recent_tools) < 3:
return False
# 最近连续 3 个工具与用户原始意图的语义相似度都低
irrelevant_count = sum(1 for t in recent_tools if intent_tool_similarity(root_intent, t) < 0.3)
return irrelevant_count >= 3
# 异常规则引擎
AGENT_ANOMALY_RULES = [
AnomalyRule(
name="agent_stuck_loop",
description="模型陷入工具调用死循环(连续 4+ 次同工具同参数)",
severity="critical",
detector=detect_stuck_loop,
action="terminate_and_alert"
),
AnomalyRule(
name="token_explosion",
description="Token 消耗每轮翻倍,疑似生成失控",
severity="warning",
detector=detect_cost_anomaly,
action="trigger_degradation"
),
AnomalyRule(
name="intent_drift",
description="工具选择偏离用户原始意图",
severity="warning",
detector=detect_irrelevant_tool_selection,
action="inject_correction_prompt"
),
]
这类基于行为序列的异常检测,比单纯监控 P99 延迟要有效得多 — 传统 APM 能告诉我们"Agent 平均响应时间从 2 秒涨到了 8 秒",但行为序列分析能告诉我们"模型正在第 7 次尝试用同一个错误参数调用 API,预算还剩 3%,建议立即介入"。
四、多 Agent 编排:Supervisor-Worker 模式的工程实现
当单个 Agent 无法胜任复杂任务时,需要引入多 Agent 编排。生产中最稳健的拓扑是 Supervisor-Worker 分层模式:
Supervisor Agent (Orchestrator)
│
├── Worker Agent "代码检索" ← 负责搜索/读取代码库
├── Worker Agent "缺陷分析" ← 负责定位 Bug 根因
├── Worker Agent "修复生成" ← 负责输出 Patch
└── Worker Agent "测试验证" ← 负责运行测试并确认修复
Supervisor 职责:
1. 任务分解与分发
2. 汇聚 Worker 结果
3. 处理 Worker 失败(重试/换 Worker/降级)
4. 合并最终输出
关键工程挑战是 Worker 间的状态传递与失败隔离。常见反模式是一个 Worker 的失败"传染"给整个任务。
import asyncio
from dataclasses import dataclass
from typing import Dict, Optional
from enum import Enum
class WorkerStatus(Enum):
PENDING = "pending"
RUNNING = "running"
SUCCESS = "success"
FAILED = "failed"
DEGRADED = "degraded" # 降级成功(部分结果可用)
@dataclass
class WorkerOutput:
agent_name: str
status: WorkerStatus
result: Any
error: Optional[str] = None
partial_result: Any = None # 降级时的部分可用结果
cost_usd: float = 0.0
class SupervisorAgent:
"""多 Agent 编排器:管理 Worker 生命周期与失败隔离"""
def __init__(self, workers: Dict, max_total_time_s: int = 300):
self.workers = workers
self.max_total_time_s = max_total_time_s
self.results: Dict = {}
self.completed = set()
async def execute_task(self, task: str) -> Dict:
start = time.time()
# Phase 1: 任务分解(由 Supervisor 模型完成)
subtasks = await self._decompose_task(task)
# Phase 2: 并发启动独立 Worker(有依赖的等待上游完成)
worker_tasks = {
name: asyncio.create_task(
self._run_worker_with_guard(name, subtask)
)
for name, subtask in subtasks.items()
}
# Phase 3: 等待完成,设置整体超时
remaining = self.max_total_time_s - (time.time() - start)
done, pending = await asyncio.wait(
worker_tasks.values(),
timeout=remaining,
return_when=asyncio.ALL_COMPLETED
)
# 处理超时未完成
for t in pending:
t.cancel()
# Phase 4: 汇总结果,处理 Worker 失败
for name, future in worker_tasks.items():
try:
output = await future
self.results[name] = output
if output.status in (WorkerStatus.SUCCESS, WorkerStatus.DEGRADED):
self.completed.add(name)
except Exception as e:
self.results[name] = WorkerOutput(
agent_name=name, status=WorkerStatus.FAILED,
result=None, error=str(e)
)
# Phase 5: 判断是否可以"降级完成"(部分 Worker 失败但不影响核心输出)
return self._aggregate_results(subtasks)
async def _run_worker_with_guard(self, name: str, subtask: str) -> WorkerOutput:
"""每个 Worker 独立隔离:一个崩溃不影响其他"""
worker = self.workers[name]
budget = AgentBudget(max_wall_time_s=60)
try:
async with worker.run(subtask, budget=budget) as session:
result = await session.get_result()
return WorkerOutput(
agent_name=name, status=WorkerStatus.SUCCESS,
result=result, cost_usd=session.total_cost
)
except WorkerCriticalError as e:
# 关键失败:尝试用降级策略产出部分结果
if worker.supports_partial:
partial = await worker.get_partial_output()
return WorkerOutput(
agent_name=name, status=WorkerStatus.DEGRADED,
result=None, partial_result=partial,
error=str(e)
)
return WorkerOutput(
agent_name=name, status=WorkerStatus.FAILED,
result=None, error=str(e)
)
def _aggregate_results(self, subtasks: Dict) -> Dict:
"""智能汇总:区分"全部成功"与"部分可用"""
critical_workers = [w for w in self.results if subtasks[w].get("critical")]
failed_critical = [w for w in critical_workers if self.results[w].status == WorkerStatus.FAILED]
if failed_critical:
# 关键 Worker 失败:不能完成,返回可解释的错误
return {
"status": "failed",
"reason": f"Critical workers failed: {failed_critical}",
"partial_results": {k: v for k, v in self.results.items()},
"suggestion": self._generate_recovery_suggestion(failed_critical)
}
# 所有关键 Worker 完成:组装最终结果,降级 Worker 的 partial_result 可参与合并
final = {}
for name, output in self.results.items():
final[name] = output.result or output.partial_result
return {
"status": "completed",
"degraded": any(v.status == WorkerStatus.DEGRADED for v in self.results.values()),
"results": final,
"total_cost_usd": sum(v.cost_usd for v in self.results.values())
}
这个设计的关键洞察是:Worker 的失败不应该自动传染。通过区分 critical 和 non-critical Worker,以及支持 partial_result 机制,我们可以在 80% Worker 失败时仍然产出"降级可用"的结果,而不是简单返回"任务失败"。
五、Agent 错误累积与自愈机制
5.1 错误分类与归档
在生产系统中,Agent 的错误不应被"吞掉",而应被结构化处理,形成自我改进的正反馈循环。
from datetime import datetime
class AgentErrorClassifier:
"""对 Agent 失败 trace 进行根因分类,驱动后续改进"""
ERROR_CATEGORIES = {
"MODEL_HALLUCINATION": {
"description": "模型输出虚假参数或幻觉工具",
"auto_retry": False, # 重试无用,需用更强模型或修正提示
"suggestion": "prompts_harden"
},
"TOOL_TIMEOUT": {
"description": "工具执行超时",
"auto_retry": True,
"suggestion": "circuit_breaker_tune"
},
"TOOL_RATE_LIMIT": {
"description": "外部 API 限流",
"auto_retry": True,
"suggestion": "backoff_increase"
},
"TOOL_INPUT_INVALID": {
"description": "工具返回参数校验失败(模型幻觉参数)",
"auto_retry": False,
"suggestion": "schema_strictening"
},
"CONTEXT_OVERFLOW": {
"description": "历史消息超出上下文窗口",
"auto_retry": True,
"suggestion": "context_pruning"
},
"INTENT_DRIFT": {
"description": "模型工具选择与用户意图偏离",
"auto_retry": False,
"suggestion": "supervisor_intervention"
},
"COST_OVERRUN": {
"description": "超出 Token 预算被强制终止",
"auto_retry": False,
"suggestion": "task_decomposition"
}
}
每集一类错误,系统自动生成一份 Postmortem 记录,包含:故障 Trace、根因分类、影响成本、建议改进措施,以及修复后的回归测试钩子。
六、生产落地的 7 条工程经验
- 永远不要把 Tool Schema 校验完全交给模型:模型会幻觉出 schema 中不存在的字段、拼错枚举值、给数字传入字符串。必须在 runtime 层做 JSON Schema 严格校验,非法参数直接返回结构化错误,让模型自己修正——而不是盲目传给工具。
- 重试预算 != 无限重试:给每个工具设置明确的 max_retries(通常 2-3 次)和一个独立的"重试预算池"。当共享预算池耗尽后,连重试也不允许,直接降级或终止。
- 工具调用的失败应该"教导"工具:每次工具调用失败,都自动将失败上下文(调用参数、错误详情、Agent 当时的决策理由)写入错误反馈库。下次同一工具被调用时,这段历史会被注入到 Agent 提示中,避免重复犯错。
- 熔断器阈值需要动态调整:生产中发现,很多工具的非 5xx 错误(如语义层面的失败)不应该触发熔断(工具本身没有故障)。仅 5xx、超时、限流类错误计入熔断计数,语义错误(工具返回"无结果")不计入。
- 多 Agent 监控需要统一的 Trace ID 串联:一个用户请求可能跨越 5 个 Agent 服务、10 次工具调用。缺少跨 Agent 的 trace_id 串联,排障时间会增加 10 倍。
- 预算守护必须是硬性的、不可绕过的:Agent 的"我觉得我应该再试一次"不应被允许——模型的贪婪倾向会使其在接近预算时仍试图继续执行,必须由基础设施层强制干预。
- 降级响应比完全失败更有价值:预算耗尽时返回"已为您检索到以下部分结果,但未完全分析"比直接报错更友好。保留 partial_result 并在前端标注"部分结果",是次优但可用的设计。
总结
AI Agent 的工程化不是堆砌更多工具调用,而是在不确定性中建立确定性保障层。本文描述的架构——熔断退避、预算代理、行为序列追踪、 Supervisor 隔离调度、错误分类自愈——构成了 Agent 从 Demo 走向生产的必要基础设施。
核心原则只有一条:不要信任模型的判断,只信任基础设施的护栏。模型的智能取决于训练数据,基础设施的可靠性则通过工程纪律保障。当两者结合时,Agent 才能在生产中真正替代人工流程。

发表评论 取消回复