Agent OS 消息总线适配器实战指南:用 Redis、Kafka、NATS 与云消息服务打通多 Agent 通信
【免费下载链接】agent-governance-toolkitAI Agent Governance Toolkit — Policy enforcement, zero-trust identity, execution sandboxing, and reliability engineering for autonomous AI agents. Covers 10/10 OWASP Agentic Top 10.项目地址: https://gitcode.com/GitHub_Trending/ag/agent-governance-toolkit
本指南面向需要在 AI Agent 之间建立解耦通信的开发者,系统讲解 Agent OS 中 Agent Message Bus(AMB)的六种内置消息代理适配器(Redis、Kafka、RabbitMQ、NATS、Azure Service Bus、AWS SQS)的安装、配置、核心 API 与常见消息模式,并结合仓库源码与测试用例说明其底层实现原理,帮助你按场景选型并落地一套可运行的多代理消息架构。
AMB 与 Broker Adapter:解耦 Agent 通信的抽象层
Agent OS 提供 Agent Message Bus(AMB,Agent 消息总线),用于在多个 AI Agent 之间进行解耦通信。与"点对点直连"不同,AMB 让消息的发送方与接收方互不感知:发送方只向某个topic发布消息,接收方按topic订阅感兴趣的消息,两者之间由 broker 中转。这套设计允许 Agent 发射信号(signal)、广播意图(intention),而无需在 Agent 之间建立紧耦合的调用链。
AMB 的核心是"broker 无关"(broker-agnostic)的抽象层。仓库中 amb_core/broker.py 定义了BrokerAdapter抽象基类,它约定了一套统一的异步接口:
connect()/disconnect():建立与释放 broker 连接;publish(message, wait_for_confirmation=False):发布消息,支持"fire and forget"与"等待确认"两种模式;subscribe(topic, handler)/unsubscribe(subscription_id):订阅与退订;request(message, timeout=30.0):请求-响应模式,等待对方回复;get_pending_messages(topic, limit):获取积压消息(各 broker 可选实现);get_backpressure_stats()/get_queue_size():背压与队列监控的可选接口。
所有内置适配器(Redis、Kafka、RabbitMQ、NATS、Azure Service Bus、AWS SQS)都实现这套接口,因此上层MessageBus的调用方式完全一致,切换 broker 只需更换传入的适配器实例。从 bus.py 可见,MessageBus.__init__默认使用InMemoryBroker(进程内实现,适合测试与开发),传入任意BrokerAdapter即可切换到真实中间件。
可用适配器一览
| 适配器 | 适用场景 | 安装命令 |
|---|---|---|
| Memory(内置) | 测试、开发、单进程应用 | 已随amb_core包含 |
| Redis | 低延迟消息、pub/sub、实时更新 | pip install agentmesh-message-bus[redis] |
| Kafka | 高吞吐、事件溯源、审计日志 | pip install agentmesh-message-bus[kafka] |
| RabbitMQ | 复杂路由、企业级消息 | pip install agentmesh-message-bus[rabbitmq] |
| NATS | 云原生、轻量级、边缘计算 | pip install agentmesh-message-bus[nats] |
| Azure Service Bus | Azure 生态、企业级消息 | pip install agentmesh-message-bus[azure] |
| AWS SQS | AWS 生态、Serverless | pip install agentmesh-message-bus[aws] |
说明:仓库 modules/amb/pyproject.toml 中该包名为
agent-governance-toolkit-message-bus,extra 依赖同样覆盖redis、rabbitmq、kafka等选项。本教程沿用了官方文档中的安装写法,实际安装时可结合你使用的发布渠道选择对应包名。
各适配器依赖的底层驱动如下(与pyproject.toml的 optional-dependencies 对应):
- Redis:
redis>=4.0.0,<5.0(异步客户端redis.asyncio); - Kafka:
aiokafka>=0.8.0,<1.0; - RabbitMQ:
aio-pika>=9.0.0,<10.0; - NATS:
nats-py; - Azure Service Bus:
azure-servicebus; - AWS SQS:
aioboto3。
若未安装对应依赖,适配器在导入时会抛出明确的 ImportError 提示安装命令,见 redis_broker.py。
Redis 快速上手:三步跑通第一个消息流
Redis 适配器同时利用了两套机制:Pub/Sub 做实时消息分发,Streams 做消息留存与请求-响应(见 redis_broker.py)。
1. 安装
pip install agentmesh-message-bus[redis]2. 启动 Redis
docker run -d -p 6379:6379 redis:latest3. 在 Agent 中使用
from amb_core.adapters import RedisBroker from amb_core import AgentMessageBus, Message # 创建 broker broker = RedisBroker(url="redis://localhost:6379/0") # 创建消息总线 bus = AgentMessageBus(broker=broker) # 连接 await bus.connect() # 订阅消息 async def handle_task(msg: Message): print(f"Received task: {msg.payload}") # 处理并回复 result = await process_task(msg.payload) # 发送响应 await bus.publish(Message( topic="results", payload=result, correlation_id=msg.correlation_id )) await bus.subscribe("tasks", handle_task) # 发布消息 await bus.publish(Message( topic="tasks", payload={"action": "analyze", "file": "data.txt"} ))从源码实现看,RedisBroker.connect()通过aioredis.from_url(url)建立客户端并取得 pub/sub 对象;publish()将消息序列化为 JSON 发布到对应 channel,同时用xadd写入stream:{topic}(保留最近 1000 条,见 redis_broker.py);subscribe()为每个 topic 启动一个后台监听任务,轮询get_message()并把消息反序列化为Message后交给 handler(redis_broker.py)。因此 Redis 模式天然支持"发布后立刻被多个订阅者收到"的实时广播,同时 Streams 为 request-response 提供了可靠的回执通道。
适配器对比与选型
Redis Adapter
最适合:低延迟消息、pub/sub 模式、实时更新
from amb_core.adapters import RedisBroker broker = RedisBroker( url="redis://localhost:6379/0" ) # 特性: # - Pub/sub 实时消息 # - Streams 消息持久化 # - 原生请求-响应支持优点:
- 延迟极低(亚毫秒级)
- 部署简单
- 借助 Redis Streams 具备内置持久化
不足:
- 与 Kafka 相比持久性有限
- 默认单节点
实现细节:Redis 适配器的 request-response 通过response:{correlation_id}流实现——发布方轮询xread阻塞读取该流直到收到响应或超时(redis_broker.py)。
Kafka Adapter
最适合:高吞吐、事件溯源、审计日志
from amb_core.adapters import KafkaBroker broker = KafkaBroker( bootstrap_servers="localhost:9092" ) # 特性: # - 高吞吐 # - 持久化消息存储 # - Consumer group 负载均衡优点:
- 吞吐量最高
- 持久性保障强
- 支持消息重放
不足:
- 部署复杂
- 延迟高于 Redis
实现细节:Kafka 适配器基于aiokafka,publish()以message.id作为 key 发送到 topic,wait_for_confirmation=True时会等待 producer 的 ack future(kafka_broker.py);每个subscribe()会创建独立的 consumer group(amb-{subscription_id},auto_offset_reset='latest'),从而天然实现多个 worker 的负载均衡与故障恢复(kafka_broker.py)。request-response 使用临时响应 topicresponse.{correlation_id}完成。
RabbitMQ Adapter
最适合:复杂路由、企业级消息
from amb_core.adapters import RabbitMQBroker broker = RabbitMQBroker( url="amqp://localhost:5672" ) # 未传 url 时默认读取环境变量 RABBITMQ_URL, # 再回退到 amqp://localhost/优点:
- 支持 topic 通配符路由(
*与#) - 企业级 AMQP 特性(mandatory 确认、排他回调队列)
- 路由灵活
不足:
- 相比 NATS 更重
实现细节:RabbitMQ 适配器声明了持久化的 topic 交换机amb.topic,发布时以message.topic作为 routing key;订阅时每个订阅者创建amb.queue.{subscription_id}自动删除队列并绑定到交换机(rabbitmq_broker.py),因此天然支持events.user.*这类通配订阅。request-response 采用经典 RPC 模式:声明排他回调队列,按 correlation_id 匹配响应(rabbitmq_broker.py)。
NATS Adapter
最适合:云原生应用、微服务、边缘计算
from amb_core.adapters import NATSBroker broker = NATSBroker( servers=["nats://localhost:4222"], use_jetstream=True # 开启持久化 ) # 特性: # - 轻量(单二进制) # - 原生请求-回复 # - JetStream 持久化优点:
- 非常轻量
- 易于部署
- 内置 request-reply
不足:
- 生态比 Kafka/RabbitMQ 小
实现细节:NATS 适配器将 topic 中的/转换为.后挂在amb.前缀下;use_jetstream=True时自动创建AMB_STREAM流(覆盖amb.>主题,上限 10 万条消息 / 100MB,见 nats_broker.py),订阅走 durable consumer;关闭 JetStream 时退化为 Core NATS 的 fire-and-forget。request-response 直接复用 NATS 原生 request-reply 机制,效率最高(nats_broker.py)。
Azure Service Bus Adapter
最适合:Azure 生态、企业级消息
from amb_core.adapters import AzureServiceBusBroker broker = AzureServiceBusBroker( connection_string="Endpoint=sb://...", topic_name="agent-messages" ) # 可选参数 subscription_name,默认 "amb-subscription" # 特性: # - 死信队列 # - Sessions 保证顺序 # - Azure AD 集成优点:
- 托管服务
- 企业级特性
- Azure 集成
不足:
- Azure 锁定
- 规模成本
实现细节:Azure 适配器将 AMB 消息映射为ServiceBusMessage,topic 写入subject字段用于过滤,并把topic、source、target、message_type写入application_properties(azure_servicebus_broker.py);接收端按主题过滤,处理失败时调用dead_letter_message进入死信(azure_servicebus_broker.py)。get_pending_messages()使用peek_messages不消费地窥视队列积压。
AWS SQS Adapter
最适合:AWS 生态、Serverless
from amb_core.adapters import AWSSQSBroker broker = AWSSQSBroker( region_name="us-east-1", queue_name="agent-messages", use_fifo=True # 开启 FIFO 保证顺序 ) # 也可直接传 queue_url 复用已有队列; # 凭证缺省时走 AWS 标准环境变量链 # 特性: # - Serverless 弹性扩展 # - FIFO 队列保证顺序 # - 死信队列优点:
- Serverless,无需自建基础设施
- 自动扩缩容
- AWS 集成
不足:
- 延迟较高
- AWS 锁定
实现细节:SQS 适配器基于aioboto3,连接时自动get_queue_url或按queue_name创建队列;FIFO 模式下自动附加FifoQueue属性,并设置MessageGroupId=topic、MessageDeduplicationId=message.id(aws_sqs_broker.py)。订阅者通过长轮询(WaitTimeSeconds=20)批量拉取消息,处理成功后删除消息,失败则保留在队列等待重试(aws_sqs_broker.py)。
四种常见消息模式
AMB 的MessageBus在底层适配器之上提供了统一的模式封装,核心 API 见 bus.py:
模式 1:请求-响应(Request-Response)
from amb_core import Message # Agent A 发送请求 response = await bus.request( Message( topic="calculate", payload={"operation": "sum", "values": [1, 2, 3]} ), timeout=30.0 ) print(f"Result: {response.payload}") # Result: 6bus.request()会自动生成correlation_id,由各适配器按各自机制实现(Redis 响应流、NATS 原生 reply、Kafka/Azure/SQS 临时响应队列),超时默认 30 秒,超时抛出TimeoutError(bus.py)。
模式 2:发布/订阅(Pub/Sub)
# Agent A 订阅事件(支持通配符,视 broker 能力而定) await bus.subscribe("events.user.*", handle_user_event) # Agent B 发布事件 await bus.publish(Message( topic="events.user.created", payload={"user_id": "123", "email": "user@example.com"} ))通配符订阅在 RabbitMQ(topic exchange 的*/#)与 NATS(subject 层级)上原生支持;Redis 适配器则需订阅实际 channel 名。
模式 3:工作队列(Work Queue)
# 多个 worker 订阅同一队列 # 每条消息只投递给一个 worker async def worker(msg: Message): result = await process_work(msg.payload) await bus.publish(Message( topic="results", payload=result, correlation_id=msg.id )) # 启动多个 worker for i in range(4): await bus.subscribe("work-queue", worker, consumer_group=f"workers")在 Kafka 上,每个subscribe()使用独立 consumer group,多个 worker 订阅同一 topic 时由 broker 自动完成分区负载均衡,实现"每条消息只被消费一次"的工作队列语义。
模式 4:事件溯源(Event Sourcing)
# 将所有事件发布到 Kafka 以获得持久性 kafka_broker = KafkaBroker(bootstrap_servers="localhost:9092") bus = AgentMessageBus(broker=kafka_broker) # 所有 Agent 行为都成为事件 await bus.publish(Message( topic="agent.events", payload={ "event_type": "document_analyzed", "agent_id": "analyzer-001", "document_id": "doc-123", "result": analysis_result, "timestamp": datetime.now(timezone.utc).isoformat() } )) # 事件可重放用于调试/审计配合MessageBus的持久化能力(persistence=True,见 bus.py),可以调用bus.replay(topic, handler, from_timestamp, to_timestamp)按时间窗口重放历史消息(bus.py),是审计与故障排查的关键能力。
多 Broker 混合部署
不同消息负载对中间件的要求不同,AMB 允许在同一应用中同时使用多个总线实例,各司其职:
from amb_core import AgentMessageBus from amb_core.adapters import RedisBroker, KafkaBroker # 实时消息走快速通道 redis_bus = AgentMessageBus( broker=RedisBroker(url="redis://localhost:6379") ) # 事件/审计走持久通道 kafka_bus = AgentMessageBus( broker=KafkaBroker(bootstrap_servers="localhost:9092") ) @kernel.register async def my_agent(task: str): # 处理任务 result = await process(task) # 通过 Redis 快速返回响应 await redis_bus.publish(Message( topic="responses", payload=result )) # 通过 Kafka 持久化事件 await kafka_bus.publish(Message( topic="events", payload={"action": "task_completed", "result": result} ))这种"热路径走低延迟、冷路径走高持久"的混合架构,是实际生产系统中最常见的 AMB 用法。
Docker Compose 一键起全套中间件
# docker-compose.yml version: '3.8' services: redis: image: redis:7 ports: - "6379:6379" kafka: image: confluentinc/cp-kafka:latest ports: - "9092:9092" environment: KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 zookeeper: image: confluentinc/cp-zookeeper:latest environment: ZOOKEEPER_CLIENT_PORT: 2181 nats: image: nats:latest ports: - "4222:4222" command: ["--js"] # 开启 JetStream agent: build: . depends_on: - redis - kafka - nats environment: REDIS_URL: redis://redis:6379 KAFKA_SERVERS: kafka:9092 NATS_URL: nats://nats:4222启动后即可在本地同时体验 Redis、Kafka(含 Zookeeper)与 NATS 三种 broker。RabbitMQ、Azure Service Bus 与 AWS SQS 的接入方式类似:RabbitMQ 可另加rabbitmq:3服务映射 5672 端口;云服务则直接在适配器中传入连接字符串或区域参数。
生产实践:环境变量、重连、DLQ 与延迟监控
1. 使用环境变量管理配置
import os broker = RedisBroker( url=os.environ.get("REDIS_URL", "redis://localhost:6379") )将连接地址外置到环境变量,便于在不同环境(本地 / CI / 生产)间切换,也避免密钥硬编码进代码。SQS 与 Azure 适配器同样支持通过标准环境变量注入凭证。
2. 处理断线重连
async def with_reconnect(bus: AgentMessageBus): while True: try: await bus.connect() break except ConnectionError: print("Connection failed, retrying in 5s...") await asyncio.sleep(5)底层各适配器均以ConnectionError表达连接失败,可统一用上述重试循环兜底;MessageBus同时实现了异步上下文管理器(async with MessageBus() as bus:),可在进入/退出时自动 connect/disconnect(bus.py)。
3. 配置死信队列(DLQ)
# 为失败消息配置 DLQ broker = RedisBroker( url="redis://localhost:6379", dead_letter_queue="dlq:agent-messages" )除适配器层级的死信机制外,MessageBus还内置了更通用的 DLQ 能力:构造时传dlq_enabled=True,订阅处理器抛异常或消息过期时,消息会被包装为DLQEntry记录原因(HANDLER_ERROR/EXPIRED)与堆栈信息进入死信队列,而不会中断总线(bus.py),可通过bus.get_dlq_stats()查看积压统计。
4. 监控消费延迟
from amb_core.observability import metrics # 跟踪消息处理延迟 @metrics.track("message_processing") async def handle_message(msg: Message): lag = time.time() - msg.timestamp metrics.gauge("message_lag_seconds", lag) await process(msg)Message模型自带 UTC 时间戳与age_seconds/remaining_ttl属性(见 models.py),可直接计算消息从发布到被处理的端到端延迟,配合优先级、TTL 等字段定位积压问题。
消息模型与进阶能力:优先级、TTL 与分布式追踪
本教程示例中的Message仅是 AMB 消息模型的冰山一角。models.py 中的完整字段包括:
id/topic/payload:消息标识、主题与负载;priority:优先级,取值BACKGROUND(0)/LOW(1)/NORMAL(5)/HIGH(8)/URGENT(10)/CRITICAL(15),InMemoryBroker会用优先级堆保证 CRITICAL 消息优先于 BACKGROUND 消费(见 memory_broker.py);correlation_id/reply_to:请求-响应关联与回执地址;ttl_seconds:生存时间,过期消息(is_expired=True)在进入 handler 前会被丢弃并转入 DLQ;trace_id/span_id/parent_span_id:分布式追踪字段,MessageBus默认自动注入当前 trace 上下文(bus.py)。
此外amb_core还提供SchemaRegistry(按 topic 对 payload 做 Pydantic 校验)、FileMessageStore等持久化存储,以及 CloudEvents 封装(见init.py),可用于构建具备消息契约校验和标准事件格式的治理型消息架构。
验证与测试
仓库在 modules/amb/tests/test_bus.py 中提供了完整的MessageBus行为测试,覆盖连接/断开、异步上下文管理器、fire-and-forget 发布、带确认发布、订阅-发布收发等核心路径,例如test_publish_fire_and_forget与test_subscribe_and_publish。运行测试:
cd agent-governance-python/agent-os/modules/amb pip install -e ".[dev]" pytest配套的示例代码位于 modules/amb/examples/,包括basic_usage.py、advanced_features.py(持久化、DLQ、schema 校验、优先级、TTL)、tracing_demo.py与backpressure_demo.py,可直接作为接入 AMB 的参考起点。
总结
Agent OS 的 AMB 通过一套BrokerAdapter抽象屏蔽了底层中间件差异:Memory 适配器让你零依赖起步,Redis 覆盖低延迟实时场景,Kafka 承载高吞吐事件溯源,RabbitMQ 提供复杂路由,NATS 适配云原生轻量部署,Azure Service Bus 与 AWS SQS 则打通两大云生态。结合MessageBus层提供的 request-response、pub/sub、工作队列、事件溯源四类模式,以及优先级、TTL、DLQ、持久化重放、分布式追踪等治理能力,你可以在不改动业务代码的前提下,按负载特征自由组合 broker,构建一套解耦、可靠、可观测的多 Agent 通信底座。
【免费下载链接】agent-governance-toolkitAI Agent Governance Toolkit — Policy enforcement, zero-trust identity, execution sandboxing, and reliability engineering for autonomous AI agents. Covers 10/10 OWASP Agentic Top 10.项目地址: https://gitcode.com/GitHub_Trending/ag/agent-governance-toolkit
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考