Протоколы связи (шины сообщений)
Разделяйте агентов с помощью шины сообщений (Redis, NATS, Kafka), чтобы они могли независимо масштабироваться и обрабатывать сбои.
«Протоколы связи (шины сообщений)» — бесплатный урок AI Agents на CoddyKit. Это урок 4 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения 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Идентификаторы корреляции
Каждое связанное сообщение содержит один и тот же correlation_id, чтобы вы могли отслеживать задачу между сервисами:
record = {'correlation_id': 'abc', 'task_id': 'def', 'span': 'llm_call'}
print(record)
Идемпотентные обработчики
Сообщения могут доставляться более одного раза — семантика доставки минимум один раз. Делайте обработчики идемпотентными: одна и та же задача с task_id обрабатывается один раз.
Очередь недоставленных сообщений
Если обработчик неоднократно завершается с ошибкой, передайте сообщение в DLQ на проверку человеком вместо бесконечного повторения:
if attempts > MAX_RETRIES:
r.rpush('queue:research:dlq', raw)
log.error('Sent to DLQ', extra={'task_id': task_id})Наблюдаемость
Отслеживайте каждое сообщение:
- Идентификаторы трассировки проходят через шину
- Журналы содержат task_id и correlation_id
- Метрики: пропускная способность, задержка и частота ошибок для каждой темы
Когда переходить
Начните с простого — агентов внутри процесса. Переходите к шине только при наличии следующих условий:
- 5 и более агентов в рабочей среде
- Необходимость независимо масштабировать агентов
- Рабочие процессы, охватывающие несколько сервисов
Доставка минимум один раз?
Что требует от ваших обработчиков «доставка минимум один раз»?
Итоги
Для простых случаев используйте Redis, для низкой задержки — NATS, для надёжного хранения — Kafka. Схемы, идентификаторы корреляции, идемпотентные обработчики, DLQ.
Часто задаваемые вопросы
Урок «Протоколы связи (шины сообщений)» бесплатный?
Да — полный текст урока «Протоколы связи (шины сообщений)» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс AI Agents, подпишись на CoddyKit PRO. Курс AI Agents содержит 4 уроков всего.
Чему я научусь в уроке «Протоколы связи (шины сообщений)»?
Разделяйте агентов с помощью шины сообщений (Redis, NATS, Kafka), чтобы они могли независимо масштабироваться и обрабатывать сбои. Ты практикуешь AI Agents с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.
Нужен ли мне опыт, чтобы начать AI Agents?
Предыдущий опыт не требуется. AI Agents на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 4 из 4.
Сколько времени занимает урок «Протоколы связи (шины сообщений)»?
Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.
Можно ли писать и запускать код в этом уроке AI Agents?
Да. Каждый урок AI Agents включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.
Все уроки этого курса
- Мультиагентная система на основе диалога (AutoGen)
- Иерархические супервизоры (координатор + исполнители)
- Роли и специализации агентов
- Протоколы связи (шины сообщений)