Очереди событий и брокеры сообщений
Очереди Redis, RabbitMQ и Kafka для раздельного взаимодействия агентов.
«Очереди событий и брокеры сообщений» — бесплатный урок AI Agents на CoddyKit. Это урок 2 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения AI Agents, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс AI Agents содержит 4 уроков всего.
Зачем агентам нужны очереди сообщений?
Очереди сообщений разделяют инициатор (производитель событий) и агента (потребителя событий). Производитель отправляет события с любой скоростью, а агент обрабатывает их в собственном темпе. Это предотвращает перегрузку и позволяет повторять попытки.
Списки Redis как очереди
Списки Redis работают как простые очереди: LPUSH добавляет элемент в начало (постановка в очередь), а BRPOP удаляет элемент с конца и блокируется, если список пуст (извлечение из очереди). Это надёжная очередь FIFO.
import redis
r = redis.Redis(host='localhost', port=6379, decode_responses=True)
QUEUE_NAME = 'agent:tasks'
# Producer: add a task to the queue
def enqueue_task(task: dict):
import json
r.lpush(QUEUE_NAME, json.dumps(task))
print(f'Enqueued task: {task["id"]}')
# Consumer: blocking pop (waits up to 5 seconds for a message)
def dequeue_task(timeout: int = 5):
import json
result = r.brpop(QUEUE_NAME, timeout=timeout)
if result:
queue_name, raw_data = result
return json.loads(raw_data)
return None # Timed out - no tasks
# Producer side
enqueue_task({'id': 'task-1', 'type': 'email', 'email_id': 'email-abc'})
enqueue_task({'id': 'task-2', 'type': 'email', 'email_id': 'email-def'})
# Check queue length
print(f'Queue length: {r.llen(QUEUE_NAME)}')Цикл работы агента с Redis
Рабочий процесс агента непрерывно проверяет очередь на наличие новых задач, обрабатывает каждую из них и повторяет цикл. Если в течение времени ожидания задачи не поступают, он снова выполняет цикл, не расходуя ресурсы.
import redis
import json
import time
import logging
logger = logging.getLogger('worker')
r = redis.Redis(host='localhost', port=6379, decode_responses=True)
def process_task(task: dict) -> bool:
task_type = task.get('type')
if task_type == 'email':
logger.info(f'Processing email: {task.get("email_id")}')
# Call email agent here
return True
elif task_type == 'file':
logger.info(f'Processing file: {task.get("filepath")}')
return True
else:
logger.warning(f'Unknown task type: {task_type}')
return False
def run_worker(queue_name: str = 'agent:tasks'):
print(f'Worker started, watching queue: {queue_name}')
while True:
try:
task = dequeue_task(timeout=5)
if task:
success = process_task(task)
if not success:
# Re-queue failed tasks for retry
r.lpush('agent:failed', json.dumps(task))
else:
logger.debug('No tasks, waiting...')
except KeyboardInterrupt:
print('Worker stopped')
break
except Exception as e:
logger.error(f'Worker error: {e}')
time.sleep(1) # Brief pause on unexpected error
print('Worker function defined')Redis Pub/Sub для событий
Redis Pub/Sub отличается от списков: он рассылает сообщения всем текущим подписчикам. Используйте Pub/Sub для уведомлений с веерной рассылкой, когда нескольким агентам нужно отреагировать на одно и то же событие.
import redis
import json
import threading
r = redis.Redis(host='localhost', port=6379, decode_responses=True)
# Publisher: send event to all subscribers
def publish_event(channel: str, event: dict):
r.publish(channel, json.dumps(event))
print(f'Published to {channel}: {event}')
# Subscriber: listen for events in a background thread
def subscribe_and_handle(channel: str, handler_fn):
pubsub = r.pubsub()
pubsub.subscribe(channel)
def listener():
for message in pubsub.listen():
if message['type'] == 'message':
event = json.loads(message['data'])
handler_fn(event)
thread = threading.Thread(target=listener, daemon=True)
thread.start()
return thread
def handle_agent_event(event):
print(f'Agent received event: {event}')
# Start subscriber
subscribe_and_handle('agent:events', handle_agent_event)
# Publish an event
publish_event('agent:events', {'type': 'user_action', 'action': 'login', 'user_id': 42})
import time
time.sleep(0.1) # Give subscriber time to receiveОсновы RabbitMQ с pika
RabbitMQ — это полнофункциональный брокер сообщений с маршрутизацией, подтверждениями и очередями недоставленных сообщений. Библиотека pika подключает Python к RabbitMQ.
import pika
import json
# Connect to RabbitMQ
connection = pika.BlockingConnection(
pika.ConnectionParameters(
host='localhost',
port=5672,
credentials=pika.PlainCredentials('guest', 'guest')
)
)
channel = connection.channel()
# Declare a durable queue (survives broker restart)
channel.queue_declare(
queue='agent_tasks',
durable=True # Queue survives RabbitMQ restart
)
# Publish a message
def publish_task(task: dict):
channel.basic_publish(
exchange='',
routing_key='agent_tasks',
body=json.dumps(task),
properties=pika.BasicProperties(
delivery_mode=2, # Make message persistent
content_type='application/json'
)
)
print(f'Published task: {task["id"]}')
publish_task({'id': 'task-1', 'type': 'process_email', 'email_id': 'email-xyz'})
connection.close()Потребитель RabbitMQ с подтверждениями
Потребители должны подтверждать сообщения после обработки. Если агент завершится до подтверждения, RabbitMQ снова поместит сообщение в очередь для другого потребителя. Это гарантирует, что сообщения не будут потеряны.
import pika
import json
def create_consumer():
connection = pika.BlockingConnection(
pika.ConnectionParameters(host='localhost')
)
channel = connection.channel()
channel.queue_declare(queue='agent_tasks', durable=True)
# Process one message at a time (fair dispatch)
channel.basic_qos(prefetch_count=1)
def process_message(ch, method, properties, body):
task = json.loads(body)
print(f'Processing: {task["id"]}')
try:
# Do agent work here
print(f'Completed task: {task["id"]}')
# Acknowledge: message is removed from queue
ch.basic_ack(delivery_tag=method.delivery_tag)
except Exception as e:
print(f'Failed task {task["id"]}: {e}')
# Negative acknowledge + requeue=True: puts message back
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True)
channel.basic_consume(queue='agent_tasks', on_message_callback=process_message)
print('Consumer ready. Waiting for messages...')
channel.start_consuming()
print('RabbitMQ consumer defined')Очереди недоставленных сообщений
Сообщения, обработка которых завершилась ошибкой слишком много раз, следует отправлять в очередь недоставленных сообщений (DLQ) для ручной проверки, а не повторять попытки бесконечно. Настройте RabbitMQ так, чтобы он автоматически направлял неудачные сообщения в DLQ.
import pika
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
# Create the dead letter queue first
channel.queue_declare(queue='agent_tasks_dlq', durable=True)
# Create main queue with dead-letter exchange config
channel.queue_declare(
queue='agent_tasks',
durable=True,
arguments={
'x-dead-letter-exchange': '',
'x-dead-letter-routing-key': 'agent_tasks_dlq',
'x-message-ttl': 3600000, # Messages expire after 1 hour
'x-max-delivery-count': 3 # Max 3 delivery attempts
}
)
# Consumer: nack without requeue sends to DLQ
def careful_consumer(ch, method, properties, body):
import json
task = json.loads(body)
try:
# Process task
ch.basic_ack(delivery_tag=method.delivery_tag)
except Exception:
# Do NOT requeue - send to DLQ
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)
print('Dead letter queue configured')
connection.close()Шаблон очереди задач
Шаблон очереди задач разделяет инициирование и выполнение: простой производитель отправляет описания задач в очередь, а рабочие процессы извлекают их и выполняют. Несколько рабочих процессов могут обрабатывать задачи параллельно.
import redis
import json
from datetime import datetime
r = redis.Redis(host='localhost', port=6379, decode_responses=True)
# Task descriptor - what needs to be done
def create_agent_task(task_type: str, params: dict, priority: int = 5) -> dict:
return {
'id': f'{task_type}-{int(datetime.utcnow().timestamp() * 1000)}',
'type': task_type,
'params': params,
'priority': priority,
'created_at': datetime.utcnow().isoformat(),
'retry_count': 0,
'max_retries': 3
}
# Enqueue with priority (multiple queues by priority)
def enqueue_with_priority(task: dict):
priority = task.get('priority', 5)
queue = f'agent:tasks:p{priority}'
r.lpush(queue, json.dumps(task))
print(f'Enqueued {task["id"]} to priority-{priority} queue')
# Worker dequeues from high-priority queue first
def dequeue_priority(timeout: int = 5):
for p in [1, 2, 3, 4, 5]: # Check priority 1 first
result = r.brpop(f'agent:tasks:p{p}', timeout=0.1)
if result:
return json.loads(result[1])
return None
enqueue_with_priority(create_agent_task('email_analysis', {'email_id': 'abc'}, priority=2))
enqueue_with_priority(create_agent_task('file_process', {'path': '/tmp/file.txt'}, priority=5))Celery: высокоуровневая очередь задач
Celery — самая популярная очередь задач для Python. Она поддерживает Redis и RabbitMQ в качестве брокеров, автоматические повторные попытки, маршрутизацию задач и мониторинг с помощью Flower. Используйте её для рабочих нагрузок агентов в production.
from celery import Celery
import os
# Create Celery app with Redis broker
app = Celery(
'agent_tasks',
broker=os.environ.get('REDIS_URL', 'redis://localhost:6379/0'),
backend=os.environ.get('REDIS_URL', 'redis://localhost:6379/0')
)
# Configure retry behavior
app.conf.update(
task_acks_late=True,
task_reject_on_worker_lost=True,
task_serializer='json',
result_expires=3600
)
@app.task(bind=True, max_retries=3, default_retry_delay=60)
def run_email_agent(self, email_id: str):
try:
# Agent logic here
print(f'Processing email: {email_id}')
return {'status': 'success', 'email_id': email_id}
except Exception as exc:
raise self.retry(exc=exc)
# Trigger a task (from any Python code)
# run_email_agent.delay('email-abc') # Fire and forget
# result = run_email_agent.apply_async(args=['email-abc'], countdown=60) # Delayed
print('Celery task defined')Мониторинг состояния очереди
Контролируйте глубину очереди, чтобы обнаруживать накопившиеся задачи. Если глубина очереди растёт, вам нужно больше рабочих процессов или агент работает слишком медленно. Настройте оповещение при превышении глубиной заданного порога.
import redis
r = redis.Redis(host='localhost', port=6379, decode_responses=True)
def check_queue_health(queue_name: str, max_depth: int = 100) -> dict:
depth = r.llen(queue_name)
oldest_raw = r.lindex(queue_name, -1) # Get last item (oldest)
oldest_age_seconds = None
if oldest_raw:
import json
from datetime import datetime
oldest = json.loads(oldest_raw)
if 'created_at' in oldest:
created = datetime.fromisoformat(oldest['created_at'])
oldest_age_seconds = (datetime.utcnow() - created).total_seconds()
health = {
'queue': queue_name,
'depth': depth,
'max_depth': max_depth,
'overloaded': depth > max_depth,
'oldest_item_age_seconds': oldest_age_seconds
}
if health['overloaded']:
print(f'ALERT: Queue {queue_name} depth={depth} exceeds max={max_depth}')
return health
print('Queue health monitor defined')
print('Usage: check_queue_health("agent:tasks")')Выбор между Redis и RabbitMQ
Используйте Redis, если вам нужна простота, вы уже используете Redis или задачи выполняются недолго и потеря нескольких сообщений при сбое допустима. Используйте RabbitMQ, если вам нужны гарантированная доставка, сложная маршрутизация, очереди недоставленных сообщений или несколько потребителей с разными подписками.
# Decision guide as code comments
# Use Redis lists when:
# - Simple FIFO queue is enough
# - Redis already in stack
# - Can tolerate rare message loss on Redis crash
# - Low operational overhead matters
# Use RabbitMQ when:
# - Message acknowledgment is critical (no data loss)
# - Need complex routing (topic exchanges, fanout)
# - Need dead-letter queues for failed messages
# - Multiple consumer types need different message subsets
# - Need message TTL and expiry
# Use Celery (on top of either) when:
# - Need task scheduling and delays
# - Need automatic retries with exponential backoff
# - Need task result storage
# - Need monitoring dashboard (Flower)
print('Redis: simple, fast, ok with rare loss')
print('RabbitMQ: guaranteed delivery, complex routing')
print('Celery: high-level abstraction over both')Проверка знаний: очереди сообщений
Проверьте своё понимание очередей сообщений и брокеров для агентов.
Итоги по очередям сообщений
Очереди сообщений отделяют инициирование агентом от выполнения: списки Redis предоставляют простые очереди FIFO с LPUSH/BRPOP; Redis Pub/Sub позволяет распространять события; RabbitMQ обеспечивает гарантированную доставку с подтверждениями и очередями недоставленных сообщений; Celery добавляет высокоуровневое планирование, повторные попытки и мониторинг. Мониторинг глубины очереди предупреждает о накоплении задач до того, как оно приведёт к сбоям.
Изучай AI Agents с ИИ-репетитором — бесплатно
Пиши и запускай код прямо в браузере, получай мгновенную помощь от ИИ-репетитора 24/7 и продолжи учиться на сайте или в приложении.
- Курсы
- 60
- Уроки
- 239
Часто задаваемые вопросы
Урок «Очереди событий и брокеры сообщений» бесплатный?
Да — полный текст урока «Очереди событий и брокеры сообщений» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс AI Agents, подпишись на CoddyKit PRO. Курс AI Agents содержит 4 уроков всего.
Чему я научусь в уроке «Очереди событий и брокеры сообщений»?
Очереди Redis, RabbitMQ и Kafka для раздельного взаимодействия агентов. Ты практикуешь AI Agents с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.
Нужен ли мне опыт, чтобы начать AI Agents?
Предыдущий опыт не требуется. AI Agents на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 2 из 4.
Сколько времени занимает урок «Очереди событий и брокеры сообщений»?
Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.
Можно ли писать и запускать код в этом уроке AI Agents?
Да. Каждый урок AI Agents включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.
Все уроки этого курса
- Асинхронный Python для разработчиков агентов
- Очереди событий и брокеры сообщений
- Неблокирующее параллельное выполнение инструментов
- Асинхронные фреймворки агентов: LangChain и не только