Ereigniswarteschlangen und Message Broker
Redis-Warteschlangen, RabbitMQ und Kafka für die entkoppelte Agentenkommunikation.
Ereigniswarteschlangen und Message Broker ist eine kostenlose AI Agents-Lektion auf CoddyKit. Dies ist Lektion 2 von 4. Du kannst die komplette Lektion unten kostenlos lesen – dann übst du sie direkt im Browser mit einem integrierten Code-Editor und einem KI-Tutor rund um die Uhr. Sie ist Teil des AI Agents-Lernpfads, und dein Fortschritt wird über Web und CoddyKit-App synchronisiert. Der AI Agents-Kurs umfasst insgesamt 4 Lektionen.
Warum Message Queues für Agenten?
Message Queues entkoppeln den Auslöser (Event-Produzent) vom Agenten (Event-Konsumenten). Der Produzent kann Ereignisse in beliebiger Geschwindigkeit auslösen, während der Agent sie in seinem eigenen Tempo verarbeitet. Dadurch werden Überlastungen verhindert und Wiederholungsversuche ermöglicht.
Redis-Listen als Queues
Redis-Listen können als einfache Queues verwendet werden: LPUSH fügt am Anfang ein (Enqueue), BRPOP entfernt am Ende ein Element und wartet bei einer leeren Liste (Dequeue). Dies ist eine zuverlässige First-in-first-out-Queue.
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)}')Agenten-Worker-Schleife mit Redis
Ein Agenten-Worker fragt die Queue kontinuierlich nach neuen Aufgaben ab, verarbeitet jede einzelne und wiederholt diesen Ablauf. Treffen innerhalb des Timeouts keine Aufgaben ein, beginnt die Schleife erneut, ohne Ressourcen zu verbrauchen.
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 für Ereignisse
Redis Pub/Sub unterscheidet sich von Listen: Es sendet Nachrichten an alle aktuell verbundenen Abonnenten. Verwenden Sie Pub/Sub für Fan-out-Benachrichtigungen, bei denen mehrere Agenten auf dasselbe Ereignis reagieren müssen.
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 receiveRabbitMQ-Grundlagen mit pika
RabbitMQ ist ein funktionsreicher Message Broker mit Routing, Bestätigungen und Dead-Letter-Queues. Die Bibliothek pika verbindet Python mit 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-Konsument mit Bestätigungen
Konsumenten müssen Nachrichten nach der Verarbeitung bestätigen. Wenn ein Agent vor der Bestätigung abstürzt, stellt RabbitMQ die Nachricht für einen anderen Konsumenten erneut in die Queue. Dadurch gehen keine Nachrichten verloren.
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')Dead-Letter-Queues
Nachrichten, deren Verarbeitung zu oft fehlschlägt, sollten zur manuellen Prüfung in eine Dead-Letter-Queue (DLQ) verschoben werden, statt sie endlos erneut zu verarbeiten. Konfigurieren Sie RabbitMQ so, dass fehlgeschlagene Nachrichten automatisch an die DLQ weitergeleitet werden.
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()Task-Queue-Muster
Das Task-Queue-Muster entkoppelt Auslöser und Ausführung: Ein schlanker Produzent stellt Aufgabenbeschreibungen in die Queue, und Worker holen sie ab und führen sie aus. Mehrere Worker können Aufgaben parallel verarbeiten.
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: High-Level-Task-Queue
Celery ist die beliebteste Task-Queue für Python. Sie unterstützt Redis und RabbitMQ als Broker, automatische Wiederholungsversuche, Task-Routing und die Überwachung mit Flower. Verwenden Sie Celery für produktive Agenten-Workloads.
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')Überwachung des Queue-Zustands
Überwachen Sie die Queue-Tiefe, um Rückstaus zu erkennen. Wenn die Queue-Tiefe wächst, benötigen Sie mehr Worker oder der Agent ist zu langsam. Lösen Sie einen Alarm aus, sobald die Tiefe einen Schwellenwert überschreitet.
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 oder RabbitMQ auswählen
Verwenden Sie Redis, wenn Sie eine einfache Lösung benötigen, Redis bereits einsetzen oder Aufgaben kurzlebig sind und der Verlust einiger Nachrichten bei einem Absturz akzeptabel ist. Verwenden Sie RabbitMQ, wenn Sie garantierte Zustellung, komplexes Routing, Dead-Letter-Queues oder mehrere Konsumenten mit unterschiedlichen Abonnements benötigen.
# 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')Wissenscheck: Message Queues
Testen Sie Ihr Verständnis von Message Queues und Brokern für Agenten.
Zusammenfassung: Message Queues
Message Queues entkoppeln Agenten-Auslöser von der Ausführung: Redis-Listen bieten einfache FIFO-Queues mit LPUSH/BRPOP; Redis Pub/Sub ermöglicht Broadcast-Ereignisse; RabbitMQ bietet mit Bestätigungen und Dead-Letter-Queues eine garantierte Zustellung; Celery ergänzt die Lösung um High-Level-Planung, Wiederholungsversuche und Monitoring. Die Überwachung der Queue-Tiefe weist auf Rückstaus hin, bevor daraus Ausfälle werden.
Häufig gestellte Fragen
Ist die Lektion „Ereigniswarteschlangen und Message Broker“ kostenlos?
Ja — der vollständige Text von „Ereigniswarteschlangen und Message Broker“ ist hier im Web kostenlos zu lesen. Um sie interaktiv zu üben (integrierter Code-Editor und 24/7 KI-Tutor) und den Rest des AI Agents-Kurses freizuschalten, upgrade auf CoddyKit PRO. Der AI Agents-Kurs umfasst insgesamt 4 Lektionen.
Was lerne ich in „Ereigniswarteschlangen und Message Broker“?
Redis-Warteschlangen, RabbitMQ und Kafka für die entkoppelte Agentenkommunikation. Du übst AI Agents mit praktischem Code, den du direkt im Browser ausführst, und ein 24/7 KI-Tutor beantwortet deine Fragen während du die Lektion bearbeitest.
Brauche ich Erfahrung, um AI Agents zu starten?
Keine Vorkenntnisse erforderlich. AI Agents auf CoddyKit ist für Anfänger bis fortgeschrittene Lernende strukturiert, sodass du hier starten oder von Anfang an beginnen und in deinem eigenen Tempo voranschreiten kannst. Dies ist Lektion 2 von 4.
Wie lange dauert die Lektion „Ereigniswarteschlangen und Message Broker“?
Die meisten CoddyKit-Lektionen dauern etwa 5–10 Minuten. Jede ist kompakt und interaktiv, sodass du stetig Fortschritte machst und genau dort weitermachst, wo du aufgehört hast – im Web und in der App.
Kann ich in dieser AI Agents-Lektion Code schreiben und ausführen?
Ja. Jede AI Agents-Lektion enthält einen integrierten Code-Editor, sodass du echten Code direkt in deinem Browser schreibst und ausführst und sofort KI-Feedback erhältst — ohne lokale Einrichtung erforderlich.
Alle Lektionen in diesem Kurs
- Asynchrones Python für Agent-Entwickler
- Ereigniswarteschlangen und Message Broker
- Nicht blockierende parallele Tool-Ausführung
- Asynchrone Agent-Frameworks: LangChain und darüber hinaus