通信协议(消息总线)
使用消息总线(Redis、NATS、Kafka)解耦智能体,使它们能够独立扩展并独立处理故障。
通信协议(消息总线) 是 CoddyKit 上的免费 AI Agents 课时。 这是第 4 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 AI Agents 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 AI Agents 课程共包含 4 节课。
本课时的部分内容尚未翻译,以英文显示。
超越进程内对话
对于小型系统,智能体会在进程内相互调用。对于大型系统,则可以通过消息总线将它们解耦,例如 Redis、NATS 或 Kafka。
为什么要解耦
- 智能体可以独立扩展
- 每个智能体的故障彼此隔离
- 支持异步工作流(无需在内联调用中等待缓慢的智能体)
- 可以重放消息以便调试
- 支持跨语言:Python 智能体和 Go 智能体都可以订阅
Simple Bus: Redis Pub/Sub
import redis
r = redis.Redis()
# Publisher
r.publish('agent.research.task', json.dumps({'task_id': 'abc', 'query': '...'}))
# Subscriber
p = r.pubsub()
p.subscribe('agent.research.task')
for msg in p.listen():
if msg['type'] == 'message':
handle_task(json.loads(msg['data']))用于广播的发布/订阅
当许多智能体都可能需要了解某个事件时使用(例如“user-question-received”)。
用于工作分发的队列
当每项任务应当恰好由一个工作进程处理时,请使用队列(Redis BLPOP、RabbitMQ):
import json, queue
r = queue.Queue()
# Producer
r.put(json.dumps({'task_id': 'abc', 'kind': 'research'}))
def process(task):
print('processing', task)
# Worker
while not r.empty():
raw = r.get()
task = json.loads(raw)
process(task)
NATS:速度优先
NATS 是一种为微服务设计的轻量级消息代理。它提供亚毫秒级延迟,并内置请求/回复功能:
import nats
nc = await nats.connect('nats://localhost:4222')
await nc.publish('agent.research', json.dumps(task).encode())
# Request/reply
response = await nc.request('agent.research', payload, timeout=10)Kafka:持久性优先
Kafka 增加了持久且可重放的日志,非常适合审计和重放:
from confluent_kafka import Producer, Consumer
producer.produce('agent-events', key=task_id, value=json.dumps(task))
producer.flush()事件模式
请使用 Pydantic 或 Protobuf 定义每种消息结构。没有模式,您的总线就会变得一团糟:
class ResearchTask(BaseModel):
task_id: str
user_id: str
query: str
deadline: datetime
class ResearchResult(BaseModel):
task_id: str
findings: list[str]
duration_ms: int关联 ID
每条相关消息都携带相同的 correlation_id,这样您就可以跨服务跟踪任务:
record = {'correlation_id': 'abc', 'task_id': 'def', 'span': 'llm_call'}
print(record)
幂等处理程序
消息可能会被投递多次(至少一次语义)。请让处理程序具备幂等性——相同的 task_id 只处理一次。
死信队列
当处理程序反复失败时,请将消息 dump 到 DLQ,交由人工审查,而不是让它无限循环:
if attempts > MAX_RETRIES:
r.rpush('queue:research:dlq', raw)
log.error('Sent to DLQ', extra={'task_id': task_id})可观测性
请跟踪每条消息:
- 跟踪 ID 在总线中流转
- 日志包含 task_id 和 correlation_id
- 指标包括每个主题的吞吐量、延迟和错误率
何时采用
请从简单方案开始——使用进程内智能体。只有在出现以下情况时,才迁移到消息总线:
- 生产环境中有 5 个或更多智能体
- 需要独立扩展智能体
- 工作流跨越多个服务
至少一次?
“至少一次投递”要求您的处理程序做到什么?
回顾
简单场景使用 Redis 解耦,重视延迟时使用 NATS,重视持久性时使用 Kafka。还需要消息模式、关联 ID、幂等处理程序和 DLQ。
用 AI 导师学习 AI Agents — 免费
在浏览器中编写并运行真实代码,获得全天候 AI 导师的即时帮助,并在网页或应用中继续学习。
- 课程
- 60
- 课程
- 239
常见问题解答
「通信协议(消息总线)」课时是免费的吗?
是的 — 「通信协议(消息总线)」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 AI Agents 课程的其余内容,请升级到 CoddyKit PRO。 AI Agents 课程共包含 4 节课。
「通信协议(消息总线)」这节课中我会学到什么?
使用消息总线(Redis、NATS、Kafka)解耦智能体,使它们能够独立扩展并独立处理故障。 你通过在浏览器中直接运行的动手代码来练习 AI Agents,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 AI Agents 需要有经验吗?
无需任何先前经验。CoddyKit 上的 AI Agents 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 4 节课,共 4 节。
「通信协议(消息总线)」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 AI Agents 课中编写并运行代码吗?
能。每节 AI Agents 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。