Kolejki zdarzeń i brokerzy komunikatów
Redis, RabbitMQ i Kafka do rozdzielonej komunikacji agentów.
Kolejki zdarzeń i brokerzy komunikatów to bezpłatna lekcja AI Agents na CoddyKit. To lekcja 2 z 4. Możesz przeczytać całą lekcję poniżej za darmo — a potem ćwiczyć ją interaktywnie w przeglądarce z wbudowanym edytorem kodu i tutorem AI dostępnym 24/7. To część ścieżki edukacyjnej AI Agents, a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs AI Agents zawiera 4 lekcji w sumie.
Dlaczego używać kolejek komunikatów dla agentów?
Kolejki komunikatów oddzielają wyzwalacz (producenta zdarzeń) od agenta (konsumenta zdarzeń). Producent emituje zdarzenia w dowolnym tempie, a agent przetwarza je we własnym tempie. Zapobiega to przeciążeniu i umożliwia ponawianie prób.
Listy Redis jako kolejki
Listy Redis mogą działać jako proste kolejki: LPUSH dodaje element na początku (umieszcza go w kolejce), a BRPOP usuwa element z końca i blokuje działanie, gdy lista jest pusta (pobiera go z kolejki). Jest to niezawodna kolejka typu first-in-first-out.
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)}')Pętla robocza agenta z Redis
Pracownik agenta nieustannie sprawdza kolejkę w poszukiwaniu nowych zadań, przetwarza każde z nich i powtarza tę czynność. Jeśli w określonym czasie nie pojawią się żadne zadania, ponownie wykonuje pętlę bez zużywania zasobów.
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 dla zdarzeń
Redis Pub/Sub różni się od list: rozgłasza komunikaty wszystkim aktualnym subskrybentom. Pub/Sub należy stosować do powiadomień typu fan-out, gdy wiele agentów musi zareagować na to samo zdarzenie.
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 receivePodstawy RabbitMQ z pika
RabbitMQ to w pełni funkcjonalny broker komunikatów z routingiem, potwierdzeniami odbioru i kolejkami dead-letter. Biblioteka pika łączy Pythona z 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()Konsument RabbitMQ z potwierdzeniami odbioru
Konsumenci muszą potwierdzać odbiór komunikatów po ich przetworzeniu. Jeśli agent ulegnie awarii przed potwierdzeniem, RabbitMQ ponownie umieści komunikat w kolejce, aby mógł go przetworzyć inny konsument. Zapewnia to, że żadne komunikaty nie zostaną utracone.
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')Kolejki dead-letter
Komunikaty, których przetwarzanie nie powiodło się zbyt wiele razy, powinny trafiać do kolejki dead-letter (DLQ) w celu ręcznej kontroli, zamiast być ponawiane bez końca. Należy skonfigurować RabbitMQ tak, aby automatycznie kierował nieudane komunikaty do kolejki 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()Wzorzec kolejki zadań
Wzorzec kolejki zadań oddziela wyzwalacz od wykonania: prosty producent umieszcza w kolejce opisy zadań, a pracownicy pobierają je i wykonują. Wielu pracowników może przetwarzać zadania równolegle.
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: wysokopoziomowa kolejka zadań
Celery to najpopularniejsza kolejka zadań dla Pythona. Obsługuje Redis i RabbitMQ jako brokery, automatyczne ponawianie prób, routing zadań oraz monitorowanie za pomocą Flower. Należy używać go w produkcyjnych obciążeniach agentów.
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')Monitorowanie kondycji kolejki
Należy monitorować głębokość kolejki, aby wykrywać zaległości. Jeśli głębokość kolejki rośnie, potrzebnych jest więcej pracowników albo agent działa zbyt wolno. Należy generować alert, gdy głębokość przekroczy określony próg.
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")')Wybór między Redis a RabbitMQ
Redis należy wybrać, gdy potrzebna jest prostota, Redis jest już używany albo zadania są krótkotrwałe i akceptowalna jest utrata kilku komunikatów w razie awarii. RabbitMQ należy wybrać, gdy potrzebne są gwarantowane dostarczanie, złożony routing, kolejki dead-letter lub wielu konsumentów z różnymi subskrypcjami.
# 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')Sprawdzenie wiedzy: kolejki komunikatów
Proszę sprawdzić swoją wiedzę na temat kolejek komunikatów i brokerów używanych przez agentów.
Podsumowanie kolejek komunikatów
Kolejki komunikatów oddzielają wyzwalacze agentów od wykonywania zadań: listy Redis zapewniają proste kolejki FIFO za pomocą LPUSH/BRPOP; Redis Pub/Sub umożliwia rozgłaszanie zdarzeń; RabbitMQ zapewnia gwarantowane dostarczanie dzięki potwierdzeniom odbioru i kolejkom dead-letter; Celery dodaje wysokopoziomowe planowanie, ponawianie prób i monitorowanie. Monitorowanie głębokości kolejki pozwala wykrywać zaległości, zanim doprowadzą do awarii.
Często zadawane pytania
Czy lekcja „Kolejki zdarzeń i brokerzy komunikatów” jest bezpłatna?
Tak — pełny tekst „Kolejki zdarzeń i brokerzy komunikatów” jest dostępny za darmo tutaj w sieci. Aby ćwiczyć ją interaktywnie (wbudowany edytor kodu i tutor AI dostępny 24/7) i odblokować resztę kursu AI Agents, przejdź na CoddyKit PRO. Kurs AI Agents zawiera 4 lekcji w sumie.
Co nauczysz się w „Kolejki zdarzeń i brokerzy komunikatów”?
Redis, RabbitMQ i Kafka do rozdzielonej komunikacji agentów. Ćwiczysz AI Agents z praktycznym kodem, który uruchamiasz bezpośrednio w przeglądarce, a tutor AI dostępny 24/7 odpowiada na Twoje pytania podczas pracy nad lekcją.
Czy potrzebuję doświadczenia, aby zacząć AI Agents?
Nie wymagamy żadnego doświadczenia. AI Agents w CoddyKit jest strukturyzowany dla początkujących i zaawansowanych użytkowników, więc możesz zacząć tutaj lub od początku i uczyć się w swoim tempie. To lekcja 2 z 4.
Ile czasu zajmuje lekcja „Kolejki zdarzeń i brokerzy komunikatów”?
Większość lekcji CoddyKit trwa około 5–10 minut. Każda lekcja to mały, interaktywny krok, dzięki czemu robisz systematyczne postępy i zawsze wracasz dokładnie do tego samego miejsca — na webie i w aplikacji.
Czy mogę pisać i uruchamiać kod w tej lekcji AI Agents?
Tak. Każda lekcja AI Agents zawiera wbudowany edytor kodu, więc piszesz i uruchamiasz prawdziwy kod bezpośrednio w przeglądarce i od razu otrzymujesz sprzężenie zwrotne od AI — bez konfiguracji na komputerze.
Wszystkie lekcje w tym kursie
- Asynchroniczny Python dla twórców agentów
- Kolejki zdarzeń i brokerzy komunikatów
- Nieblokujące równoległe wykonywanie narzędzi
- Asynchroniczne frameworki agentów: LangChain i nie tylko