Code di eventi e message broker
Code Redis, RabbitMQ e Kafka per una comunicazione disaccoppiata tra agenti.
Code di eventi e message broker è una lezione AI Agents gratuita su CoddyKit. Questa è la lezione 2 di 4. Puoi leggere la lezione completa qui gratuitamente — poi esercitati direttamente nel browser con un editor di codice integrato e un tutor IA disponibile 24/7. Fa parte del percorso di apprendimento AI Agents, e i tuoi progressi si sincronizzano tra il web e l'app CoddyKit. Il corso AI Agents include 4 lezioni in totale.
Perché utilizzare le code di messaggi per gli agenti?
Le code di messaggi disaccoppiano il trigger (produttore di eventi) dall’agente (consumatore di eventi). Il produttore genera eventi a qualsiasi frequenza, mentre l’agente li elabora al proprio ritmo. Questo previene il sovraccarico e abilita i tentativi automatici.
Le liste Redis come code
Le liste Redis possono essere utilizzate come code semplici: LPUSH aggiunge in testa (accodamento), mentre BRPOP rimuove dalla coda e si blocca se è vuota (estrazione). Si tratta di una coda FIFO affidabile.
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)}')Ciclo di lavoro di un agente con Redis
Un worker dell’agente interroga continuamente la coda per cercare nuove attività, ne elabora una alla volta e ripete il ciclo. Se non arriva alcuna attività entro il timeout, ripete il ciclo senza consumare risorse.
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 per gli eventi
Redis Pub/Sub è diverso dalle liste: trasmette i messaggi a tutti i sottoscrittori attivi. Utilizzi Pub/Sub per le notifiche fan-out, quando più agenti devono reagire allo stesso evento.
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 receiveNozioni di base su RabbitMQ con pika
RabbitMQ è un broker di messaggi completo, con routing, conferme di ricezione e code dei messaggi non recapitabili. La libreria pika collega Python a 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()Consumer RabbitMQ con conferme di ricezione
I consumer devono confermare la ricezione dei messaggi dopo averli elaborati. Se un agente si arresta prima di confermare la ricezione, RabbitMQ rimette il messaggio in coda per un altro consumer. In questo modo si garantisce che nessun messaggio vada perso.
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')Code dei messaggi non recapitabili
I messaggi la cui elaborazione fallisce troppe volte devono essere spostati in una coda dei messaggi non recapitabili (DLQ) per un’ispezione manuale, invece di essere ritentati all’infinito. Configuri RabbitMQ in modo che indirizzi automaticamente i messaggi non riusciti alla 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()Pattern della coda di attività
Il pattern della coda di attività disaccoppia il trigger dall’esecuzione: un produttore leggero inserisce nella coda le descrizioni delle attività, mentre i worker le prelevano e le eseguono. Più worker possono elaborare le attività in parallelo.
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: coda di attività di alto livello
Celery è la coda di attività Python più diffusa. Supporta Redis e RabbitMQ come broker, i tentativi automatici, il routing delle attività e il monitoraggio con Flower. Lo utilizzi per i carichi di lavoro degli agenti in produzione.
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')Monitoraggio dello stato della coda
Monitori la profondità della coda per rilevare gli arretrati. Se la profondità aumenta, sono necessari più worker oppure l’agente è troppo lento. Configuri un avviso quando la profondità supera una soglia.
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")')Scegliere tra Redis e RabbitMQ
Utilizzi Redis quando ha bisogno di semplicità, utilizza già Redis oppure le attività sono di breve durata e la perdita di alcuni messaggi in caso di arresto è accettabile. Utilizzi RabbitMQ quando ha bisogno di consegna garantita, routing complesso, code dei messaggi non recapitabili o più consumer con sottoscrizioni diverse.
# 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')Verifica delle conoscenze: code di messaggi
Verifichi la Sua comprensione delle code di messaggi e dei broker per gli agenti.
Riepilogo delle code di messaggi
Le code di messaggi disaccoppiano i trigger degli agenti dall’esecuzione: le liste Redis forniscono semplici code FIFO con LPUSH/BRPOP; Redis Pub/Sub abilita la trasmissione di eventi; RabbitMQ offre la consegna garantita con conferme di ricezione e code dei messaggi non recapitabili; Celery aggiunge pianificazione di alto livello, tentativi automatici e monitoraggio. Il monitoraggio della profondità della coda consente di ricevere avvisi sugli arretrati prima che causino interruzioni del servizio.
Impara AI Agents con un tutor IA — gratis
Scrivi ed esegui vero codice nel tuo browser, ricevi aiuto istantaneo da un tutor IA disponibile 24/7, e riprendi da dove hai lasciato sul web o nell'app.
- Corsi
- 60
- Lezioni
- 239
Domande Frequenti
La lezione «Code di eventi e message broker» è gratuita?
Sì — il testo completo di «Code di eventi e message broker» è gratuito qui sul web. Per esercitarvi in modo interattivo (un editor di codice integrato e un tutor IA 24/7) e sbloccare il resto del corso AI Agents, passa a CoddyKit PRO. Il corso AI Agents include 4 lezioni in totale.
Cosa imparerò in «Code di eventi e message broker»?
Code Redis, RabbitMQ e Kafka per una comunicazione disaccoppiata tra agenti. Eserciti AI Agents con codice pratico che esegui direttamente nel browser, e un tutor IA 24/7 risponde alle tue domande mentre lavori sulla lezione.
Ho bisogno di esperienza per iniziare AI Agents?
Non è richiesta alcuna esperienza precedente. AI Agents su CoddyKit è strutturato per principianti e studenti avanzati, quindi puoi iniziare da qui o dall'inizio e procedere al tuo ritmo. Questa è la lezione 2 di 4.
Quanto tempo richiede la lezione «Code di eventi e message broker»?
La maggior parte delle lezioni CoddyKit richiede circa 5–10 minuti. Ogni lezione è breve e interattiva, quindi fai progressi costanti e riprendi esattamente da dove hai lasciato su web e app.
Posso scrivere ed eseguire codice in questa lezione AI Agents?
Sì. Ogni lezione AI Agents include un editor di codice integrato, quindi scrivi ed esegui codice reale direttamente nel tuo browser e ricevi feedback istantaneo dall'IA — nessuna configurazione locale necessaria.
Tutte le lezioni di questo corso
- Python asincrono per sviluppatori di agenti
- Code di eventi e message broker
- Esecuzione parallela non bloccante degli strumenti
- Framework asincroni per agenti: LangChain e oltre