AI Agents · 课时

通信协议(消息总线)

使用消息总线(Redis、NATS、Kafka)解耦智能体,使它们能够独立扩展并独立处理故障。

第 4 / 4 课15 个步骤

通信协议(消息总线) 是 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 反馈 — 无需本地设置。

此课程中的所有课时

  1. 基于对话的多智能体(AutoGen)
  2. 分层监督者(编排器 + 工作器)
  3. 智能体角色与专门化
  4. 通信协议(消息总线)
← 返回 AI Agents