lección 3
Producir y consumir eventos: tu primer flujo event-driven
Construye un pipeline event-driven completo: un productor genera eventos de pedidos y múltiples consumidores reaccionan.
⏱ 55 min
### De las piezas al flujo completo
En la lección anterior montaste la infraestructura: buses, colas, reglas. Ahora vamos a darle VIDA construyendo un flujo event-driven realista de principio a fin. Vamos a simular un e-commerce donde un pedido nuevo dispara TRES acciones simultáneas: actualizar inventario, enviar email de confirmación y notificar al almacén. Todo desacoplado. Todo reactivo.
La belleza de este patrón es que el servicio de pedidos NO SABE que existen los otros tres servicios. Solo dice "oye mundo, acaba de pasar un pedido". Y quien quiera enterarse, se entera. Si mañana añades un cuarto consumidor (analytics, por ejemplo), el productor NO CAMBIA. Cero código nuevo en el productor. Esa es la magia del desacoplamiento.
### Arquitectura del flujo
Nuestro sistema tendrá: un script productor que genera eventos de pedidos, un bus EventBridge que los enruta, y tres colas SQS (inventario, email, almacén) con un consumidor Python para cada una. Es como un periódico: el editor publica una noticia y los suscriptores la reciben según sus intereses.
### El productor: generar eventos de pedidos
1import boto32import json3import random4import time5from datetime import datetime67eb = boto3.client("events", endpoint_url="http://localhost:4566",8 region_name="eu-west-1", aws_access_key_id="test",9 aws_secret_access_key="test")1011productos = ["Laptop", "Teclado", "Monitor", "Auriculares", "Webcam", "SSD 1TB"]12clientes = [f"CLI-{i}" for i in range(1, 20)]1314def generar_pedido():15 """Genera un evento de pedido aleatorio."""16 items = [{"producto": p, "precio": round(random.uniform(20, 500), 2)}17 for p in random.sample(productos, k=random.randint(1, 3))]18 total = round(sum(i["precio"] for i in items), 2)19 return {20 "Source": "servicio-pedidos",21 "DetailType": "PedidoCreado",22 "Detail": json.dumps({23 "pedido_id": f"ORD-{random.randint(1000, 9999)}",24 "cliente_id": random.choice(clientes),25 "items": items,26 "total": total,27 "timestamp": datetime.utcnow().isoformat() + "Z",28 }),29 "EventBusName": "ecommerce-bus",30 }3132# Producir 5 pedidos con 2s de separación33for i in range(5):34 evento = generar_pedido()35 resp = eb.put_events(Entries=[evento])36 detail = json.loads(evento["Detail"])37 print(f"[{i+1}] Pedido {detail['pedido_id']} - Total: {detail['total']}EUR")38 time.sleep(2)3940print("\nProductor finalizado. 5 eventos publicados.")
Un productor que simula pedidos aleatorios cada 2 segundos
### El consumidor: reaccionar a los eventos
Cada consumidor es un script Python independiente que lee de SU cola específica. Vamos a escribir uno genérico que podemos apuntar a cualquier cola. En producción real, cada consumidor haría algo diferente: uno actualiza base de datos, otro envía emails, otro llama a una API externa.
1import boto32import json3import time45sqs = boto3.client("sqs", endpoint_url="http://localhost:4566",6 region_name="eu-west-1", aws_access_key_id="test",7 aws_secret_access_key="test")89def consumir(queue_name: str, handler_name: str):10 """Consume mensajes de una cola indefinidamente."""11 queue_url = f"http://sqs.eu-west-1.localhost.localstack.cloud:4566/000000000000/{queue_name}"12 print(f"[{handler_name}] Escuchando cola: {queue_name}...")1314 while True:15 resp = sqs.receive_message(16 QueueUrl=queue_url,17 MaxNumberOfMessages=5,18 WaitTimeSeconds=20,19 )20 for msg in resp.get("Messages", []):21 body = json.loads(msg["Body"])22 detail = body.get("detail", body)23 print(f" [{handler_name}] Procesando pedido {detail.get('pedido_id', '?')}")24 # Simular procesamiento25 time.sleep(0.5)26 sqs.delete_message(QueueUrl=queue_url, ReceiptHandle=msg["ReceiptHandle"])27 print(f" [{handler_name}] OK - confirmado")2829# Ejecutar: python consumidor.py inventario30# Ejecutar: python consumidor.py emails31# Ejecutar: python consumidor.py almacen32if __name__ == "__main__":33 import sys34 cola = sys.argv[1] if len(sys.argv) > 1 else "pedidos-nuevos"35 consumir(cola, cola.upper())
Consumidor genérico que escucha una cola SQS con long polling
### El patrón request-reply vs fire-and-forget
En nuestro flujo, el productor usa "fire-and-forget": lanza el evento y se desentiende. No espera confirmación de los consumidores. Esto es lo habitual en arquitecturas event-driven. Pero a veces necesitas una respuesta (¿se confirmó el pago?). Para eso existe el patrón request-reply: publicas un evento con un "correlation_id" y el consumidor responde con otro evento que incluye ese mismo ID.
1# Patrón Request-Reply con correlation_id2import uuid34# El productor incluye un correlation_id5evento_request = {6 "Source": "servicio-checkout",7 "DetailType": "VerificarPago",8 "Detail": json.dumps({9 "correlation_id": str(uuid.uuid4()),10 "pedido_id": "ORD-5555",11 "monto": 299.99,12 "tarjeta_hash": "xxx-1234"13 }),14 "EventBusName": "ecommerce-bus",15}1617# El consumidor (servicio-pagos) responde con OTRO evento18evento_reply = {19 "Source": "servicio-pagos",20 "DetailType": "PagoVerificado",21 "Detail": json.dumps({22 "correlation_id": "mismo-uuid-del-request",23 "pedido_id": "ORD-5555",24 "aprobado": True,25 "referencia_banco": "BNK-789456"26 }),27 "EventBusName": "ecommerce-bus",28}
Request-reply: correlacionar petición y respuesta a través de eventos
### Ordenación y consistencia eventual
Un tema CRÍTICO que muchos juniors ignoran: en un sistema event-driven, los eventos pueden llegar EN CUALQUIER ORDEN. SQS Standard no garantiza orden. Incluso con colas FIFO, si tienes múltiples productores, el orden "global" no existe. Esto significa que tu consumidor debe ser tolerante al desorden. Puede llegar "PedidoEnviado" antes que "PedidoCreado" si hay latencias diferentes.
Esto se llama CONSISTENCIA EVENTUAL: eventualmente todos los sistemas tendrán la misma vista del mundo, pero en un momento dado, diferentes servicios pueden tener datos ligeramente diferentes. Es el precio que pagas por el desacoplamiento y la escalabilidad. Y es un precio que VALE LA PENA para la inmensa mayoría de los casos.
Consejo de senior: diseña tus consumidores para tolerar mensajes fuera de orden. Usa timestamps y versiones en los eventos. Si recibes un evento "v3" y ya procesaste "v5", descarta el v3. Nunca asumas que los eventos llegan en el orden en que se produjeron.
Lo que le diría a mi yo de hace 5 años: antes de añadir un consumidor nuevo, pregúntate: ¿qué pasa si este consumidor está caído 2 horas? ¿Los mensajes se acumulan en la cola? ¿Hay backpressure? ¿La cola tiene retención suficiente? Piensa en los escenarios de fallo ANTES.
No confundas "event-driven" con "tiempo real". Un sistema event-driven reacciona a eventos, pero la latencia depende de muchos factores: congestión de la cola, velocidad del consumidor, reintentos. "Event-driven" no significa mágicamente milisegundos de latencia.
### Testing de sistemas event-driven
Testear un flujo event-driven es más complejo que testear una API síncrona. No puedes simplemente llamar y esperar una respuesta. Estrategias: 1) Test end-to-end con LocalStack: publica un evento y verifica que llega a la cola destino. 2) Test unitario del consumidor: pasa un mensaje directamente al handler sin infraestructura. 3) Test de contrato: verifica que el JSON del evento cumple el schema esperado.
1# Test unitario del handler (sin infraestructura)23def handler_inventario(evento: dict) -> dict:4 """Procesa un evento de pedido actualizando inventario."""5 items = evento["detail"]["items"]6 resultados = []7 for item in items:8 resultados.append({"producto": item["producto"], "stock_restado": item.get("cantidad", 1)})9 return {"status": "ok", "actualizaciones": resultados}1011def test_handler_inventario():12 evento = {13 "detail-type": "PedidoCreado",14 "detail": {15 "pedido_id": "ORD-TEST",16 "items": [17 {"producto": "Laptop", "cantidad": 1},18 {"producto": "Mouse", "cantidad": 2},19 ]20 }21 }22 resultado = handler_inventario(evento)23 assert resultado["status"] == "ok"24 assert len(resultado["actualizaciones"]) == 225 assert resultado["actualizaciones"][0]["stock_restado"] == 1
Testear handlers de eventos como funciones puras, sin infraestructura
## ejercicios
Construir un flujo de 3 consumidores
Escribe el setup completo: crea 3 colas (inventario, emails, almacen), configura las reglas de EventBridge para que un evento PedidoCreado llegue a las tres colas, y escribe un productor que publique 3 pedidos. Qué deberías ver: cada create_queue imprime "Cola creada: X", put_targets no lanza excepción, y cada put_events imprime "Evento ORD-N publicado" con FailedEntryCount=0.
Handler con lógica de reintento manual
Escribe un handler que procese un evento de envío de email. Si el envío falla (simula un 30% de fallos aleatorios), debe reintentar hasta 3 veces con backoff exponencial antes de rendirse.
💡 Resultado esperado
a@test.com: enviado
Reintento 1/3 de b@test.com en 1s...
Reintento 2/3 de b@test.com en 2s...
Reintento 3/3 de b@test.com en 4s...Validar el contrato de un evento
Escribe una función que valide que un evento cumple el contrato esperado: debe tener source (string), detail-type (string), time (ISO 8601), y detail (dict con pedido_id). Retorna True/False con lista de errores.
💡 Resultado esperado
Valido: True | Errores: [] Valido: False | Errores: ["Falta 'source' o está vacío", "'time' no es ISO 8601 válido: bad-date", "Falta 'detail.pedido_id'"] Valido: False | Errores: ["Falta 'source' o está vacío", "Falta 'detail-type' o está vacío"] Valido: False | Errores: ["'detail' no es un diccionario"]
Verificar que el fan-out funciona
Después de publicar un evento, verifica que ha llegado a las 3 colas. Escribe un script que publique 1 evento y luego lea de las 3 colas para confirmar la entrega. Qué deberías ver: las tres colas deben decir "OK (1 mensaje(s))". Si alguna dice "VACÍA (¿regla mal configurada?)", la regla de esa cola no está apuntando al target correcto.
Regístrate para guardar tu progreso.
## comentarios
Reporta erratas, ayuda a otros o comparte tu opinión. Sé constructivo.
Inicia sesión para comentar y responder.
cargando comentarios...