Machine Learning Mastery

The End-to-End Agentic AI Pipeline

8.5内容质量

TL;DR · AI 摘要

生产级代理AI系统依赖七个关键组件构成闭环反馈循环,确保系统稳定性与扩展性。

核心要点

  • 生产系统需包含感知、记忆、推理等7个独立组件,而非简单脚本
  • 防护机制需贯穿所有操作步骤,而非特定阶段
  • 组件解耦设计可应对并发用户和复杂任务场景

结构提纲

按章节快速跳转。

  1. 揭示演示脚本与生产系统的架构差距,强调组件化设计的必要性

  2. 列出感知、记忆、推理等核心组件及其在反馈循环中的角色

  3. 将原始输入转化为结构化数据以供推理处理

  4. 负责任务分解和决策路径生成

  5. 跨组件的安全机制和监控体系设计

思维导图

用一张图看清主题之间的关系。

查看大纲文本(无障碍 / 无 JS 友好)
  • 端到端代理AI流水线
    • 核心组件
      • 感知
      • 记忆
      • 推理/规划
      • 工具执行
      • 编排
      • 防护
      • 可观测性
    • 设计原则
      • 组件解耦
      • 闭环反馈

金句 / Highlights

值得收藏与分享的关键句。

#AI架构#代理AI#系统设计#机器学习
打开原文

端到端代理AI流水线 - MachineLearningMastery.com

端到端代理AI流水线

作者:

Shittu Olumide

发布于:

2026年7月31日

分类:

人工智能

3

分享

文章

在本文中,你将学习区分生产级代理AI系统与演示脚本的七个架构组件,以及每个组件如何融入代理的核心反馈循环。

我们将涵盖的主题包括:

  • 每个组件(感知、记忆、推理与规划、工具执行、编排、防护措施和可观测性)的具体职责。
  • 每个组件在实际系统中容易出错的位置,以及为何必须与其他组件分离。
  • 针对每个组件职责的可运行Python代码示例。

介绍

大多数“构建AI代理”的教程展示的是一个40行的脚本,该脚本在循环中调用LLM并声称任务完成。这个脚本对于演示来说运行良好,但无法应对第二个并发用户、不可靠的第三方API或需要十二个步骤而非两个步骤的任务。

演示与生产系统之间的差距并非巧妙的提示工程,而是架构设计。生产级代理系统由一组一致的互联组件构建:感知、推理、规划、记忆、工具执行、编排和防护措施。这种组件划分在去年几乎所有严肃的架构文章、调查论文和生产后分析中均有体现,无论撰写内容的是哪个框架或供应商。

所有这些背后的循环是一致的:目标 → 感知 → 推理 → 规划 → 行动 → 观察 → 记忆更新 → 返回推理,重复直到目标达成、触发停止条件或代理决定需要人工干预。本文将逐步解析这个循环中的每个环节——每个组件的职责、常见故障点以及聚焦代码片段来具体说明责任划分。这里没有任何内容是硬编码到某个运行流水线中的。每个组件都是独立展示的,这同样是你在决定系统需求时应该采用的思考方式。

七个组件概览

架构调查普遍聚焦于相同的内核组件:感知、记忆、推理/规划、工具执行和编排构成一个封闭的反馈循环——实际运行的循环,逐步执行。防护措施和可观测性作为跨切面关注点环绕整个循环,而非序列中的步骤。你不会在第4步“执行”防护措施;防护措施位于每个拟议行动与现实世界之间,监控每一步。

这一区别塑造了本文的其余内容。前五个部分按照数据实际流动的顺序逐步解析循环。最后两个部分则覆盖了使循环在涉及真实资金、真实客户和真实副作用时仍能运行的封装层。

感知模块的任务是将原始输入(文本、语音、API数据负载、传感器数据和文件上传)转换为推理引擎可以实际处理的结构化表示。这是大多数教程完全跳过的部分,因为在演示中,“用户只需输入文本”且无需进行标准化处理。实际系统需要从Webhook、结构化API调用、文件上传和多个渠道同时接收输入,而所有这些输入在下游任何组件信任之前都必须转换为相同的结构。

perception.py

依赖条件:仅需Python标准库

运行方式:python perception.py

from dataclasses import dataclass, field from typing import Any from enum import Enum import json from datetime import datetime, timezone

class InputSource(Enum): USER_TEXT = "user_text" WEBHOOK = "webhook" FILE_UPLOAD = "file_upload"

@dataclass class AgentInput: """所有下游组件消费的标准化内部结构,无论原始输入实际来自何处。感知层的全部意义就在于此:所有后续组件只能看到这种单一结构。""" source: InputSource content: str metadata: dict[str, Any] = field(default_factory=dict) received_at: str = field(default_factory=lambda: datetime.now(timezone.utc).isoformat())

def perceive_user_text(raw_text: str) -> AgentInput: """原始聊天输入--最简单的情况,但仍需要标准化处理。""" return AgentInput( source=InputSource.USER_TEXT, content=raw_text.strip(), metadata={"channel": "chat"}, )

def perceive_webhook(raw_payload: str) -> AgentInput: """Webhook传递的是结构化JSON而非纯文本。感知模块会提取代理需要推理的部分,并丢弃传输层噪声(如头部和签名)。""" payload = json.loads(raw_payload) event_type = payload.get("event_type", "unknown") description = payload.get("description", "") return AgentInput( source=InputSource.WEBHOOK, content=f"收到事件 '{event_type}': {description}", metadata={"event_type": event_type, "raw_payload": payload}, )

def perceive_file_upload(filename: str, file_size_bytes: int, mime_type: str) -> AgentInput: """文件上传事件完全没有任何自然语言内容--感知模块必须构建推理引擎实际可用的内容。""" return AgentInput( source=InputSource.FILE_UPLOAD, content=f"用户上传文件 '{filename}' ({mime_type}, {file_size_bytes} 字节)", metadata={"filename": filename, "mime_type": mime_type, "size_bytes": file_size_bytes}, )

if __name__ == "__main__": text_input = perceive_user_text(" 请问我的退款状态如何? ") webhook_input = perceive_webhook(json.dumps({ "event_type": "payment_failed", "description": "订单 #4821 付款被拒绝", })) file_input = perceive_file_upload("invoice_q3.pdf", 184320, "application/pdf") for inp in [text_input, webhook_input, file_input]: print(f"[{inp.source.value}] content='{inp.content}'") print(f" metadata keys: {list(inp.metadata.keys())}\n")

InputSource

(

)

:

USER_TEXT

=

"用户文本"

WEBHOOK

"网络钩子"

FILE_UPLOAD

"文件上传"

@

AgentInput

""

"

所有下游组件消耗的标准化内部结构,无论原始输入实际来自何处。感知层的全部意义就在于此:从此之后的所有组件只能看到这种单一结构。

source

content

str

metadata

dict

[

]

default_factory

received_at

lambda

.

now

utc

isoformat

def

perceive_user_text

raw_text

->

"最简单的原始聊天输入,但仍需要标准化。"

return

strip

{

"channel"

"chat"

}

perceive_webhook

raw_payload

网络钩子传递的是结构化JSON而非纯文本。感知层会提取代理需要推理的部分,并丢弃传输层的噪声(如头部和签名)。

payload

loads

event_type

get

"event_type"

"unknown"

description

"description"

f

"收到事件 '{event_type}':{description}"

"raw_payload"

perceive_file_upload

filename

file_size_bytes

int

mime_type

文件上传事件完全没有自然语言内容——感知层必须构建出推理引擎实际可用的结构。

"用户上传文件 '{filename}' ({mime_type}, {file_size_bytes} 字节)"

"filename"

"mime_type"

"size_bytes"

if

__name__

==

"__main__"

text_input

" 我的退款状态如何? "

webhook_input

dumps

"payment_failed"

"订单 #4821 付款失败:信用卡被拒"

file_input

"invoice_q3.pdf"

184320

"application/pdf"

for

inp

print

"[{inp.source.value}] content='{inp.content}'"

" metadata keys: {list(inp.metadata.keys())}\n"

运行方式:python perception.py ,无需依赖项。

三种完全不同的原始结构——纯文本、网络钩子JSON负载和文件上传事件——都会坍缩为相同的AgentInput结构。下游的推理组件永远不需要知道或关心信息是通过哪个渠道到达的。这就是将感知层作为独立组件处理,而非在输入进入系统时随意内联解析的全部价值所在。

工作上下文与实际持久化内容

这是最需要细致处理的组件,演示代码经常错误地将“内存”简单理解为“到目前为止的对话内容”。生产级内存架构将工作内存(当前任务的即时上下文窗口)与长期内存区分开来,长期内存本身又分为情景记忆(发生了什么)、语义记忆(学到的事实)和程序性记忆(技能和操作知识)。短期记忆存在于上下文中且本质上是免费的;长期记忆通常存储在向量数据库中,按语义检索而非精确匹配进行索引。

这种操作区分非常重要:工作内存快速且可丢弃——会话结束的瞬间就会消失。情景记忆为代理提供了工作内存结构上无法实现的东西:跨会话的后见之明,能够回忆“我们以前处理过类似的情况,当时发生了以下事情”。

memory.py # 依赖项:仅需 Python 标准库 # 运行方式:python memory.py from dataclasses import dataclass, field from datetime import datetime, timezone @dataclass class WorkingMemory: """ 工作记忆:仅用于当前任务的即时上下文。在进程内运行,受对话轮次限制,会话结束即消失。这便是大多数演示代码所称的 "memory" —— 但只是完整图景中的一环 """ max_turns: int = 10 turns: list[dict] = field(default_factory=list) def add_turn(self, role: str, content: str) -> None: self.turns.append({"role": role, "content": content}) if len(self.turns) > self.max_turns: self.turns.pop(0) # 轮次超过限制时移除最早记录 def as_context(self) -> str: return "\n".join(f"{t['role']}: {t['content']}" for t in self.turns) @dataclass class EpisodicMemoryEntry: """单个存储的事件记录 —— 发生了什么、发生时间及用于后续检索的向量表示""" timestamp: str summary: str embedding: list[float] # 生产环境会使用真实向量模型生成 class EpisodicMemory: """ 情景记忆:跨会话持久化存储,存于外部存储(生产环境使用向量数据库),通过语义相似度而非时间顺序进行检索。这使代理具备"事后认知"能力 —— 工作记忆结构无法实现此功能,因为会话结束即消失 """ def __init__(self): self._store: list[EpisodicMemoryEntry] = [] def record_episode(self, summary: str, embedding: list[float]) -> None: self._store.append(EpisodicMemoryEntry( timestamp=datetime.now(timezone.utc).isoformat(), summary=summary, embedding=embedding, )) def retrieve_similar(self, query_embedding: list[float], top_k: int = 2) -> list[EpisodicMemoryEntry]: """实际实现会使用向量索引进行余弦相似度计算""" def dot(a, b): return sum(x * y for x, y in zip(a, b)) ranked = sorted(self._store, key=lambda e: dot(e.embedding, query_embedding), reverse=True) return ranked[:top_k] if __name__ == "__main__": wm = WorkingMemory(max_turns=3) wm.add_turn("user", "What's my refund status?") wm.add_turn("agent", "Let me check that for you.") wm.add_turn("user", "It's order 4821") wm.add_turn("agent", "Found it -- refund is processing") # 超出限制时移除最早记录 print("Working memory (bounded to last 3 turns):") print(wm.as_context()) em = EpisodicMemory() em.record_episode("User asked about refund for order 4821, resolved successfully", [0.9, 0.1, 0.0]) em.record_episode("User asked about shipping delay for order 1190", [0.1, 0.9, 0.0]) em.record_episode("User asked about refund eligibility for order 7734", [0.85, 0.15, 0.0]) similar = em.retrieve_similar([0.88, 0.12, 0.0], top_k=2) print("\nEpisodic memory -- retrieved by similarity to a new refund query:") for entry in similar: print(f" {entry.summary}")

75

76

77

78

memory.py

运行方式:python memory.py

WorkingMemory

工作记忆:仅用于当前任务的即时上下文。

在进程内运行,受对话轮次限制,一旦

会话结束即消失。这便是大多数演示代码所称的 "

memory

" —

但只是完整图景中的一环。

max_turns

turns

list

add_turn

self

role

None

append

"role"

"content"

len

pop

0

轮次超过限制时移除最早记录

as_context

"\n"

join

"{t['role']}: {t['content']}"

t

EpisodicMemoryEntry

"单个存储的事件记录 —— 发生了什么、发生时间及用于后续检索的向量表示。"

timestamp

summary

embedding

float

/

在生产环境中,这来自真正的嵌入模型

EpisodicMemory

情景记忆:跨会话持久化,外部存储

(生产环境中为向量存储),通过语义相似度而非近期性进行检索。这正是代理具有"

事后认知

"能力的原因——这是工作记忆结构上无法实现的,因为一旦会话结束,信息就会消失。

__init__

_store

record_episode

retrieve_similar

query_embedding

top_k

"真实实现会对向量索引进行余弦相似度计算。"

dot

a

b

sum

x *

y

x

zip

ranked

sorted

key

e

reverse

True

wm

"user"

"我的退款状态如何?"

"agent"

"让我为您查询。"

"是订单4821"

"找到了——退款正在处理中"

将第一轮对话推出

"工作记忆(限制为最近3轮对话):"

em

"用户询问了订单4821的退款问题,已成功解决"

0.9

0.1

0.0

"用户询问了订单1190的发货延迟问题"

"用户询问了订单7734的退款资格问题"

0.85

0.15

similar

0.88

0.12

"\n情景记忆——通过与新退款查询的相似度检索:"

entry

" {entry.summary}"

运行方式:python memory.py ,无需依赖项。

工作记忆在达到限制时会丢弃最旧的对话记录,到会话结束时关于查询状态的第一轮对话已消失。情景记忆则相反:在三个存储条目中,它会根据语义而非发生时间,优先展示两个与退款相关的条目。这就是两者的结构性区别——一个是滑动窗口,另一个是可搜索的档案库。

推理与规划(决定下一步该做什么)

推理与规划模块会接收当前目标、感知输入以及检索到的记忆,生成计划——有时是单一的下一步动作,有时是多步骤的分解。这是代理的认知核心,通过咨询记忆和知识资源,合成可传递给执行模块的动作方案。

关键设计点(容易被忽视):规划模块的责任仅限于生成计划。它不会调用工具、接触API,也不会产生任何副作用。这种分离是刻意为之,正是这种设计使得下一个组件——工具执行——可以独立测试和独立防护。

planning.py # 先决条件:仅需 Python 标准库 # 运行方式:python planning.py 从 dataclasses import dataclass, field import json @dataclass class PlanStep: step_number: int description: str tool_tag: str # 该步骤需要的工具类别 -- 执行决定如何操作 @dataclass class Plan: goal: str steps: list[PlanStep] = field(default_factory=list) def mock_llm_plan(goal: str) -> str: """ 模拟 LLM 调用,替代真实的规划调用。要演示的重点是结构性的:规划会生成一个计划对象,但不会执行任何操作。执行是另一个独立组件(下一节),计划会传递给该组件处理。 """ if "refund" in goal.lower(): return json.dumps({ "steps": [ {"step_number": 1, "description": "通过ID查找订单", "tool_tag": "database_lookup"}, {"step_number": 2, "description": "根据政策检查退款资格", "tool_tag": "policy_check"}, {"step_number": 3, "description": "如果符合条件则发放退款", "tool_tag": "payment_api"}, {"step_number": 4, "description": "通知客户结果", "tool_tag": "email"}, ] }) return json.dumps({"steps": [ {"step_number": 1, "description": "在知识库中搜索答案", "tool_tag": "search"}, ]}) def create_plan(goal: str) -> Plan: """接收目标,生成结构化计划。此处没有任何副作用。""" parsed = json.loads(mock_llm_plan(goal)) return Plan(goal=goal, steps=[PlanStep(**s) for s in parsed["steps"]]) if __name__ == "__main__": for goal in ["Process a refund for order 4821", "What are your business hours?"]: plan = create_plan(goal) print(f"Goal: {plan.goal}") for step in plan.steps: print(f" Step {step.step_number}: {step.description} [tool_tag={step.tool_tag}]") print()

planning.py

运行方式:python planning.py

PlanStep

step_number

tool_tag

该步骤需要的工具类别 -- 执行决定如何操作

Plan

goal

steps

mock_llm_plan

模拟 LLM 调用,替代真实的规划调用。要演示的重点是结构性的:规划会生成一个计划对象 -- 它不会执行任何操作。执行是另一个独立组件(下一节),该计划会传递给该组件处理。

"refund"

lower

"steps"

"step_number"

"通过ID查找订单"

"tool_tag"

"database_lookup"

"根据政策检查退款资格"

"policy_check"

"如果符合条件则发放退款"

"payment_api"

"通知客户结果"

"email"

"在知识库中搜索答案"

"search"

create_plan

"接收目标,生成结构化计划。此处没有任何副作用。"

parsed

*

s

"Process a refund for order 4821"

"What are your business hours?"

"Goal: {plan.goal}"

step

" Step {step.step_number}: {step.description} [tool_tag={step.tool_tag}]"

运行方式:python planning.py,无需依赖项。

退款目标会生成一个四步计划;营业时间问题会生成一个步骤。两者都没有执行任何工具——它们只是返回了一个描述下一步应发生事项的Plan对象。该对象是推理与其他流水线组件之间的交接凭证,这正是为什么后续的编排(稍后会讲到)可以在任何步骤实际运行前选择暂停、修改或拒绝计划。

工具执行将代理连接到外部系统——API、数据库和服务——处理调用能力的机制并将结果反馈到推理过程中。这也是大多数生产环境事故的真正源头,因为它是整个循环中唯一具有真实外部副作用的组件。

用具体数字说明这一限制很有必要:如果每个操作的失败率是5%,那么一个在运行中执行20个操作的代理在没有防护措施的情况下会频繁失败到无法使用。正是这个统计数据解释了为什么工具执行不能只是“调用API并祈祷”——它必须具备验证、超时和幂等性作为基本要求,而不是可有可无的附加功能。

tool_execution.py

依赖条件:仅需 Python 标准库

运行方式:python tool_execution.py

import time import hashlib from dataclasses import dataclass from typing import Callable, Any, Optional

@dataclass class ToolResult: success: bool output: Any = None error: Optional[str] = None idempotency_key: Optional[str] = None from_cache: bool = False

class ToolExecutor: """ 工具执行的不可妥协的基本原则:

  1. 在调用任何外部操作前验证输入
  2. 强制超时机制,防止慢速调用阻塞整个运行流程
  3. 使具有副作用的操作具备幂等性,确保重试不会导致重复扣款或重复发送邮件

"""

def __init__(self, timeout_seconds: float = 5.0): self.timeout_seconds = timeout_seconds self._idempotency_cache: dict[str, ToolResult] = {}

def _make_idempotency_key(self, tool_name: str, args: dict) -> str: raw = f"{tool_name}:{sorted(args.items())}" return hashlib.sha256(raw.encode()).hexdigest()[:16]

def execute(self, tool_name: str, tool_fn: Callable, args: dict, required_args: list[str], idempotent: bool = False) -> ToolResult:

1. 在接触任何外部操作前验证参数

missing = [a for a in required_args if a not in args] if missing: return ToolResult(success=False, error=f"缺少必要参数: {missing}")

2. 幂等性 - 相同参数的重试将返回缓存结果

idem_key = self._make_idempotency_key(tool_name, args) if idempotent else None if idem_key and idem_key in self._idempotency_cache: cached = self._idempotency_cache[idem_key] return ToolResult( success=cached.success, output=cached.output, idempotency_key=idem_key, from_cache=True )

3. 在超时预算内执行

start = time.monotonic() try: output = tool_fn(**args) except Exception as e: return ToolResult(success=False, error=str(e), idempotency_key=idem_key)

if (time.monotonic() - start) > self.timeout_seconds: return ToolResult(success=False, error=f"工具执行超时({self.timeout_seconds}s)", idempotency_key=idem_key)

result = ToolResult(success=True, output=output, idempotency_key=idem_key) if idem_key: self._idempotency_cache[idem_key] = result return result

def issue_refund(order_id: str, amount: float) -> str: return f"已退款 ${amount}(订单号:{order_id})"

if __name__ == "__main__": executor = ToolExecutor(timeout_seconds=1.0)

缺少参数测试 - 在执行前被拦截

r1 = executor.execute("issue_refund", issue_refund, {"order_id": "4821"}, ["order_id", "amount"]) print(f"缺少参数测试: success={r1.success}, error='{r1.error}'")

第一次调用执行,相同参数重试将命中幂等性缓存

args = {"order_id": "4821", "amount": 49.99} r2a = executor.execute("issue_refund", issue_refund, args, ["order_id", "amount"], idempotent=True) r2b = executor.execute("issue_refund", issue_refund, args, ["order_id", "amount"], idempotent=True)

print(f"\n首次调用: from_cache={r2a.from_cache}, output='{r2a.output}'") print(f"重试(相同参数): from_cache={r2b.from_cache}, output='{r2b.output}'")

_make_idempotency_key

tool_name

args

raw

"{tool_name}:{sorted(args.items())}"

sha256

encode

hexdigest

execute

tool_fn

required_args

idempotent

1. 在接触任何外部系统之前进行验证

missing

not

"Missing required args: {missing}"

2. 幂等性 -- 使用相同参数重试调用时返回缓存结果

idem_key

else

and

cached

3. 在超时预算内执行

start

monotonic

try

except

Exception

as

-

"Tool exceeded {self.timeout_seconds}s timeout"

result

issue_refund

order_id

amount

"Refunded ${amount} for order {order_id}"

executor

1.0

缺少参数 -- 在执行前被捕获

r1

"issue_refund"

"order_id"

"4821"

"amount"

"Missing arg test: success={r1.success}, error='{r1.error}'"

首次调用执行,使用相同参数重试时命中幂等性缓存

49.99

r2a

r2b

"\nFirst call: from_cache={r2a.from_cache}, output='{r2a.output}'"

"Retry, same args: from_cache={r2b.from_cache}, output='{r2b.output}'"

运行方式:python tool_execution.py ,无需依赖项。

使用相同参数重试时会返回缓存结果,而不是再次调用issue_refund。即使上层协调器因短暂网络问题重试步骤,客户也只会收到一次退款。这就是为什么要在执行层实现幂等性,而不是依赖协调器永不重试。

协调

协调层在多个步骤之间(在多智能体系统中跨多个智能体)保持循环,决定何时继续、何时根据步骤结果改变路径、以及何时实际完成运行。这是近年来发展最快的层级,LangGraph、CrewAI和AutoGen现在已能处理生产级协调,而不是每个团队都从零开始手动实现循环。

orchestrator.py # 前提条件:仅需Python标准库 # 运行方式:python orchestrator.py from dataclasses import dataclass, field @dataclass class PlanStep: step_number: int description: str tool_tag: str @dataclass class Plan: goal: str steps: list[PlanStep] = field(default_factory=list) @dataclass class StepOutcome: step_number: int success: bool output: str # 用于替代上文中的规划和工具执行组件 -- # 协调器的任务是按顺序调用这些组件,而不是执行它们的实际工作 def mock_create_plan(goal: str) -> Plan: return Plan(goal=goal, steps=[ PlanStep(1, "查询订单", "database_lookup"), PlanStep(2, "检查资格", "policy_check"), PlanStep(3, "发放退款", "payment_api"), ]) def mock_execute_tool(step: PlanStep) -> StepOutcome: if step.tool_tag == "policy_check": # 模拟失败以测试停止逻辑 return StepOutcome(step.step_number, success=False, output="资格检查失败:订单过期") return StepOutcome(step.step_number, success=True, output=f"完成:{step.description}") class Orchestrator: """ 在多个步骤中保持流程连贯性,决定是否继续或停止,并通过明确的步骤限制来控制执行。它不进行规划,也不执行工具操作 -- 它只是调用实际执行这些操作的组件。 """ def __init__(self, max_steps: int = 10): self.max_steps = max_steps def run(self, goal: str) -> list[StepOutcome]: plan = mock_create_plan(goal) outcomes: list[StepOutcome] = [] for step in plan.steps: if len(outcomes) >= self.max_steps: print(f" 步骤限制({self.max_steps})已达到 -- 停止执行。") break print(f" 执行步骤 {step.step_number}: {step.description}") outcome = mock_execute_tool(step) outcomes.append(outcome) if not outcome.success: print(f" 步骤 {step.step_number} 失败: {outcome.output}") print(f" 停止执行 -- 出现失败步骤将阻断后续计划。") break return outcomes if __name__ == "__main__": outcomes = Orchestrator(max_steps=10).run("处理订单4821的退款") print(f"\n尝试执行的步骤总数: {len(outcomes)}") print(f"最终结果: success={outcomes[-1].success}, output='{outcomes[-1].output}'")

orchestrator.py

运行方式:python orchestrator.py

StepOutcome

用于替代上文中的规划和工具执行组件 --

协调器的任务是按顺序调用这些组件,而不是执行它们的实际工作。

mock_create_plan

"查询订单"

"检查资格"

"发放退款"

mock_execute_tool

模拟失败以测试停止逻辑

"资格检查失败:订单过期"

"完成:{step.description}"

Orchestrator

在多个步骤中保持流程连贯性,决定是否继续或停止,并通过明确的步骤限制来控制执行。它不进行规划,也不执行工具操作 -- 它只是调用实际执行这些操作的组件。

max_steps

run

outcomes

=

" 步骤限制({self.max_steps})已达到 -- 停止执行。"

break

" 执行步骤 {step.step_number}: {step.description}"

outcome

" 步骤 {step.step_number} 失败: {outcome.output}"

" 停止执行 -- 出现失败步骤将阻断后续计划。"

"\n尝试执行的步骤总数: {len(outcomes)}"

"最终结果: success={outcomes[-1].success}, output='{outcomes[-1].output}'"

如何运行 : python orchestrator.py ,无需依赖项。

输出 :

执行步骤 1:查找订单 执行步骤 2:检查资格 步骤 2 失败:策略检查失败:订单过期 停止执行——失败的步骤会阻断当前计划的后续步骤 已尝试步骤总数:2 最终结果:success=False,output='策略检查失败:订单过期'

执行 查找 订单 检查 资格 失败 策略 过期 停止 —— 阻断 后续 步骤 总数 尝试 最终 '策略检查失败:订单过期'

步骤 3——实际退款过程——从未执行。这不是模拟环境的偶然现象,而是协调器在履行其特定职责。计划器在生成三步计划时并不知道步骤 2 是否会成功。工具执行器运行了步骤 2 并报告失败。决定在此停止而非盲目继续对刚刚失去资格的订单发起退款,这既不属于计划器也不属于执行器的职责范畴,而是属于协调机制的职责范围。

安全防护机制

安全防护机制定义了行为规则:工具和域的允许/禁止列表、隐私与数据驻留控制、成本上限、速率限制,以及针对高风险或不可逆操作的升级路径。这不是产品上线后的附加功能,而是区分演示中表现惊艳的代理与能够安全对接真实客户账户和真实支付系统的代理的关键所在。

当前生产环境的最佳实践已收敛于相同的核心模式:通过代码定义策略、对不可逆操作实施强制审批关卡,以及防范提示注入攻击(防止不可信的检索内容被误认为指令)。

guardrails.py # 依赖条件:仅需 Python 标准库 # 运行方式:python guardrails.py from dataclasses import dataclass from enum import Enum class GuardrailVerdict(Enum): ALLOW = "允许" DENY = "拒绝" REQUIRE_APPROVAL = "需人工审批" @dataclass class ProposedAction: tool_name: str args: dict estimated_cost: float irreversible: bool @dataclass class GuardrailResult: verdict: GuardrailVerdict reason: str class GuardrailEngine: """ 位于规划模块的拟执行动作与工具执行器之间 实现工具白名单控制、单次操作成本上限控制, 对所有不可逆操作强制要求人工审批(无论规划模块的置信度多高) """ def __init__(self, allowed_tools: set[str], cost_ceiling: float): self.allowed_tools = allowed_tools self.cost_ceiling = cost_ceiling def check(self, action: ProposedAction) -> GuardrailResult: if action.tool_name not in self.allowed_tools: return GuardrailResult(GuardrailVerdict.DENY, f"工具 '{action.tool_name}' 不在白名单中") if action.estimated_cost > self.cost_ceiling: return GuardrailResult(GuardrailVerdict.DENY, f"预估成本 ${action.estimated_cost} 超过上限 ${self.cost_ceiling}") if action.irreversible: return GuardrailResult(GuardrailVerdict.REQUIRE_APPROVAL, "不可逆操作 - 需人工审批后方可执行") return GuardrailResult(GuardrailVerdict.ALLOW, "通过所有安全校验") if __name__ == "__main__": guardrails = GuardrailEngine( allowed_tools={"database_lookup", "policy_check", "payment_api", "email"}, cost_ceiling=100.0, ) test_actions = [ ProposedAction("database_lookup", {"order_id": "4821"}, 0.0, False), ProposedAction("delete_customer_account", {"id": "u_991"}, 0.0, True), ProposedAction("payment_api", {"order_id": "4821", "amount": 5000.0}, 5000.0, True), ProposedAction("payment_api", {"order_id": "4821", "amount": 49.99}, 49.99, True), ] for action in test_actions: result = guardrails.check(action) print(f"{action.tool_name} -> {result.verdict.value}: {result.reason}")

guardrails.py

运行方式:python guardrails.py

GuardrailVerdict

允许

"允许"

拒绝

"拒绝"

需人工审批

"需人工审批"

ProposedAction

estimated_cost

irreversible

GuardrailResult

verdict

reason

GuardrailEngine

位于规划模块的拟执行动作与工具执行器之间。

实现工具白名单控制、单次操作成本上限控制,

对所有不可逆操作强制要求人工审批 -- 无论

规划模块的置信度多高。

allowed_tools

set

cost_ceiling

action

"工具 '{action.tool_name}' 不在白名单中"

"预估成本 ${action.estimated_cost} 超过上限 ${self.cost_ceiling}"

"不可逆操作 -- 需人工审批后方可执行"

"通过所有安全校验"

guardrails

100.0

test_actions

"delete_customer_account"

"id"

"u_991"

5000.0

"{action.tool_name} -> {result.verdict.value}: {result.reason}"

运行方式:python guardrails.py ,无需依赖项。

database_lookup -> 允许: 通过所有安全校验 delete_customer_account -> 拒绝: 工具 'delete_customer_account' 不在白名单中 payment_api -> 拒绝: 预估成本 $5000.0 超过上限 $100.0 payment_api -> 需人工审批: 不可逆操作 -- 需人工审批后方可执行

database_lookup

通过

所有

安全校验

delete_customer_account

工具

'delete_customer_account'

不在

payment_api

预估

费用

$

超过

上限

需要

人工

批准

在执行之前

最后一个案例值得我们关注:一笔 $49.99 的退款,远低于 $100 的费用上限,仍然需要人工批准,因为它是不可逆的——完全无法例外。即使在预算范围内,也无法推翻这一规则。这正是那种看似简单易写的规则,但如果未将防护机制视为独立组件并进行专门检查,就很容易被忽略。

可观测性

可观测性意味着对每个步骤进行追踪级别的日志记录——不仅仅是最终输出,还包括工具选择和中间推理过程——因为没有追踪信息就无法调试或改进代理行为。这个组件能将「代理做了错误的事」转化为「代理在第 4 步调用了 policy_check,返回了 success=False,而协调器正确地在此处停止了运行」。这是可迭代改进的系统与只能重启并寄希望于运气的系统之间的关键区别。

observability.py # 先决条件:仅需 Python 标准库 # 运行方式:python observability.py from dataclasses import dataclass, field from datetime import datetime, timezone from typing import Any import json @dataclass class TraceEntry: step_number: int component: str event: str detail: dict[str, Any] timestamp: str = field(default_factory=lambda: datetime.now(timezone.utc).isoformat()) class TraceLogger: """ 为运行的每一步骤添加结构化、带时间戳的日志条目。刻意与协调器分离:协调器决定发生什么,追踪日志器只是观察并记录——这正是你能直接定位导致偏差的具体步骤、组件和输入的原因,而无需重新运行任何内容。 """ def __init__(self, run_id: str): self.run_id = run_id self.entries: list[TraceEntry] = [] self._step_counter = 0 def log(self, component: str, event: str, **detail) -> None: self._step_counter += 1 self.entries.append(TraceEntry(self._step_counter, component, event, detail)) def export_json(self) -> str: return json.dumps({ "run_id": self.run_id, "trace": [ {"step": e.step_number, "component": e.component, "event": e.event, "detail": e.detail, "timestamp": e.timestamp} for e in self.entries ], }, indent=2) def find_failure_point(self) -> TraceEntry | None: """追踪日志中最有用的查询:第一次出错的位置在哪里。""" for entry in self.entries: if entry.detail.get("success") is False: return entry return None if __name__ == "__main__": tracer = TraceLogger(run_id="run_2026_06_20_001") tracer.log("perception", "input_normalized", source="user_text", content="Refund order 4821") tracer.log("planning", "plan_created", step_count=3, goal="Process a refund for order 4821") tracer.log("tool_execution", "tool_called", tool="database_lookup", success=True, output="Order found") tracer.log("tool_execution", "tool_called", tool="policy_check", success=False, output="Order too old", error="Policy check failed: order too old") tracer.log("orchestrator", "run_stopped", reason="step_failed", step_number=4) failure = tracer.find_failure_point() print(f"Failure point: component='{failure.component}', step={failure.step_number}") print(f" detail: {failure.detail}")

observability.py

运行方式:python observability.py

TraceEntry

组件

事件

细节

TraceLogger

为运行的每一步骤添加结构化、带时间戳的日志条目。

与编排器刻意分离:编排器决定发生什么,追踪记录器只是观察并记录过程——这正是让你能直接定位导致偏差的具体步骤、组件和输入,而无需重新运行任何内容。

run_id entries _step_counter log += export_json

"run_id" "trace" "step" "component" "event" "detail" "timestamp" indent find_failure_point

| "追踪中最实用的查询:首次出错的位置在哪里。" "success" tracer

"run_2026_06_20_001" "perception" "input_normalized" "退款订单 4821" "planning" "plan_created" step_count "tool_execution" "tool_called" "订单找到" "订单过期" "orchestrator" "run_stopped" "step_failed" failure "故障点:组件='{failure.component}',步骤={failure.step_number}" "详情:{failure.detail}"

运行方式:python observability.py,无需依赖项。

find_failure_point() 会遍历追踪记录,直接定位到第 4 步的 policy_check 调用——精确到具体工具、具体参数和具体失败原因。无需重新运行代理,无需猜测五个步骤中哪一个出了问题。这就是将可观测性作为结构性组件封装在循环之外,而非在编排器中随意添加 print() 语句的实践价值。

总结

一旦系统脱离演示阶段,这些组件都将成为必需,尽管大多数教程只会展示其中两三个。感知和记忆为推理核心提供输入。推理和规划将计划对象传递给工具执行。编排负责在各步骤间保持流程连贯,并决定何时停止。护栏和可观测性应封装整个循环而非嵌入其中——前者限制循环可执行的操作,后者记录循环实际执行的内容。

生产级智能体系统最终看起来更像软件架构而非提示工程,是因为在大语言模型调用之下,它们确实是软件架构。模型只是七个处理推理和规划的组件之一,而非整个系统。理解每个组件的具体独立职责,才能实现故障调试、风险操作防护以及流水线的规模化扩展,而不仅仅是希望那40行脚本持续运行。

更多相关内容

  • 构建端到端情感分析流水线
  • 构建长期运行智能体的上下文裁剪流水线
  • 如何结合 LLM 嵌入 + TF-IDF + 元数据
  • 使用多 GPU 训练大型模型
  • 5 个提升工作效率的 Scikit-learn 流水线技巧
  • 使用 llama.cpp 在 Python 中构建 RAG 流水线

/.entry