智能体原型如何走到可用功能
用几段 Python 脚本把多个 Agent 串起来,能很快验证协作思路;它并不能说明链路已经适合接入真实任务。原型通常把状态放在进程内,默认工具会成功,也默认模型总能返回可解析的字段。接入并发、重启和异常输入后,这些默认条件会逐个失效:任务可能停在半途,重复调用可能没有出口,解析错误也可能一路传到后续步骤。
原型走向可用功能,重点不是继续堆 Prompt,而是补上状态、边界和失败路径。本文围绕状态机、结构化校验和人工接管说明这条链路;代码中的阈值仅用于演示,接入系统前应按任务类型和可接受成本重新设定。
1. 原型崩溃的三大现场:为什么 Demo 进不了生产
把 Agent 原型接入业务前,先检查三个容易被演示掩盖的问题。
首先是状态丢失与不可追溯。许多 Demo 把上下文放在内存列表里,进程退出后就没有续跑依据。至少应保存任务标识、当前步骤、输入摘要、工具结果和版本信息;这样出现偏差时,排查人员能回到同一份状态,而不是靠记忆拼出过程。
其次是重复调用没有出口。工具返回失败后,模型可能沿用同一参数再次尝试。轮次、单次任务预算和相同调用的去重应由程序控制;超过边界后给出可解释的失败状态,交给人工或其他队列处理。
最后是契约过于宽松。模型输出可能带代码块、缺字段或字段类型不对。调用方不应把它直接交给下游,而应先解析、校验并返回具体错误;无法修复时终止当前步骤,比带着错误状态继续执行更安全。
2. 生产级多 Agent 状态机与防线架构
要解决非确定性带来的混乱,核心思路是用确定性的状态机(Finite State Machine)和校验拦截器把 Agent 限制在可控的安全网内。
多 Agent 协作的本质是分布式任务调度。每个 Agent 不再是自由发挥的个体,而是状态机上的节点。节点之间的转换由强类型 Schema 校验器与人工确认门禁共同驱动。
在上面这套架构中,引入了三个硬核治理机制:
- 强类型数据总线:Agent 之间传递的数据必须经过 Pydantic 校验,拒绝任何任意格式的字符串直接流转。
- 状态持久化快照:每次 Agent 完成阶段性思考或工具调用,必须将 Snapshot 写入 Redis 或数据库,确保宕机可恢复。
- 带退避的自我修复(Auto-repair)与人工接管机制:格式错误优先由轻量 Prompt 自纠,纠错失败则抛出事件入队等待人工审核。
3. 生产级 Agent 调度与防线代码实现
下面是一个不依赖高层抽象框架、基于 Python 原生 asyncio、Pydantic 和 Redis 思想实现的生产级多 Agent 状态调度器与验证防线。代码展示了强类型契约、状态持久化与轮次限制的具体落地方式。
import asyncio import json import logging from typing import Dict, Any, Optional, List from pydantic import BaseModel, Field, ValidationError logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s") # 1. 强类型交互契约定义 class AgentOutputSchema(BaseModel): task_id: str = Field(..., description="任务唯一 ID") next_action: str = Field(..., description="下一步动作:SEARCH, EXECUTE, FINISH") payload: Dict[str, Any] = Field(default_factory=dict, description="结构化参数") thought_summary: str = Field(..., description="决策思考摘要") class ExecutionState(BaseModel): task_id: str current_step: int = 0 max_steps: int = 5 accumulated_tokens: int = 0 max_token_budget: int = 4096 status: str = "RUNNING" # RUNNING, FAILED, COMPLETED, NEED_HUMAN history: List[Dict[str, Any]] = [] # 2. 模拟生产环境中的代理执行与防线网关 class ProductionAgentRunner: def __init__(self, state: ExecutionState): self.state = state async def _mock_llm_call(self, prompt: str) -> str: """模拟 LLM 返回,故意包含非标准格式以测试防线""" await asyncio.sleep(0.1) if self.state.current_step == 0: # 正常响应 return json.dumps({ "task_id": self.state.task_id, "next_action": "SEARCH", "payload": {"query": "高并发架构方案"}, "thought_summary": "先检索相关技术文档" }) elif self.state.current_step == 1: # 故意缺少必要字段以测试 Schema 校验 return json.dumps({ "task_id": self.state.task_id, "next_action": "EXECUTE", # 缺少 thought_summary }) else: return json.dumps({ "task_id": self.state.task_id, "next_action": "FINISH", "payload": {"result": "处理完成"}, "thought_summary": "任务收尾" }) async def step(self, user_prompt: str) -> bool: # 防线检查 1:步骤与 Token 预算检查 if self.state.current_step >= self.state.max_steps: logging.warning(f"任务 {self.state.task_id} 达到最大轮次 {self.state.max_steps},触发熔断") self.state.status = "FAILED" return False if self.state.accumulated_tokens >= self.state.max_token_budget: logging.warning(f"任务 {self.state.task_id} 超出 Token 预算,触发中断") self.state.status = "FAILED" return False self.state.current_step += 1 logging.info(f"执行 Step {self.state.current_step}/{self.state.max_steps}...") raw_response = await self._mock_llm_call(user_prompt) # 估算并累加 Token(模拟) self.state.accumulated_tokens += len(raw_response) // 4 # 防线检查 2:强类型 Schema 校验与自修复逻辑 validated_output: Optional[AgentOutputSchema] = None try: raw_json = json.loads(raw_response) validated_output = AgentOutputSchema(**raw_json) except (json.JSONDecodeError, ValidationError) as e: logging.error(f"Schema 校验失败: {e},尝试自动修复...") validated_output = await self._auto_repair_schema(raw_response, str(e)) if not validated_output: logging.error("自动修复失败,转移至人工接管队列") self.state.status = "NEED_HUMAN" return False # 状态持久化快照(模拟保存至 Redis) self.state.history.append(validated_output.model_dump()) logging.info(f"Step {self.state.current_step} 状态校验成功,决策动作: {validated_output.next_action}") if validated_output.next_action == "FINISH": self.state.status = "COMPLETED" return False return True async def _auto_repair_schema(self, bad_output: str, error_msg: str) -> Optional[AgentOutputSchema]: """简化的纠错尝试组件""" logging.info("向修补 Agent 发送修复请求...") await asyncio.sleep(0.05) # 补全缺失的字段 try: data = json.loads(bad_output) data["thought_summary"] = f"自动修补字段,修复原错误: {error_msg[:30]}" return AgentOutputSchema(**data) except Exception: return None # 3. 运行验证 async def main(): state = ExecutionState(task_id="task_prod_9001", max_steps=4) runner = ProductionAgentRunner(state) logging.info(f"开始运行生产级 Agent 任务,初始状态: {state.status}") active = True while active: active = await runner.step("请协助完成系统巡检") logging.info(f"任务最终结束状态: {state.status},总消耗 Token: {state.accumulated_tokens}") logging.info(f"持久化历史记录数: {len(state.history)}") if __name__ == "__main__": asyncio.run(main())运行上述代码,可以在控制台看到清晰的日志,展示当 LLM 输出残缺数据时,系统如何拦截、尝试修补并在无法恢复时优雅降级为NEED_HUMAN状态,而不是直接引发 Unhandled Exception。
2026-08-26 10:15:01 [INFO] 开始运行生产级 Agent 任务,初始状态: RUNNING 2026-08-26 10:15:01 [INFO] 执行 Step 1/4... 2026-08-26 10:15:01 [INFO] Step 1 状态校验成功,决策动作: SEARCH 2026-08-26 10:15:01 [INFO] 执行 Step 2/4... 2026-08-26 10:15:01 [ERROR] Schema 校验失败: 1 validation error for AgentOutputSchema... 2026-08-26 10:15:01 [INFO] 向修补 Agent 发送修复请求... 2026-08-26 10:15:01 [INFO] Step 2 状态校验成功,决策动作: EXECUTE 2026-08-26 10:15:01 [INFO] 执行 Step 3/4... 2026-08-26 10:15:01 [INFO] Step 3 状态校验成功,决策动作: FINISH 2026-08-26 10:15:01 [INFO] 任务最终结束状态: COMPLETED,总消耗 Token: 1284. 落地生产环境的物理验收清单
在准备将多 Agent 系统切入线上真实流量之前,必须对照以下工程清单逐一过一遍:
第一,状态持久化与无状态节点验证。任意杀掉一个正在执行任务的 Agent 容器节点,新的容器恢复后能否读取 Redis 中的状态快照继续执行,而不是重新重头启动任务。
第二,硬性预算隔离测试。在模拟环境注入死循环 Tool Call,验证 Token 计数器与轮次限额能否在到达阈值时 100% 触发中断并产生报警告警,杜绝资金击穿风险。
第三,脏数据抗性测试。向 Agent 节点并发发送非 JSON 文本、空字符串、超长乱码以及 SQL 注入攻击文本,系统应当返回优雅的结构化错误提示,且状态机保持稳定。
把 Demo 变成可用功能,从来不是再加几句 Prompt 的事。把确定性的规矩立起来,让非确定性的 AI 在工程轨道里跑,系统才能真正稳如磐石。