
第一次用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默认的单线程实现按事件循环调度单核CPUEvent/ 消息类跨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职责分离的系统。协议式消息把一组相关的处理合并成一个消息类里面用HeaderPayload的组合。比如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_factorydict) dataclass class WorkflowResponse(Event): action: str request_id: str result: dict field(default_factorydict)这样设计的好处是你可以为action的每种取值写一个专属handler同时订阅关系只需要绑定到WorkflowRequest这一个类型上面。2.3 路由判定的两种模式Core Runtime本身支持按类型路由也就是订阅时指定message_type。但如果你的业务里同一种消息需要不同Agent按内容决定谁来处理这就是按内容路由。按类型路由适合任务类型天然划分清晰的场景。比如Writer只处理WritingRequestReviewer只处理ReviewRequest订阅关系写死即可。按内容路由适合需要动态分发的场景比如“一条任务消息根据target_role字段决定转发给谁”。这种情况下通常需要一个RouterAgent来充当转发节点它订阅一个总Topic收到消息后读取内容字段然后把消息重新发布到不同的子Topic。我比较推荐的组合是外部消息统一进总TopicRouterAgent做内容判断子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_idDefaultTopicId(writer, sourcerouter), ) elif message.target_role reviewer: await self.publish_message( message, topic_idDefaultTopicId(reviewer, sourcerouter), )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_idmessage.request_id, draftdraft), topic_idDefaultTopicId(draft, sourcewriter), )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_idmessage.request_id, final_textmessage.draft), topic_idDefaultTopicId(complete, sourcereviewer), ) else: await self.publish_message( ReviewFeedback( request_idmessage.request_id, feedback需要精简, approvedFalse, ), topic_idDefaultTopicId(review, sourcereviewer), )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_typeentry, message_typeTaskRequest, agent_typerouter, )) await runtime.add_subscription(TypeSubscription( topic_typewriter, message_typeTaskRequest, agent_typewriter, )) await runtime.add_subscription(TypeSubscription( topic_typedraft, message_typeDraftCreated, agent_typereviewer, )) runtime.start() await runtime.publish_message( TaskRequest(request_idreq-001, content写一篇关于AutoGen的实战文章, target_rolewriter), topic_idDefaultTopicId(entry, sourcemain), ) await runtime.stop_when_idle() if __name__ __main__: asyncio.run(main())你可能会问writer这条Topic上订阅的也是TaskRequest那Router转发出去的消息Writer收到后怎么知道这是给自己的因为Router发布时指定了topic_idDefaultTopicId(writer)只有订阅了writer这个topic的WriterAgent才会收到。这里的威力在于Router不需要知道Writer的实例甚至不需要知道Writer存在它只需要知道有一条叫writer的路由信道。3.5 看日志理解事件流转顺序我把这个Demo跑起来以后事件流转顺序是这样的main向entry发布TaskRequest(req-001)。RouterAgent收到TaskRequest因为target_rolewriter转发到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_idmessage.request_id, draftdraft), topic_idDefaultTopicId(draft, sourcewriter), ) except Exception as e: await self.publish_message( ErrorMessage(request_idmessage.request_id, errorstr(e), source_agentself.id.type), topic_idDefaultTopicId(errors, sourcewriter), )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_idrequest_id, final_text聚合完成), topic_idDefaultTopicId(complete, sourcecoordinator), )这种聚合逻辑在对话式框架里很难写因为对话流天然是线性的在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(), timeout60) except asyncio.TimeoutError: await runtime.publish_message( ErrorMessage(request_idreq-001, errortask timeout, source_agentmain), topic_idDefaultTopicId(errors, sourcemain), )5.3 状态外置重启之后还能恢复最后这点是我在生产环境里最看重的能力任务状态一定要能外置。Core Runtime的单线程版本跑在进程内如果服务重启所有Agent内存里的状态都会丢失。我的做法是把每个Agent已消费的消息记录到持久化存储数据库或消息表每次重启后从数据库读取任务快照再重新把未完成的消息投递到对应Topic。因为消息是事件驱动的只要每个handler是幂等的——处理同样两条消息不会产生重复副作用——整个系统就可以从崩溃中恢复。这里有一个可供落地的方案在publish_message前把消息序列化到数据库的outbox表标记为“待发布”消费端收到消息后在数据库的inbox表记录“已处理”的message_id。如果重启扫描两张表找出没处理完的消息再次投递。这本质上就是事件溯源和消息队列的思路。Core Runtime本身不提供这个能力但它的消息协议足够清晰让我可以在外部实现完整的状态恢复这是旧版对话式框架做不到的。我个人在实际操作中的体会是Core Runtime初看比0.2版本复杂一旦你习惯了“消息类型即协议、Topic即边界、订阅即关系”这套思维再去设计多智能体系统反而轻松得多。以前我要操心三个人之间彼此怎么喊话现在只需要定义好快递包裹的规格摆好收货柜剩下的事交给运行时自己流动。最后再分享一个小技巧刚开始上手时别急着写复杂业务先搭两个Agent和一个Router把消息流转日志打全盯着看十分钟你对这套事件路由的感觉会完全不一样。