第一次用AutoGen搭多智能体应用时,我还在0.2版本里打转,两个Agent之间靠initiate_chat你一言我一语地串流程。当时觉得只要把对话链设计得足够长,什么任务都能跑通。直到我把一个真实项目迁移到新版Core Runtime上,才意识到这套以消息路由为核心的内核,才是多智能体协作真正的地基,也是旧版那种“人形对话”模式撑不住的架构瓶颈。
这篇实战记录聚焦AutoGen Core Runtime的底层机制:事件总线、消息路由、Topic与订阅的关系、多智能体之间的消息协议设计,以及我在真实项目中遇到的几个坑。适合已经跑通过官方demo、想深入理解Runtime运作方式、准备把AutoGen用到实际业务集成里的开发者。下面讲的API形态基于autogen-core 0.4/0.5系列,不同小版本在方法名和包名上会有差异,但核心机制是一致的。
1. 先理解Core Runtime为什么要把“对话”改成“事件路由”
1.1 旧版对话式架构的三个软肋
AutoGen 0.2时代的编程模型很好上手:你创建一个ConversableAgent,然后调用initiate_chat,让两个Agent像两个人聊天一样完成一个任务。这个模型在demo里非常惊艳,但一旦进入真实业务,问题就出来了。
第一个问题是发起者和接收者强耦合。A想给B发一条消息,必须知道B的实例、B的回复方式、B的对话终止条件。想在中间插入一个C来审核对话,就得重构整条对话链路,而不是简单加一个订阅者。
第二个问题是对话历史被当作唯一状态。旧版把消息列表直接作为Agent记忆,但真实业务里往往需要“任务状态”而不只是“聊了些什么”。比如一次数据处理任务进行到哪一步、哪个子任务失败、该重试还是跳过,这些在对话链里很难表达清楚。
第三个问题是扩展性受限。三个Agent协作已经显得混乱,十个Agent协作时,调用关系直接变成一团乱麻。你很难回答一个问题:这条消息到底是谁处理的?谁该处理?如果没人处理该怎么办?
1.2 Core Runtime的核心对象一览
新版Core Runtime本质上是一个事件驱动的消息分发系统。它的核心不再是“Agent之间互相说话”,而是“Agent向运行时发布事件,运行时根据订阅关系把事件路由给对应的Agent”。我刚接触时最大的感受是:原来的代码思路是“直接呼叫”,现在变成了“发布订阅”,一开始很不习惯,但用久了会发现边界清晰得多。
这里涉及几个关键对象,我建议你先在脑子里建立一幅地图:
| 对象 | 角色定位 | 类比 |
|---|---|---|
AgentRuntime | 运行时抽象,管理Agent生命周期、消息分发 | 操作系统内核 |
SingleThreadedAgentRuntime | 默认的单线程实现,按事件循环调度 | 单核CPU |
Event/ 消息类 | 跨Agent传递的数据载体 | 快递包裹 |
TopicId | 消息发布的主题地址 | 快递单上的收货地址 |
Subscription | 某个Agent对某个Topic/消息类型的订阅绑定 | 订阅报纸的登记表 |
MessageContext | 消息的元信息(发送者、Topic、来源等) | 快递面单 |
从这张表可以看出,Agent和Agent之间其实不需要互相认识。A发布消息到Topic,谁订阅了这个Topic,谁就会收到消息。这个解耦方式让多智能体系统的扩展变得非常自然。
1.3 你需要转换的三个核心思维
根据我迁移项目的实际体验,从旧版思维切到Core Runtime,有三个转变至关重要。
从“调用别人”到“发布消息”。旧版里A调用B的generate_reply,是同步等待B的结果。新版里A只负责publish_message,发完就继续干自己的事,或者挂起等待后续事件。结果什么时候回来、由谁回来,完全由运行时决定。这听起来麻烦,但却是并行和扩展的前提。
从“对话历史”到“消息日志”。旧版里每个Agent的内存里都存着一份对话历史,用来决定下一步说什么。新版里消息是点在运行时层面的,每条消息都是不可变的快照。你要做的是设计好消息类型、定义好交互协议,而不是操心“对方上句话说了什么”。
从“Agent内部状态”到“运行时管理状态”。旧版里Agent自己记住自己干到哪了,新版里推荐把任务状态放到一个专门的协调Agent里,或者外置到存储层。这样好处是Agent可以随时重建,任务进度不会丢。
2. 消息路由的完整链路:Topic、订阅与消息类型
2.1 一次完整消息分发要经过几个环节
很多人第一次写Core Runtime,代码能跑但不知道消息到底是怎么从A到B的。先把这个链路搞清楚,后面遇到问题才有排查头绪。
一次完整的消息分发包含这几个环节:
- 某个Agent调用
publish_message(message, topic_id)提交消息到运行时。 - 运行时根据
topic_id找到该Topic下所有的Subscription。 - 每个
Subscription内部有一个message_type,运行时用isinstance检查消息类型是否匹配。 - 匹配成功的消息被投递到对应Agent的消息队列。
- Agent运行时取出消息,调用带有
@message_handler装饰器的方法,并把MessageContext一起传进去。 - handler执行完毕,消息生命周期结束。
这个流程和旧版的generate_reply有本质区别。旧版是“你直接打电话给某人问答案”,新版是“你把信投进邮筒,邮局负责分拣,谁订了这封类型的信息,谁就收到”。如果你发布了一条消息但没有任何Agent订阅匹配的消息类型,它会被运行时直接丢弃,不会报错。这个特性后面还要踩坑,先记着。
2.2 消息类型设计的三种风格
消息类型是Core Runtime里最该花心思的地方。它的设计决定了你的路由粒度、扩展成本和调试难度。
我在项目里总结出三种风格:
粗粒度消息:整个系统只有一两个消息类,比如TaskMessage,里面加一个type字段区分场景。好处是简单,坏处是所有handler都绑到同一个类上,路由判断退化成if-else,订阅关系变得毫无意义。适合只有一个Agent在消费的小Demo。
细粒度消息:每个业务动作一个消息类,比如TranslateRequest、SummarizeRequest、ReviewRequest。好处是订阅关系一目了然,坏处是类会爆炸。适合任务边界清晰、Agent职责分离的系统。
协议式消息:把一组相关的处理合并成一个消息类,里面用Header+Payload的组合。比如WorkflowRequest里有action字段、request_id字段、data字段。这种方式兼具扩展性和可控性,是目前我认为最工程化的做法。
from dataclasses import dataclass, field from autogen_core import Event @dataclass class WorkflowRequest(Event): action: str request_id: str payload: dict = field(default_factory=dict) @dataclass class WorkflowResponse(Event): action: str request_id: str result: dict = field(default_factory=dict)这样设计的好处是:你可以为action的每种取值写一个专属handler,同时订阅关系只需要绑定到WorkflowRequest这一个类型上面。
2.3 路由判定的两种模式
Core Runtime本身支持按类型路由,也就是订阅时指定message_type。但如果你的业务里,同一种消息需要不同Agent按内容决定谁来处理,这就是按内容路由。
按类型路由适合任务类型天然划分清晰的场景。比如Writer只处理WritingRequest,Reviewer只处理ReviewRequest,订阅关系写死即可。
按内容路由适合需要动态分发的场景,比如“一条任务消息根据target_role字段决定转发给谁”。这种情况下,通常需要一个RouterAgent来充当转发节点:它订阅一个总Topic,收到消息后读取内容字段,然后把消息重新发布到不同的子Topic。
我比较推荐的组合是:外部消息统一进总Topic,RouterAgent做内容判断,子Topic做职责隔离。这样路由规则集中在一个地方,方便审阅和修改。
3. 实战:用消息路由搭一个三角色协作系统
3.1 场景设计:一个三角色内容工坊
我拿一个“内容工坊”场景来讲,它足够小,能说清楚原理;又足够典型,能推广到大多数多Agent业务。
系统里有三个角色:
- Planner:收到任务后拆解需求,确定由谁执行。
- Writer:负责写初稿。
- Reviewer:负责审稿并给出修改意见。
传统写法里这三个角色要互相持有对方的引用,消息链路绕来绕去。在Core Runtime里,我让它们之间完全解耦,只依赖消息类型和Topic。
3.2 定义消息协议
根据前面说的协议式消息思路,我来定义消息。TaskRequest是入口消息,DraftCreated是Writer的输出,ReviewFeedback是Reviewer的反馈。
from dataclasses import dataclass from autogen_core import Event @dataclass class TaskRequest(Event): request_id: str content: str target_role: str @dataclass class DraftCreated(Event): request_id: str draft: str @dataclass class ReviewFeedback(Event): request_id: str feedback: str approved: bool @dataclass class TaskComplete(Event): request_id: str final_text: str这里每个消息都带request_id,目的很明确:消息在异步路由过程中会四处漂移,只有靠request_id才能把同一次任务的多个事件关联起来。没有这个字段,后面做聚合和终态判定都要抓瞎。
3.3 实现RouterAgent与WorkerAgent
先写RouterAgent。它的任务是订阅总入口Topic,读取target_role字段,把TaskRequest转发到对应子Topic。
from autogen_core import ( AgentRuntime, RoutedAgent, DefaultTopicId, MessageContext, message_handler, ) class RouterAgent(RoutedAgent): def __init__(self, runtime: AgentRuntime): super().__init__("router", runtime) @message_handler async def on_task_request(self, message: TaskRequest, ctx: MessageContext) -> None: # 按内容字段动态决定路由目标 if message.target_role == "writer": await self.publish_message( message, topic_id=DefaultTopicId("writer", source="router"), ) elif message.target_role == "reviewer": await self.publish_message( message, topic_id=DefaultTopicId("reviewer", source="router"), )WriterAgent订阅writer这条Topic,收到TaskRequest后执行生成动作,产出DraftCreated,再发布到draft这条Topic。
class WriterAgent(RoutedAgent): def __init__(self, runtime: AgentRuntime): super().__init__("writer", runtime) @message_handler async def on_task_request(self, message: TaskRequest, ctx: MessageContext) -> None: # 这里替换成真实的模型调用或业务逻辑 draft = f"这里是 {message.content} 的初稿" await self.publish_message( DraftCreated(request_id=message.request_id, draft=draft), topic_id=DefaultTopicId("draft", source="writer"), )ReviewerAgent订阅draftTopic,收到草稿后给出反馈。如果确认通过,就发布TaskComplete。
class ReviewerAgent(RoutedAgent): def __init__(self, runtime: AgentRuntime): super().__init__("reviewer", runtime) @message_handler async def on_draft_created(self, message: DraftCreated, ctx: MessageContext) -> None: # 模拟评审逻辑 if len(message.draft) < 10: await self.publish_message( TaskComplete(request_id=message.request_id, final_text=message.draft), topic_id=DefaultTopicId("complete", source="reviewer"), ) else: await self.publish_message( ReviewFeedback( request_id=message.request_id, feedback="需要精简", approved=False, ), topic_id=DefaultTopicId("review", source="reviewer"), )3.4 启动运行时并注册订阅
主程序里把Agent注册进去,并添加订阅关系。这里有一点需要注意:订阅关系是绑定topic_type和message_type的,不要漏掉。
import asyncio from autogen_core import ( SingleThreadedAgentRuntime, TypeSubscription, ) async def main(): runtime = SingleThreadedAgentRuntime() await runtime.try_register_agent("router", lambda: RouterAgent(runtime)) await runtime.try_register_agent("writer", lambda: WriterAgent(runtime)) await runtime.try_register_agent("reviewer", lambda: ReviewerAgent(runtime)) # 外部任务先进总入口Topic await runtime.add_subscription(TypeSubscription( topic_type="entry", message_type=TaskRequest, agent_type="router", )) await runtime.add_subscription(TypeSubscription( topic_type="writer", message_type=TaskRequest, agent_type="writer", )) await runtime.add_subscription(TypeSubscription( topic_type="draft", message_type=DraftCreated, agent_type="reviewer", )) runtime.start() await runtime.publish_message( TaskRequest(request_id="req-001", content="写一篇关于AutoGen的实战文章", target_role="writer"), topic_id=DefaultTopicId("entry", source="main"), ) await runtime.stop_when_idle() if __name__ == "__main__": asyncio.run(main())你可能会问:writer这条Topic上订阅的也是TaskRequest,那Router转发出去的消息,Writer收到后怎么知道这是给自己的?因为Router发布时指定了topic_id=DefaultTopicId("writer"),只有订阅了writer这个topic的WriterAgent才会收到。这里的威力在于:Router不需要知道Writer的实例,甚至不需要知道Writer存在,它只需要知道有一条叫writer的路由信道。
3.5 看日志理解事件流转顺序
我把这个Demo跑起来以后,事件流转顺序是这样的:
main向entry发布TaskRequest(req-001)。- RouterAgent收到
TaskRequest,因为target_role=writer,转发到writer。 - WriterAgent收到
TaskRequest,生成初稿,发布DraftCreated到draft。 - ReviewerAgent收到
DraftCreated,给出评审反馈或TaskComplete。
整个流程里,没有任何一个Agent持有另一个Agent的引用。你把Writer换成另一个完全不同的实现,只要它还订阅writer、还发送DraftCreated,系统其余部分完全不用动。这种可替换性,对真实业务太重要了。
4. 实战中的踩坑实录:路由静默丢失、阻塞与异常吞噬
4.1 订阅注册晚于消息发布,消息被静默丢弃
我第一次在多Agent系统中加入消息路由时,遇到一个“明明发布了消息,但没有任何Agent响应”的问题。查了半天发现:订阅是在runtime.start()之后才注册的。
Core Runtime的处理逻辑是:在消息发布的那一瞬间,运行时去查找当前已有的订阅关系,如果找不到匹配订阅,消息就被丢弃,而不会给你任何报错。这跟数据库外键约束不同,更像是UDP包,没人要就扔。
解决方案很简单:在runtime.start()之前完成所有add_subscription调用。如果你是在Agent的初始化方法里做订阅,务必确保那个Agent在发送第一条消息之前已经注册完毕。后期排查此类问题时,我最常用的手段是给每个消息加日志,打印消息的id和topic_id,这样能快速定位是哪一环没接上。
4.2 Agent内部异常被吞掉,事件链无声断裂
另一个让我头疼的问题,是某个Agent在处理消息时抛了异常,但整个程序完全不崩溃,后续Agent也收不到任何数据。看起来就像消息进入了黑洞。
原因是Core Runtime在调度handler时,对于未捕获的异常默认只记录日志,并不会向消息发送方返回错误。如果你没看日志,就会误以为“没人处理这条消息”。这种情况下,最糟糕的做法是到处加print,最有效的做法是设计一个错误传播协议。
我通常这样处理:在每个重要的handler上用try/except包裹,捕获后发布一条ErrorMessage,携带原始消息的request_id和异常信息,让协调者可以感知失败并决定重试或终止。
from autogen_core import Event, MessageContext, message_handler @dataclass class ErrorMessage(Event): request_id: str error: str source_agent: str class WriterAgent(RoutedAgent): @message_handler async def on_task_request(self, message: TaskRequest, ctx: MessageContext) -> None: try: draft = self._generate(message.content) await self.publish_message( DraftCreated(request_id=message.request_id, draft=draft), topic_id=DefaultTopicId("draft", source="writer"), ) except Exception as e: await self.publish_message( ErrorMessage(request_id=message.request_id, error=str(e), source_agent=self.id.type), topic_id=DefaultTopicId("errors", source="writer"), )4.3 单线程运行时里做耗时同步调用,整个系统“卡死”
Core Runtime默认的SingleThreadedAgentRuntime基于单个事件循环调度。它本身又快又轻,但有代价:如果你在handler里用requests.get、time.sleep这种阻塞调用,整个运行时都会被卡住。
我踩过这个坑之后,总结出的规则是:
- 纯CPU密集任务,用
asyncio.to_thread把它扔到线程池。 - I/O密集任务(HTTP、数据库),优先用
httpx.AsyncClient、aiomysql这类异步SDK。 - 如果一个Agent要做的事情确实很重,考虑把它拆成两个Agent,一个负责接收,一个负责处理。
这里有个容易忽略的点:即便你用了asyncio.create_task把阻塞任务丢到后台,如果后台任务里又用了同步请求,还是会阻塞事件循环。检查手段很简单:在日志里记录每个handler的耗时,一旦发现某个agent处理时间异常,基本就是它内部有同步阻塞。
4.4 消息体里放了可变对象,导致数据被意外修改
Core Runtime里消息按引用传递时,如果消息类里有一个list或dict字段,多个Agent拿到的是同一个对象引用。A agent往里面塞了数据,B agent看到的就是被修改后的东西。表面上看是“共享状态”,实际上会让消息日志完全失真,排查问题时非常痛苦。
安全做法是:消息字段尽量用不可变类型或深拷贝。发布数据时,如果是自己构造的对象,可以copy.deepcopy后再放入消息。虽然有点损耗,但换来的是消息不可变性和可追溯性,在分布式多Agent系统里非常值。
5. 让它更像生产系统:终态判定、超时与外置状态
5.1 终态判定需要聚合层,而不仅是单个Agent
上面那个三角色Demo里,Reviewer发布了TaskComplete,看起来流程结束了。但在真实系统里,“任务完成”往往不是某个Agent一拍脑袋决定的,而是需要根据多个子事件聚合判断。
比如一个内容工坊系统,Writer要产出文案,GraphicAgent要产出配图,只有两者都完成,任务才算真正结束。这种情况下,需要一个CoordinatorAgent专门维护任务状态表,收到DraftCreated就标记“文案完成”,收到ImageReady就标记“配图完成”,两个标记都到位后,才对外宣称“任务完成”。
class CoordinatorAgent(RoutedAgent): def __init__(self, runtime: AgentRuntime): super().__init__("coordinator", runtime) self._task_states: dict[str, set[str]] = {} @message_handler async def on_draft_created(self, message: DraftCreated, ctx: MessageContext) -> None: state = self._task_states.setdefault(message.request_id, set()) state.add("draft") await self._try_complete(message.request_id) @message_handler async def on_image_ready(self, message: ImageReady, ctx: MessageContext) -> None: state = self._task_states.setdefault(message.request_id, set()) state.add("image") await self._try_complete(message.request_id) async def _try_complete(self, request_id: str) -> None: state = self._task_states[request_id] if "draft" in state and "image" in state: await self.publish_message( TaskComplete(request_id=request_id, final_text="聚合完成"), topic_id=DefaultTopicId("complete", source="coordinator"), )这种聚合逻辑在对话式框架里很难写,因为对话流天然是线性的;在Core Runtime里却非常自然,因为每条消息都是独立的事件,调度器只需要按request_id合并就行。
5.2 用超时兜底,别让运行时永远等下去
事件驱动系统有一个通病:如果某个Agent没响应,整个任务就可能无限挂起。比如外部大模型API超时、Agent内部死循环,都会导致最终没有人发布TaskComplete。
工程上我习惯做两件事:
第一,给任务加超时。外层用asyncio.wait_for包住整个流程,超时后直接标记任务失败并对外返回。第二,给stop_when_idle加监听逻辑。当运行时进入空闲状态但任务还没完成,多半是某条消息没被消费或者某个Agent异常退出。此时应该触发兜底逻辑,而不是简单地认为“没事做了”。
try: await asyncio.wait_for(wait_all_complete(), timeout=60) except asyncio.TimeoutError: await runtime.publish_message( ErrorMessage(request_id="req-001", error="task timeout", source_agent="main"), topic_id=DefaultTopicId("errors", source="main"), )5.3 状态外置:重启之后还能恢复
最后这点是我在生产环境里最看重的能力:任务状态一定要能外置。
Core Runtime的单线程版本跑在进程内,如果服务重启,所有Agent内存里的状态都会丢失。我的做法是:把每个Agent已消费的消息记录到持久化存储(数据库或消息表),每次重启后,从数据库读取任务快照,再重新把未完成的消息投递到对应Topic。因为消息是事件驱动的,只要每个handler是幂等的——处理同样两条消息不会产生重复副作用——整个系统就可以从崩溃中恢复。
这里有一个可供落地的方案:在publish_message前把消息序列化到数据库的outbox表,标记为“待发布”;消费端收到消息后,在数据库的inbox表记录“已处理”的message_id。如果重启,扫描两张表找出没处理完的消息再次投递。这本质上就是事件溯源和消息队列的思路。Core Runtime本身不提供这个能力,但它的消息协议足够清晰,让我可以在外部实现完整的状态恢复,这是旧版对话式框架做不到的。
我个人在实际操作中的体会是:Core Runtime初看比0.2版本复杂,一旦你习惯了“消息类型即协议、Topic即边界、订阅即关系”这套思维,再去设计多智能体系统反而轻松得多。以前我要操心三个人之间彼此怎么喊话,现在只需要定义好快递包裹的规格,摆好收货柜,剩下的事交给运行时自己流动。最后再分享一个小技巧:刚开始上手时,别急着写复杂业务,先搭两个Agent和一个Router,把消息流转日志打全,盯着看十分钟,你对这套事件路由的感觉会完全不一样。