资讯动态

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

发布时间:2026/9/18 9:04:54 来源:尧图企业网站定制
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 BusAMB的六种内置消息代理适配器Redis、Kafka、RabbitMQ、NATS、Azure Service Bus、AWS SQS的安装、配置、核心 API 与常见消息模式并结合仓库源码与测试用例说明其底层实现原理帮助你按场景选型并落地一套可运行的多代理消息架构。AMB 与 Broker Adapter解耦 Agent 通信的抽象层Agent OS 提供 Agent Message BusAMBAgent 消息总线用于在多个 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_confirmationFalse)发布消息支持fire and forget与等待确认两种模式subscribe(topic, handler)/unsubscribe(subscription_id)订阅与退订request(message, timeout30.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-busextra 依赖同样覆盖redis、rabbitmq、kafka等选项。本教程沿用了官方文档中的安装写法实际安装时可结合你使用的发布渠道选择对应包名。各适配器依赖的底层驱动如下与pyproject.toml的 optional-dependencies 对应Redisredis4.0.0,5.0异步客户端redis.asyncioKafkaaiokafka0.8.0,1.0RabbitMQaio-pika9.0.0,10.0NATSnats-pyAzure Service Busazure-servicebusAWS SQSaioboto3。若未安装对应依赖适配器在导入时会抛出明确的 ImportError 提示安装命令见 redis_broker.py。Redis 快速上手三步跑通第一个消息流Redis 适配器同时利用了两套机制Pub/Sub 做实时消息分发Streams 做消息留存与请求-响应见 redis_broker.py。1. 安装pip install agentmesh-message-bus[redis]2. 启动 Redisdocker run -d -p 6379:6379 redis:latest3. 在 Agent 中使用from amb_core.adapters import RedisBroker from amb_core import AgentMessageBus, Message # 创建 broker broker RedisBroker(urlredis://localhost:6379/0) # 创建消息总线 bus AgentMessageBus(brokerbroker) # 连接 await bus.connect() # 订阅消息 async def handle_task(msg: Message): print(fReceived task: {msg.payload}) # 处理并回复 result await process_task(msg.payload) # 发送响应 await bus.publish(Message( topicresults, payloadresult, correlation_idmsg.correlation_id )) await bus.subscribe(tasks, handle_task) # 发布消息 await bus.publish(Message( topictasks, payload{action: analyze, file: data.txt} ))从源码实现看RedisBroker.connect()通过aioredis.from_url(url)建立客户端并取得 pub/sub 对象publish()将消息序列化为 JSON 发布到对应 channel同时用xadd写入stream:{topic}保留最近 1000 条见 redis_broker.pysubscribe()为每个 topic 启动一个后台监听任务轮询get_message()并把消息反序列化为Message后交给 handlerredis_broker.py。因此 Redis 模式天然支持发布后立刻被多个订阅者收到的实时广播同时 Streams 为 request-response 提供了可靠的回执通道。适配器对比与选型Redis Adapter最适合低延迟消息、pub/sub 模式、实时更新from amb_core.adapters import RedisBroker broker RedisBroker( urlredis://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_serverslocalhost:9092 ) # 特性 # - 高吞吐 # - 持久化消息存储 # - Consumer group 负载均衡优点吞吐量最高持久性保障强支持消息重放不足部署复杂延迟高于 Redis实现细节Kafka 适配器基于aiokafkapublish()以message.id作为 key 发送到 topicwait_for_confirmationTrue时会等待 producer 的 ack futurekafka_broker.py每个subscribe()会创建独立的 consumer groupamb-{subscription_id}auto_offset_resetlatest从而天然实现多个 worker 的负载均衡与故障恢复kafka_broker.py。request-response 使用临时响应 topicresponse.{correlation_id}完成。RabbitMQ Adapter最适合复杂路由、企业级消息from amb_core.adapters import RabbitMQBroker broker RabbitMQBroker( urlamqp://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_jetstreamTrue # 开启持久化 ) # 特性 # - 轻量单二进制 # - 原生请求-回复 # - JetStream 持久化优点非常轻量易于部署内置 request-reply不足生态比 Kafka/RabbitMQ 小实现细节NATS 适配器将 topic 中的/转换为.后挂在amb.前缀下use_jetstreamTrue时自动创建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_stringEndpointsb://..., topic_nameagent-messages ) # 可选参数 subscription_name默认 amb-subscription # 特性 # - 死信队列 # - Sessions 保证顺序 # - Azure AD 集成优点托管服务企业级特性Azure 集成不足Azure 锁定规模成本实现细节Azure 适配器将 AMB 消息映射为ServiceBusMessagetopic 写入subject字段用于过滤并把topic、source、target、message_type写入application_propertiesazure_servicebus_broker.py接收端按主题过滤处理失败时调用dead_letter_message进入死信azure_servicebus_broker.py。get_pending_messages()使用peek_messages不消费地窥视队列积压。AWS SQS Adapter最适合AWS 生态、Serverlessfrom amb_core.adapters import AWSSQSBroker broker AWSSQSBroker( region_nameus-east-1, queue_nameagent-messages, use_fifoTrue # 开启 FIFO 保证顺序 ) # 也可直接传 queue_url 复用已有队列 # 凭证缺省时走 AWS 标准环境变量链 # 特性 # - Serverless 弹性扩展 # - FIFO 队列保证顺序 # - 死信队列优点Serverless无需自建基础设施自动扩缩容AWS 集成不足延迟较高AWS 锁定实现细节SQS 适配器基于aioboto3连接时自动get_queue_url或按queue_name创建队列FIFO 模式下自动附加FifoQueue属性并设置MessageGroupIdtopic、MessageDeduplicationIdmessage.idaws_sqs_broker.py。订阅者通过长轮询WaitTimeSeconds20批量拉取消息处理成功后删除消息失败则保留在队列等待重试aws_sqs_broker.py。四种常见消息模式AMB 的MessageBus在底层适配器之上提供了统一的模式封装核心 API 见 bus.py模式 1请求-响应Request-Responsefrom amb_core import Message # Agent A 发送请求 response await bus.request( Message( topiccalculate, payload{operation: sum, values: [1, 2, 3]} ), timeout30.0 ) print(fResult: {response.payload}) # Result: 6bus.request()会自动生成correlation_id由各适配器按各自机制实现Redis 响应流、NATS 原生 reply、Kafka/Azure/SQS 临时响应队列超时默认 30 秒超时抛出TimeoutErrorbus.py。模式 2发布/订阅Pub/Sub# Agent A 订阅事件支持通配符视 broker 能力而定 await bus.subscribe(events.user.*, handle_user_event) # Agent B 发布事件 await bus.publish(Message( topicevents.user.created, payload{user_id: 123, email: userexample.com} ))通配符订阅在 RabbitMQtopic exchange 的*/#与 NATSsubject 层级上原生支持Redis 适配器则需订阅实际 channel 名。模式 3工作队列Work Queue# 多个 worker 订阅同一队列 # 每条消息只投递给一个 worker async def worker(msg: Message): result await process_work(msg.payload) await bus.publish(Message( topicresults, payloadresult, correlation_idmsg.id )) # 启动多个 worker for i in range(4): await bus.subscribe(work-queue, worker, consumer_groupfworkers)在 Kafka 上每个subscribe()使用独立 consumer group多个 worker 订阅同一 topic 时由 broker 自动完成分区负载均衡实现每条消息只被消费一次的工作队列语义。模式 4事件溯源Event Sourcing# 将所有事件发布到 Kafka 以获得持久性 kafka_broker KafkaBroker(bootstrap_serverslocalhost:9092) bus AgentMessageBus(brokerkafka_broker) # 所有 Agent 行为都成为事件 await bus.publish(Message( topicagent.events, payload{ event_type: document_analyzed, agent_id: analyzer-001, document_id: doc-123, result: analysis_result, timestamp: datetime.now(timezone.utc).isoformat() } )) # 事件可重放用于调试/审计配合MessageBus的持久化能力persistenceTrue见 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( brokerRedisBroker(urlredis://localhost:6379) ) # 事件/审计走持久通道 kafka_bus AgentMessageBus( brokerKafkaBroker(bootstrap_serverslocalhost:9092) ) kernel.register async def my_agent(task: str): # 处理任务 result await process(task) # 通过 Redis 快速返回响应 await redis_bus.publish(Message( topicresponses, payloadresult )) # 通过 Kafka 持久化事件 await kafka_bus.publish(Message( topicevents, 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( urlos.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/disconnectbus.py。3. 配置死信队列DLQ# 为失败消息配置 DLQ broker RedisBroker( urlredis://localhost:6379, dead_letter_queuedlq:agent-messages )除适配器层级的死信机制外MessageBus还内置了更通用的 DLQ 能力构造时传dlq_enabledTrue订阅处理器抛异常或消息过期时消息会被包装为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.pycorrelation_id/reply_to请求-响应关联与回执地址ttl_seconds生存时间过期消息is_expiredTrue在进入 handler 前会被丢弃并转入 DLQtrace_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),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价