资讯详情

Agent间数据流与控制流分离:可复用协作架构的落地实践

📅 2026/9/10 8:44:02 | 华诺云谱 👁 阅读
Agent间数据流与控制流分离:可复用协作架构的落地实践
新年回来第一天就遇到一个典型问题项目里三个 Agent 协作跑了几周字段越加越多路由判断越写越长到最后谁都不敢动其中一个 Agent——改一处后面全乱。我把代码翻开一看发现每个 Agent 内部都堆满了这一步跑完该把结果发给谁下游 Agent 想要什么格式如果报错就把数据原样回调给上游这类逻辑。数据在业务代码里传着传着控制逻辑也跟着长在了数据上。这篇文章就是冲着这个问题来的。标题给出的方向是Agent 间数据流与控制流分离我的目标是把它拆成能直接落地的东西为什么分离能解决多 Agent 协作的混乱、分离的具体边界画在哪、怎么用一套简单的消息协议和编排器实现以及分离之后还会踩哪些坑。适合正在从零搭多 Agent 项目、或者已经被 Agent 编排代码折磨过的开发者和架构师。内容不偏理论给的是我自己跑过、验证过的方案和经验。1. 一次多 Agent 翻车现场控制逻辑和数据在代码里打成死结1.1 我最初的写法一个字典传到底刚开始搭协作系统时我按最朴素的方式想Agent A 产出的结果总要传给 Agent B那就定义一个万能结果字典里面放上业务数据、错误标记、下一步要发给谁。后来代码演变成了这样def agent_a_process(task): result do_search(task[query]) # 这里就开始帮下游做决定了 if result[score] 0.8: task[next] write_agent task[data] result else: task[error] retry task[next] search_agent return task看起来没什么问题对吧问题藏在两点上第一Agent A 开始充当调度器它必须知道整个流程里还有谁、谁适合处理低置信度的结果第二业务数据和控制信息共用同一个结构传递过程中任何一方格式变化另一方就看不懂了。这还只是两个 Agent 的场景。等第四个、第五个 Agent 加入A 已经不认识全部流程了于是又有人加了一层统一消息转发每个 Agent 返回后先走到转发器由转发器解析 result 里那个 next 字段再决定发给谁。表面上解耦了实际上转发器里全是针对具体字段的 if-else和把流程塞进 Agent 没有本质区别只是集中到了更难维护的地方。1.2 控制与数据纠缠的三种典型症状我在这个项目里观察到的症状很有代表性基本可以归纳成三类症状表现最终后果隐式路由Agent 返回值里带 next、target 这类字段由调用方解析新增流程时旧 Agent 必须知道新流程改动范围失控格式耦合下游直接依赖上游返回的字段名和层级换掉上游 Agent 时下游无条件崩溃Agent 无法替换职责污染Agent 内部处理超时、重试、优先级等编排逻辑Agent 无法单独测试业务逻辑和系统治理逻辑搅在一起最让人头疼的是职责污染。想想看Agent 的核心职责是做某件具体的事——检索、写作、审核——但一旦它开始关心失败后要不要重试重试几次消息优先级多高它就不再是一个可复用的单元而是一段绑死在特定流程里的脚本。这也是为什么标题里强调可复用。真正可复用的 Agent 应该像一台只认识原料和成品的机器给它什么协议的数据它处理完按协议返回。至于这台机器在流水线上放在哪个工位、上游是谁、下游是谁、出了故障是重启还是报废都不应该是机器自己关心的事。2. 快递包裹和分拣系统先给数据流与控制流画清边界2.1 不再绕圈子的精确定义这个类比我用了很久对团队里新人也很有效数据流是快递包裹本身控制流是物流分拣系统。包裹从 A 地到 B 地里面装什么由寄件人决定包裹运到哪个分拣中心、下一站发往哪里、运输延误多少天算超时、签收失败后退回还是重新派送这些由物流系统决定。放到 Agent 系统里数据流指的是 Agent 之间传递的真实业务信息用户输入、检索结果、中间产物、生成文本、参数配置。控制流指的是系统内的编排决策下一个执行哪个 Agent、条件满足时走哪条分支、并行任务什么时候汇聚、失败重试怎么处理、超时阈值是多少。很多人其实分得清这两个词的字面意思但落到代码里立刻混在一起。原因很简单——为了省事。既然调 Agent A 的代码马上要调 Agent B顺手把接下来调谁写成 A 的返回值不是最方便吗确实方便但短期方便换来的是长期脆弱。2.2 控制流必须知道数据流数据流绝不能反向依赖控制流这是我在实践中总结出的最重要一条原则。控制流要不要感知数据流要。编排器决定下一步发给谁之前需要知道当前任务的业务结果——比如写稿 Agent 完成后要根据文章是否通过审核决定进入结束节点还是退回修改节点这个审核是否通过就是数据流里的信息。所以控制层读数据、判断数据、决定路由是完全必须的。但数据流绝不能反向依赖控制流。Agent 处理业务时不应该去查下一步是谁我这次属于第几次重试之后还有几个 Agent 要接力这些信息跟它的业务逻辑无关。一旦 Agent 业务逻辑里出现对控制信息的访问就破坏了分层。这有点类似操作系统里的用户态和内核态也像 HTTP 请求里的请求体和请求头的关系帮你传输数据的人负责路由和可靠性你写的业务数据只表达业务本身。设计 Agent 间的消息结构时这个分层要直接体现在数据结构上而不是停留在口头约定上。2.3 主流框架隐藏的同一思想现在市面上的主流 Agent 编排框架比如 LangGraph、AutoGen、CrewAI其实内嵌的编排与执行机制就已经把这个思想落地了——框架里的流程图、状态机、Harness 就是控制流的载体而节点间传递的消息内容是数据流。框架使用者如果心里没有这根弦很容易写出违背框架设计意图的代码比如在 Agent 内部直接修改全局编排状态或者在消息正文里夹带路由标记。所以这篇文章虽然给的是手写轻量框架的实现思路但它和主流框架的底层模型是一致的。理解了这个分离原则再去读 LangGraph 这类框架的源码原来觉得晦涩的节点边状态快照概念会清晰很多它们都是在帮你把数据流搬上车、把控制流做成轨道。3. 落地分离架构信封消息、消息总线与编排器的分工3.1 消息信封往数据外面套一层壳分离的第一步不是去画架构图而是重新设计消息结构。我采用的做法是给所有 Agent 之间的通信定义统一的信封业务数据装在信封里控制信息写在信封表面。信封包含两部分一部分是给编排系统看的元信息另一部分是给业务 Agent 入参出参的载荷。一个比较顺手的信封字段设计如下msg_id消息唯一 ID用于幂等、重试和去重trace_id整条业务链路的追踪 ID所有 Agent 日志里都打印它route本消息的目标 Agent 标识由编排器写入msg_type消息类型常见的有task、result、error、controlcontrol控制元信息例如重试次数、超时时间、优先级payload真正的业务数据关键点在于control和payload分开。控制层只读control和route不解析payload的具体业务结构Agent 只处理payload不关心route和control。两边通过信封这条约定各取所需互不侵入。3.2 控制总线与数据总线分开承载有了信封还需要两个承载通道可以称为控制总线与数据总线。控制总线专门传递编排决策消息——任务分配、节点完成通知、重试指令、终止信号数据总线承载真正的业务内容通常是 Agent 产出的中间结果。为什么要拆两条不是矫情而是因为两类消息的生命周期和传输要求完全不同。控制消息要求低延迟、小体积、高可靠投递丢失一条可能导致整个流程卡死数据消息可能体积大、需要持久化、可以异步消费。把它们混在同一条通道里处理策略很难同时满足两边的需求——这在引入消息队列时尤其明显。我实际项目的做法是控制总线用内存中的同步调用或轻量任务队列数据总线用共享存储加引用传递。也就是说数据不直接塞在信封里长途运输而是写入一个可被全局访问的存储后信封里只带一个data_keyAgent 按需去取。3.3 业务上下文与编排状态必须分家第三个容易犯错的地方是状态管理。多 Agent 协作中A 的产出要能被 B 访问所以天然需要一份全局业务上下文但同时编排器需要记录流程走到哪一步、完成了哪些节点、哪些节点重试了几次。这两类状态的变更频率和访问方完全不同。业务上下文跟着数据流走被 Agent 读写编排状态跟着控制流走只被编排器读写。如果把它们合在一个大全局对象里Agent 为了取业务上下文不得不拿到整个编排状态的写权限危险程度可想而知——一个 Agent 手滑改了current_step整个流程状态就崩了。我采用的方式是用两个独立的上下文对象一个存业务上下文一个存编排状态两者通过trace_id关联。Agent 注入业务上下文编排器操作编排状态互不交叉。业务上下文(DataContext): 以 trace_id 为 key存业务中间结果Agent 可读可写 编排状态(ControlState): 以 trace_id 为 key存当前节点、已完成列表、重试计数仅编排器可写这个拆分看起来只是代码结构上的小调整但它直接决定了 Agent 能不能被安全地并行替换和复用。4. 极简可跑的参考实现用 200 行代码搭一个分离式协作框架理论说得再多不如给一段能跑的实现。我用 Python 写了一个最基本的分离式协作骨架代码量不大但把信封、编排器、Agent 注册表、数据上下文几个部分都包含进去了。你可以直接抄回去改。4.1 核心数据结构首先是信封类from dataclasses import dataclass, field from typing import Any, Dict, Optional from uuid import uuid4 dataclass class Envelope: msg_id: str field(default_factorylambda: str(uuid4())) trace_id: str field(default_factorylambda: str(uuid4())) route: str unknown msg_type: str task # task / result / error / control control: Dict[str, Any] field(default_factorydict) payload: Any None然后是业务上下文和编排状态dataclass class DataContext: trace_id: str store: Dict[str, Any] field(default_factorydict) def put(self, key: str, value: Any): self.store[key] value def get(self, key: str): return self.store.get(key) dataclass class ControlState: current_node: str start finished: list field(default_factorylist) retries: Dict[str, int] field(default_factorydict) status: str running # running / done / error / timeout4.2 编排器只要控制信息不看业务内容编排器的职责只有一个拿到信封后决定下一步动作。它从上一步返回的信封里读取msg_type和route判断当前流程处于什么状态然后确定下一个要发给谁。class Coordinator: def __init__(self, workflow: Dict[str, list], registry: Dict[str, Any]): self.workflow workflow self.registry registry self.state ControlState() def step(self, envelope: Envelope, data_ctx: DataContext) - Optional[Envelope]: # 流程开始进入入口节点 if self.state.current_node start: return self._to_agent(search_agent, data_ctx) # 出错时根据重试策略决定是重试还是终止 if envelope.msg_type error: route envelope.control.get(retry_target) current self.state.current_node self.state.retries[current] self.state.retries.get(current, 0) 1 if self.state.retries[current] 3: self.state.status error return None return self._to_agent(route, data_ctx) # 节点完成后取 workflow 中配置的下一跳 if envelope.msg_type result: current self.state.current_node self.state.finished.append(current) next_candidates self.workflow.get(current, []) if not next_candidates: self.state.status done return None # 简单规则取第一个可满足条件的节点实际项目可换成条件路由 next_node next_candidates[0] if next_node __done__: self.state.status done return None return self._to_agent(next_node, data_ctx) return None def _to_agent(self, route: str, data_ctx: DataContext) - Envelope: payload_key f{data_ctx.trace_id}:{route}:input self.state.current_node route return Envelope( trace_iddata_ctx.trace_id, routeroute, msg_typetask, payloadpayload_key, )这里注意一个设计编排器从不直接拿业务数据信封里的payload只是一个data_keyAgent 通过这个 key 去业务上下文里取数据。这样编排器就能保持对业务内容完全无知无论 Agent 里跑的是检索还是写作还是审核编排器只认流程节点名。4.3 Agent Worker管好自己的一亩三分地Agent 的写法被刻意做成没有脾气的样子从上下文里取自己的输入 key处理完往上下文里写输出 key返回一个表示成功的result信封。它不需要知道自己会被谁调用、下一步去哪里。class SearchAgent: name search_agent def run(self, envelope: Envelope, data_ctx: DataContext) - Envelope: input_key envelope.payload query data_ctx.get(input_key) # 模拟检索动作 result {query: query, materials: [material_a, material_b]} out_key f{data_ctx.trace_id}:{self.name}:output data_ctx.put(out_key, result) return Envelope( trace_idenvelope.trace_id, routeself.name, msg_typeresult, payloadout_key, )WriteAgent、ReviewAgent 的写法完全一样只是业务不同。它们之间没有任何两两杂交的逻辑全部通过业务上下文间接协作。这就是可复用的含义以后想换掉 SearchAgent只要新的 SearchAgent 满足从上下文读输入、向上下文写输出的协议直接替换不牵连其他 Agent。4.4 跑一下完整流程把流程图和 Agent 组装起来WORKFLOW { search_agent: [write_agent], write_agent: [review_agent], review_agent: [__done__], } REGISTRY { search_agent: SearchAgent(), write_agent: WriteAgent(), review_agent: ReviewAgent(), } def run_pipeline(initial_data: dict): trace_id str(uuid4()) data_ctx DataContext(trace_idtrace_id) data_ctx.put(f{trace_id}:start:input, initial_data) coordinator Coordinator(WORKFLOW, REGISTRY) env Envelope(trace_idtrace_id, routestart, msg_typeresult) env.msg_type task while coordinator.state.status running: agent_env coordinator.step(env, data_ctx) if agent_env is None: break agent REGISTRY[agent_env.route] env agent.run(agent_env, data_ctx) print(status:, coordinator.state.status) print(finished:, coordinator.state.finished) print(data keys:, list(data_ctx.store.keys())) run_pipeline({query: agent architecture})这段代码虽然简单但已经把数据流与控制流分离的骨架搭出来了业务 Agent 只处理上下文里的数据编排器只管理流程状态信封上的控制信息和载荷信息各就各位。5. 分离架构的实测坑位序列化、死锁、幂等与 Trace 排查5.1 序列化和消息体积没管住总线先撑爆我最初的设计里Envelope 的payload字段直接放完整数据对象Agent A 拿到查询结果后把整份文档塞进信封传给 Agent B。项目规模小的时候一切正常直到某次让 Agent 处理一个很大的项目文档内存和队列同时告急排查下来发现是几个大结果对象在队列里排着队每一份都带着完整副本。解决方式后来就是前面提过的引用传递数据落到共享存储信封只传data_key。这个改动带来的额外好处是Agent 的输入输出天然变成可寻址的想重放某次 Agent 处理过程直接把对应时间点的data_key再喂进去就行不用重新跑上游。5.2 循环依赖和超时控制流也会死锁流程图一旦复杂起来很容易出现两个 Agent 互为上下游的循环路径尤其是我用了审核不通过就退回重写这种带有回退边的结构时。编排器如果没有超时保护遇到循环逻辑就直接卡死。我的第一个教训流程定义加载时必须做环路检测。哪怕你的图设计时是 DAG运行时的条件路由也可能产生环最好在每条消息进入编排器时检查它的trace_id在当前节点上的访问次数超过阈值立刻按错误处理。第二个教训是给每个节点配置独立的超时时间Agent 处理超过时限后编排器直接把该节点标记为timeout并走重试或终止分支。超时参数应该放在信封的control里而不是写死在 Agent 里这样同一套 Agent 在不同业务场景下可以对超时表现不同。5.3 重试不幂等副作用会翻倍控制流天然要考虑失败重试但重试最容易被忽略的副作用是重复执行。比如一个 Agent 里包含调用外部 API 生成报告这个动作如果网络抖动导致响应丢失、编排器触发重试报告可能被生成两次。解决思路是让所有 Agent 对msg_id幂等。Agent 在执行业务前先检查这个msg_id是否处理过处理过就直接返回上次结果。实现上我习惯在业务上下文里维护一个processed_msg_ids集合配合msg_id做去重。这是数据流层面的要求和业务逻辑无关适合做成基类或装饰器而不是让每个 Agent 自己实现。5.4 没有 trace_id 贯穿日志流程问题查一次肝硬化一次分离之后 Agent 之间不再直连好处是解耦坏处是问题链路变长了。以前 Agent A 里报错顺着调用栈就能看到 B 和 C现在每个 Agent 的日志散落在不同模块没有统一的关联标识排查就像在深海捞针。trace_id就是用来解决这个问题的。所有入站消息生成时分配一个trace_id之后的每一步信封、每一条日志、每一次外部调用都带上它。我还在日志统一出口加了一个过滤器把trace_id、agent_name、msg_id提取成结构化字段这样排查时只要按trace_id做一次倒序检索整条链路的执行轨迹就出来了。看起来是个小改动但对多 Agent 系统的可维护性提升是决定性的。6. 从分离走向复用让 Agent 协作从写死走向装配6.1 协议先行所有 Agent 穿同一件马甲分离做到位后复用的第一步是定义一套 Agent 的统一接入协议。我习惯用适配器这个词来理解不管业务方拿来的 Agent 是什么形式哪怕是别人项目里的一个类、一个 HTTP 接口只要包一层适配器让它满足读信封、写上下文、返回结果信封的约定就能直接挂进协作系统。这样做带来一个很直接的好处Agent 的接入、移除、替换不再需要改动流程配置之外的任何代码。团队里新来的同事想加一个摘要 Agent只需要写一个类实现协议然后在配置文件里把它挂到流程图的某个节点上完毕。协作骨架完全不用碰。6.2 编排策略改为声明式配置既然控制流已经集中到编排器上流程定义就不该再藏在代码里而是应该变成一份看得见、改得动的声明式配置。这几乎是水到渠成的演变当每个 Agent 都不再关心下一步时下一步的逻辑就只剩下编排器里的一张路由表把它抽成 YAML 或 JSON 非常自然。我的配置结构参考了有向图模型每个节点声明type、next、retry、timeout条件分支用简单的if_payload_field表达式表示。人工审批节点、并行 fan-out 节点都可以在配置里加进去。编排器不需要为每种业务场景写一套代码它只是一个解释执行配置的引擎。到这里可复用才真正兑现复用的不是某个 Agent 的业务逻辑而是整套流程编排和消息传输的通用能力。6.3 这套架构目前还缺什么讲到这里这套分离式协作架构的核心已经完整了但它离生产级多 Agent 平台还有几段路。我目前正在补的几个方向也分享给你作参考Agent 注册中心不只是本地注册表而是带版本管理、健康检查、灰度发布能力的服务发现中心让运行中的流程可以安全升级节点实现。动态编排目前流程是配置里写死的下一步想做运行时动态规划——根据当前业务上下文让编排器内部嵌入策略决策逻辑自动决定下一个调谁而不只是读静态配置。人工审批节点很多业务不能全自动跑完需要插入一个人来处理高风险决策这个节点的控制流要能被编排器管理同时给人留出操作入口。这些方向本质上都没有跳出控制与数据分离这个底座。底座的稳定性决定了上层能长多大多复杂。我现在每次给别人讲多 Agent 协作时都会强调这句话一份干净的数据流配上一条可配置的控制流就是 Agent 协作系统能持续演化的底线。踩过几次坑之后你会发现最先做对这一步的项目后面加 Agent 的速度和信心完全不一样。
📝

华诺云谱内容团队

资深建站顾问 · 行业研究员

10年+企业数字化服务经验,专注智能建站、SEO优化与品牌营销,持续输出建站技巧、行业洞察与营销干货,已帮助5000+企业实现数字化增长。

你可能需要的服务

订阅华诺云谱资讯周报

每周一封,精选建站技巧、SEO与营销干货,直达邮箱。已有 8,000+ 企业主订阅,助你少走弯路。