AutoGen(Python)消息广播核心概念完全指南:Topic、Subscription 与 Type-Based Subscription 的深入解析
本指南围绕仓库 python/docs/src/user-guide/core-user-guide/core-concepts/topic-and-subscription.md 展开,系统讲解 AutoGen Agent 运行时(agent runtime)中的广播消息模型——Topic(主题)与 Subscription(订阅)。读完本文,你将理解 Topic 的类型/来源二元结构与命名规则、Subscription 到 Agent 的映射机制、为何基于类型的订阅是推荐做法,以及如何将单租户单主题、单租户多主题、多租户等典型场景落地为可运行的订阅与发布代码。
两种消息投递方式:点对点(Direct Messaging)与广播(Broadcast)
AutoGen 的 Agent 运行时向 Agent 投递消息有两种方式,它们在运行时抽象 _agent_runtime.py中分别对应两类核心 API:
- 直接消息(Direct Messaging):一对一投递。发送方必须明确给出接收方的 Agent ID,对应运行时抽象中的
send_message类接口; - 广播(Broadcast):一对多投递。发送方不需要指定任何接收方 Agent ID,只需要指明广播到哪个 Topic,由运行时根据既有的订阅关系决定消息送到谁手上。
Topic 本质上是 Agent ID 之上的一层"间接寻址"(indirection):广播让生产者和消费者完全解耦,发送方不感知、也不依赖接收方的具体身份。
许多场景天然适合广播。典型的是事件驱动的工作流:Agent 之间往往互不感知、互不依赖,谁会对某类消息感兴趣、谁会来消费,是由发布者所不知道的订阅关系动态决定的。因此,理解广播模型,核心就是理解 Topic(广播的边界) 与 Subscription(边界到 Agent 的映射) 这两个概念。
Topic:一次广播消息的作用范围
Topic 的二元构成
根据关联文档的定义:
Topic = (Topic Type, Topic Source)
一个 Topic 由两部分组成:
- Topic Type(主题类型):通常由应用代码定义,用来标记该主题承载的是哪一类消息。例如,一个处理 GitHub 事件的 Agent,在发布"新 issue"消息时可以用
"GitHub_Issues"作为 topic type; - Topic Source(主题来源):同一 topic type 下唯一标识一个 Topic 的值,通常由应用数据决定。例如该 GitHub Agent 可以用
"github.com/{repo_name}/issues/{issue_number}"这样的字符串作为 topic source,从而唯一锁定"某一个仓库的某一个 issue"。
Topic Source 是发布者划定消息边界、制造数据隔离(silo)的关键手段:同一个 topic type 下不同的 source 之间天然互相隔离,互不串扰。
Topic 的字符串表示与命名约束
Topic ID 可以与字符串互相转换,格式为:
Topic_Type/Topic_Source
关联文档给出的可读性规则是:
- Topic Type:必须是 UTF-8 文本,仅允许字母(a-z)、数字(0-9)和下划线(_);不能以数字开头,不能包含任何空格;
- Topic Source:必须是 UTF-8 文本,字符范围须落在 ASCII 32(空格)到 126(
~)之间(含端点)。
在此基础上,Python 端实际的校验规则可以在 TopicId 实现中得到印证。autogen_core._topic.TopicId 是一个 @dataclass(eq=True, frozen=True),其 type 字段在构造时用正则 ^[\w\-\.\:\=]+\Z 校验,命中正则之外的内容会抛出 ValueError。可以看到 Python 实现实际允许的字符比上文描述更宽(额外允许 -、.、:、=,且允许数字开头)。从字段 docstring 看,type 与 source 均对齐 CloudEvents 规范 的 type 与 source 语义。实际开发时以当前仓库实现为准。
同时,TopicId 提供两个双向转换方法:
str(TopicId(type="github_issues", source="github.com/example/autogen/issues/1"))
# "github_issues/github.com/example/autogen/issues/1"
TopicId.from_str("github_issues/github.com/example/autogen/issues/1")
# TopicId(type="github_issues", source="github.com/example/autogen/issues/1")
实现见 python/packages/autogen-core/src/autogen_core/_topic.py#L37-L47:__str__ 以 type/source 拼接,from_str 则用 split("/", maxsplit=1) 仅切分第一段,因此 source 内部即使包含 / 也不会破坏解析。
Subscription:把 Topic 映射到 Agent
如果说 Topic 回答"这条广播是关于什么的",那么 Subscription 回答"这条广播应该送到谁"——它把 Topic 映射到 Agent ID:
运行时维护着一张订阅表,并依据这张表把广播消息投递给各 Agent:
- 某个 Topic 没有任何订阅:发布到该 Topic 的消息不会送达任何 Agent(消息对当前应用而言无消费者);
- 某个 Topic 有多个订阅:消息会依据全部匹配的订阅被投递,且同一接收方只会收到一次;
- 应用可以通过运行时 API 动态添加或移除订阅。
在 Python 抽象层 _subscription.py中,Subscription 被定义为一个可运行时检查的协议(@runtime_checkable Protocol),它描述任何订阅实现都必须具备的契约:
| 成员 | 职责 |
|---|---|
id |
订阅的唯一标识,通常是一个 UUID,用于去重与移除 |
is_match(topic_id) |
判断某个 Topic 是否与该订阅匹配 |
map_to_agent(topic_id) |
把匹配的 Topic 映射到应处理它的 AgentId;仅应在 is_match 返回 True 时调用,否则实现可抛 CantHandleException |
__eq__ |
默认按 id 判定两个订阅是否相等 |
运行时是如何维护与匹配订阅的
在单进程实现的内部,SubscriptionManager 承担订阅簿记工作,可以从源码看到三条关键逻辑:
add_subscription:若已存在相同订阅则抛ValueError,否则追加并重建"Topic → 接收方 AgentId"索引(源码 L42-L48);remove_subscription(id):按订阅 id 移除,不存在时抛ValueError,移除后同样重建索引(源码 L50-L61);- 懒构建索引:当一个从未见过的 Topic 第一次被查询时,
_build_for_new_topic会遍历全部订阅,凡is_match(topic)为真,就把map_to_agent(topic)的结果追加进该 Topic 的接收方列表(源码 L63-L78)。
也就是说,"多订阅命中同一 Topic、每个订阅各映射出一个接收方"的汇总正是在这里完成的。对应用开发者而言,只需面向运行时抽象 API 操作订阅即可,例如 AgentRuntime 接口中声明的 add_subscription(subscription) 与 remove_subscription(id)(见 python/packages/autogen-core/src/autogen_core/_agent_runtime.py#L268-L285)。
Type-Based Subscription:基于类型的订阅
上述 Subscription 只是协议,框架为绝大多数场景准备了一个推荐使用的具体实现:TypeSubscription(在 Python 端位于 autogen_core 顶层命名空间,即 from autogen_core import TypeSubscription,源码见 python/packages/autogen-core/src/autogen_core/_type_subscription.py)。
它的声明形式是:
Type-Based Subscription = Topic Type --> Agent Type
即把一个 topic type 映射到一个 agent type。它是一种"无界"的声明式映射——无需预先知道任何具体的 topic source 或 agent key:
- 凡是 topic type 与该订阅匹配的 Topic,都会被映射到一个 Agent;
- 这个 Agent 的类型取订阅声明的 agent type;
- 这个 Agent 的 key(agent key)取该 Topic 的 source 值。
关联文档中还同步给出了相关概念 Agent ID 的说明(Agent ID 同样由 type 与 key 两部分组成,其字符串形态为 AgentType/Key,详见 Agent 标识与生命周期文档)。
从 TypeSubscription 的实现(python/packages/autogen-core/src/autogen_core/_type_subscription.py#L53-L60)可以清晰看到这一机制的两个核心方法:
def is_match(self, topic_id: TopicId) -> bool:
return topic_id.type == self._topic_type
def map_to_agent(self, topic_id: TopicId) -> AgentId:
if not self.is_match(topic_id):
raise CantHandleException("TopicId does not match the subscription")
return AgentId(type=self._agent_type, key=topic_id.source)
is_match只比较 topic type,与 source 无关;map_to_agent直接构造AgentId(type=订阅的agent_type, key=Topic的source)——topic source 原样下沉成为 agent key。
类文档字符串里的例子更直观地展示了这一点:TypeSubscription(topic_type="t1", agent_type="a1") 一旦注册,source 为 s1 的 Topic ("t1","s1") 会被 a1/s1 处理,source 为 s2 的 Topic ("t1","s2") 会被 a1/s2 处理——每个 source 都拥有自己独立的 Agent 实例,天然形成数据隔离。
细节:
TypeSubscription.__init__(topic_type, agent_type, id=None)中id缺省时自动生成 UUID;agent_type同时接受字符串或AgentType对象。两个订阅的相等判定为"id 相同,或 agent_type 与 topic_type 都相同"(源码 L62-L66)。
为什么它更受推荐?因为它可移植且与具体数据解耦:开发者无需编写依赖特定 Agent ID 的应用代码,即使 topic source 千变万化,一条声明即可覆盖无穷多个具体 Topic。
典型应用场景:基于类型订阅的落地模式
把 TypeSubscription 用在实际工程中,通常会从两个维度去拆解需求:
- 是否多租户(single-tenant / multi-tenant);
- 每个租户下是单个 Topic 还是多个 Topic(single / multiple topics per tenant)。
这里的"租户(tenant)"通常指负责某一特定用户会话或某一特定请求的一组 Agent。
场景一:单租户 + 单 Topic
这是最简单、也最常见的场景,适合命令行工具、单用户应用等只有一套会话的场合。做法是:
- 为每一个 agent type 各创建一个
TypeSubscription; - 所有订阅使用同一个 topic type;
- 发布时始终使用同一个 Topic(即相同的 topic type 与 topic source)。
延续关联文档的例子:假设有 "triage_agent"、"coder_agent"、"reviewer_agent" 三个 agent type,topic type 取 "default",则订阅如下:
# Type-based Subscriptions for single-tenant, single topic scenario
TypeSubscription(topic_type="default", agent_type="triage_agent")
TypeSubscription(topic_type="default", agent_type="coder_agent")
TypeSubscription(topic_type="default", agent_type="reviewer_agent")
发布时所有消息都用 source "default",于是 Topic 恒为 ("default", "default")。发到该 Topic 的消息会送达上面三种类型的全部 Agent,即产生以下 Agent ID:
# The agent IDs created based on the topic source
AgentID("triage_agent", "default")
AgentID("coder_agent", "default")
AgentID("reviewer_agent", "default")
一种常见的落地形态是:三个 Agent 各自承担三阶段工作(分流 → 编码 → 审查),全部订阅同一 Topic,任何发布到 ("default", "default") 的消息都会以广播形式同时进入这三个 Agent 的消息流。若目标 Agent ID 尚不存在,运行时会为其创建实例(前提是该 agent type 已在运行时注册了相应的工厂与处理逻辑,相关约定见 Agent 标识与生命周期文档)。
场景二:单租户 + 多 Topic
当同一个应用里希望控制不同 Agent 各管一类 Topic、制造订阅层面的"分工与隔离"时,就用多 Topic 模式。做法是:
- 为每个 agent type 创建订阅,但使用不同的 topic type;
- 若希望若干 agent type 共享同一类 Topic,就把同一个 topic type 映射给多个 agent type;
- topic source 在发布时依旧统一使用同一个值。
继续沿用三个 agent type 的例子:
# Type-based Subscriptions for single-tenant, multiple topics scenario
TypeSubscription(topic_type="triage", agent_type="triage_agent")
TypeSubscription(topic_type="coding", agent_type="coder_agent")
TypeSubscription(topic_type="coding", agent_type="reviewer_agent")
于是:
- 发布到
("triage", "default")的消息 → 只投递给triage_agent; - 发布到
("coding", "default")的消息 → 同时投递给coder_agent和reviewer_agent。
场景三:多租户场景
单租户场景里 topic source 是"写死在代码里"的常量(如 "default");一旦进入多租户场景,topic source 就变成了随数据而变的变量。
一个判断你是否处于多租户场景的可靠信号:你需要同一 agent type 的多个实例。例如用不同 Agent 实例分别处理不同用户会话以隔离私有数据,或让同一 agent type 的多个实例并行分担重负载。
以 GitHub issue 处理为例。假如只为 "triage_agent" 声明了订阅:
TypeSubscription(topic_type="github_issues", agent_type="triage_agent")
那么:
- 消息发布到 Topic
("github_issues", "github.com/microsoft/autogen/issues/1")→ 投递给 Agent ID("triage_agent", "github.com/microsoft/autogen/issues/1"); - 消息发布到 Topic
("github_issues", "github.com/microsoft/autogen/issues/9")→ 投递给 Agent ID("triage_agent", "github.com/microsoft/autogen/issues/9")。
这里的 Agent ID 完全由数据决定:不同 issue 的广播被路由给不同的 triage_agent 实例,运行时会为尚不存在的目标实例创建新 Agent,从而做到按 issue 粒度的处理隔离——这正是"topic source 下沉为 agent key"这一设计在多租户场景下的威力。
如果每个租户内部还想进一步细分(即"每个租户多个 Topic"),做法与单租户多 Topic 一致:使用不同的 topic type 即可。例如每个 issue 租户下同时有 "github_issues" 与 "github_discussions" 两个 topic type,分别路由给负责 issue 与 discussion 的 Agent。
发布与订阅:一个端到端的组合示例
把上述概念串起来,一次完整的广播消息消费大致长这样(方法签名与语义均取自当前仓库的 AgentRuntime 抽象与 TypeSubscription 实现):
from autogen_core import AgentId, TopicId, TypeSubscription
# 1. 向运行时注册"类型化订阅":github_issues 类消息一律交给 triage_agent
await runtime.add_subscription(
TypeSubscription(topic_type="github_issues", agent_type="triage_agent")
)
# 2. 此后,向任何以 github_issues 为 type 的 Topic 广播
topic = TopicId(type="github_issues", source="github.com/microsoft/autogen/issues/1")
await runtime.publish_message(message, topic_id=topic, sender=AgentId("poller_agent", "default"))
publish_message 是无响应的广播(不期待任何接收方回复),其可传关键字参数还包括 sender、cancellation_token 与 message_id。运行时会遍历所有匹配订阅,把消息逐一送入由 TypeSubscription 推导出的目标 Agent。当业务不再需要某条订阅时,用注册时记录下来的订阅 id(默认为自动生成的 UUID)调用 runtime.remove_subscription(id) 即可摘除该路由。
上述四种场景可以归纳为一张速查表,方便在不同需求下快速选型:
| 场景 | 订阅策略 | 发布使用的 Topic | 实际投递到的 Agent ID |
|---|---|---|---|
| 单租户 + 单 Topic | 每种 agent type 各一订阅,共用一个 topic type | 固定,如 ("default", "default") |
各 agent type 的 (agent_type, "default") |
| 单租户 + 多 Topic | 各 agent type 用不同 topic type(可多对一共享) | 按类别切换 topic type,source 固定 | (agent_type, "default") |
| 多租户 | 每个要并发的 agent type 一条订阅 | topic type 固定,source 取租户唯一标识 | (agent_type, <tenant/数据标识>) |
| 多租户 + 多 Topic | 每租户按需注册多个 topic type 的订阅 | 同时切换 topic type 与租户 source | (agent_type, <tenant/数据标识>) |
进一步阅读
- 关联文档原文:Topic and Subscription(core-concepts)
- Agent ID 的构成与生命周期:agent-identity-and-lifecycle.md
- 广播/直连、Agent 与运行时如何组成完整应用栈:application-stack.md
- Python 端核心实现:TopicId、Subscription 协议、TypeSubscription、SubscriptionManager、AgentRuntime 接口
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0631
MiniCPM5-2BMiniCPM5-2B 是一款面向端侧、本地部署和资源受限场景的 2B 稠密 Transformer,能够达到同尺寸开源模型 SOTA 水平。Markdown00
video-shotcraftAI宣传片skill,使用 Remotion 制作电影级产品视频:提供106 张镜头配方卡和可复用的视频魔板。适用于 Claude Code 与 Codex以及所有其他智能体Markdown00
HivisionIDPhotos⚡️HivisionIDPhotos: a lightweight and efficient AI ID photos tools. 一个轻量级的AI证件照制作算法。Python09
DragonOSDragonOS is an operating system developed from scratch using Rust, with Linux compatibility. It is designed for **Serverless** scenarios. 使用Rust从0自研内核,具有Linux兼容性的操作系统,面向云计算Serverless场景而设计。Rust00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00