news 2026/9/18 9:04:48

Agent OS 消息总线适配器实战指南:用 Redis、Kafka、NATS 与云消息服务打通多 Agent 通信

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Agent OS 消息总线适配器实战指南:用 Redis、Kafka、NATS 与云消息服务打通多 Agent 通信

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 BusAzure 生态、企业级消息pip install agentmesh-message-bus[azure]
AWS SQSAWS 生态、Serverlesspip install agentmesh-message-bus[aws]

说明:仓库 modules/amb/pyproject.toml 中该包名为agent-governance-toolkit-message-bus,extra 依赖同样覆盖redisrabbitmqkafka等选项。本教程沿用了官方文档中的安装写法,实际安装时可结合你使用的发布渠道选择对应包名。

各适配器依赖的底层驱动如下(与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:latest

3. 在 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 适配器基于aiokafkapublish()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字段用于过滤,并把topicsourcetargetmessage_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=topicMessageDeduplicationId=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: 6

bus.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_forgettest_subscribe_and_publish。运行测试:

cd agent-governance-python/agent-os/modules/amb pip install -e ".[dev]" pytest

配套的示例代码位于 modules/amb/examples/,包括basic_usage.pyadvanced_features.py(持久化、DLQ、schema 校验、优先级、TTL)、tracing_demo.pybackpressure_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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/18 9:02:53

音乐数据分析:从用户画像到精准运营

1. 项目背景与核心价值音乐产业正经历着前所未有的数字化转型浪潮。根据国际唱片业协会(IFPI)最新报告&#xff0c;全球数字音乐收入已突破300亿美元大关&#xff0c;其中流媒体收入占比超过67%。在这个背景下&#xff0c;各大音乐平台和版权方都面临着一个共同的挑战&#xff…

作者头像 李华
网站建设 2026/9/18 9:01:39

如何5分钟安装Kap并完成第一次录屏:新手快速上手指南

如何5分钟安装Kap并完成第一次录屏&#xff1a;新手快速上手指南 【免费下载链接】Kap An open-source screen recorder built with web technology 项目地址: https://gitcode.com/gh_mirrors/ka/Kap Kap 是一款免费开源的 Mac 屏幕录制工具&#xff08;screen recorde…

作者头像 李华
网站建设 2026/9/18 9:00:48

2026年AI写作工具市场现状与专业评测

1. 2026年AI写作工具市场现状2026年的AI写作领域已经进入成熟期&#xff0c;各类工具在细分场景的应用呈现出明显的差异化特征。根据最新行业调研数据&#xff0c;全球AI写作工具市场规模已达到320亿美元&#xff0c;年复合增长率保持在28%左右。这个快速增长的市场背后&#x…

作者头像 李华
网站建设 2026/9/18 9:00:46

YuE2:基于Hugging Face的AR-NAR混合Transformer生成框架

1. 项目概述&#xff1a;从“YuE”到AR–NAR混合架构的落地实践最近在Hugging Face上频繁看到“YuE”和“YuE2”这两个词&#xff0c;尤其在文本生成、语音合成和多模态建模相关的Spaces和Model Cards里反复出现。起初我以为是某个新出的开源模型缩写&#xff0c;查了一圈才发现…

作者头像 李华
网站建设 2026/9/18 8:58:43

数据库原理及应用复习指南:关系模型、SQL与事务核心解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华