Files d’événements et courtiers de messages
Files Redis, RabbitMQ et Kafka pour une communication découplée entre agents.
Files d’événements et courtiers de messages est une leçon AI Agents gratuite sur CoddyKit. Ceci est la leçon 2 sur 4. Tu peux lire la leçon complète ci-dessous gratuitement — puis la pratiquer en direct dans le navigateur avec un éditeur de code intégré et un tuteur IA 24/7. Elle fait partie du parcours d'apprentissage AI Agents, et ta progression se synchronise sur le web et l'application CoddyKit. Le cours AI Agents comprend 4 leçons au total.
Pourquoi utiliser des files de messages pour les agents ?
Les files de messages découplent le déclencheur (producteur d’événements) de l’agent (consommateur d’événements). Le producteur émet des événements à n’importe quel rythme ; l’agent les traite à son propre rythme. Cela évite la surcharge et permet les nouvelles tentatives.
Les listes Redis comme files d’attente
Les listes Redis peuvent servir de files d’attente simples : LPUSH ajoute un élément au début (mise en file), tandis que BRPOP retire un élément à la fin et se bloque si la liste est vide (retrait de la file). Il s’agit d’une file fiable selon le principe premier entré, premier sorti.
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)}')Boucle de travail d’un agent avec Redis
Un agent exécutant interroge continuellement la file d’attente à la recherche de nouvelles tâches, traite chacune d’elles, puis recommence. Si aucune tâche n’arrive avant l’expiration du délai d’attente, il recommence sans consommer de ressources.
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')Publication/abonnement Redis pour les événements
Le mécanisme de publication/abonnement de Redis diffère des listes : il diffuse les messages à tous les abonnés actuels. Utilisez-le pour les notifications en diffusion multiple lorsque plusieurs agents doivent réagir au même événement.
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 receiveBases de RabbitMQ avec pika
RabbitMQ est un courtier de messages complet qui prend en charge le routage, les accusés de réception et les files d’attente des messages non distribuables. La bibliothèque pika connecte 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()Consommateur RabbitMQ avec accusés de réception
Les consommateurs doivent accuser réception des messages après leur traitement. Si un agent plante avant d’accuser réception, RabbitMQ remet le message dans la file pour un autre consommateur. Cela garantit qu’aucun message n’est perdu.
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')Files d’attente des messages non distribuables
Les messages dont le traitement échoue trop souvent doivent être envoyés dans une file d’attente des messages non distribuables (DLQ) pour inspection manuelle, plutôt que d’être retentés indéfiniment. Configurez RabbitMQ pour acheminer automatiquement les messages échoués vers la 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()Modèle de file d’attente des tâches
Le modèle de file d’attente des tâches découple le déclenchement de l’exécution : un producteur léger ajoute les descriptions des tâches à la file ; les travailleurs les retirent et les exécutent. Plusieurs travailleurs peuvent traiter les tâches en parallèle.
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 : file d’attente des tâches de haut niveau
Celery est la file d’attente de tâches Python la plus populaire. Elle prend en charge Redis et RabbitMQ comme courtiers, les nouvelles tentatives automatiques, le routage des tâches et la surveillance avec Flower. Utilisez-la pour les charges de travail d’agents en 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')Surveillance de l’état des files d’attente
Surveillez la profondeur de la file d’attente pour détecter les accumulations. Si elle augmente, vous avez besoin de davantage de travailleurs ou l’agent est trop lent. Déclenchez une alerte lorsque la profondeur dépasse un seuil.
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")')Choisir entre Redis et RabbitMQ
Utilisez Redis si vous avez besoin de simplicité, si vous utilisez déjà Redis ou si les tâches sont de courte durée et que la perte de quelques messages en cas de panne est acceptable. Utilisez RabbitMQ si vous avez besoin d’une livraison garantie, d’un routage complexe, de files d’attente des messages non distribuables ou de plusieurs consommateurs avec des abonnements différents.
# 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')Vérification des connaissances : files de messages
Vérifiez votre compréhension des files de messages et des courtiers pour les agents.
Résumé des files de messages
Les files de messages découplent les déclencheurs des agents de leur exécution : les listes Redis fournissent des files simples selon le principe premier entré, premier sorti avec LPUSH/BRPOP ; la publication/abonnement Redis permet de diffuser des événements ; RabbitMQ garantit la livraison grâce aux accusés de réception et aux files d’attente des messages non distribuables ; Celery ajoute la planification de haut niveau, les nouvelles tentatives et la surveillance. La surveillance de la profondeur des files vous alerte en cas d’accumulation avant qu’elle ne provoque une interruption de service.
Questions Fréquemment Posées
La leçon « Files d’événements et courtiers de messages » est-elle gratuite ?
Oui — le texte complet de « Files d’événements et courtiers de messages » est gratuit à lire ici sur le web. Pour la pratiquer de manière interactive (un éditeur de code intégré et un tuteur IA 24/7) et déverrouiller le reste du cours AI Agents, passe à CoddyKit PRO. Le cours AI Agents comprend 4 leçons au total.
Qu'est-ce que j'apprendrai dans « Files d’événements et courtiers de messages » ?
Files Redis, RabbitMQ et Kafka pour une communication découplée entre agents. Tu pratiques AI Agents avec du code pratique que tu exécutes directement dans le navigateur, et un tuteur IA 24/7 répond à tes questions au fur et à mesure que tu avances dans la leçon.
Dois-je avoir de l'expérience pour commencer AI Agents ?
Aucune expérience préalable n'est requise. AI Agents sur CoddyKit est structuré pour les débutants jusqu'aux apprenants avancés, donc tu peux commencer ici ou depuis le début et avancer à ton rythme. Ceci est la leçon 2 sur 4.
Combien de temps prend la leçon « Files d’événements et courtiers de messages » ?
La plupart des leçons CoddyKit prennent environ 5–10 minutes. Chacune est courte et interactive, tu progresses régulièrement et tu repiques exactement où tu t'es arrêté sur le web et l'app.
Peux-tu écrire et exécuter du code dans cette leçon AI Agents ?
Oui. Chaque leçon AI Agents inclut un éditeur de code intégré, tu écris et exécutes du vrai code directement dans ton navigateur et tu reçois des retours IA instantanés — aucune configuration locale requise.
Toutes les leçons de ce cours
- Python asynchrone pour les développeurs d’agents
- Files d’événements et courtiers de messages
- Exécution parallèle non bloquante des outils
- Frameworks d’agents asynchrones : LangChain et au-delà