Saltar al contenido

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.

Fan-out: un evento produce múltiples reacciones independientes

### El productor: generar eventos de pedidos

1import boto3
2import json
3import random
4import time
5from datetime import datetime
6
7eb = 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")
10
11productos = ["Laptop", "Teclado", "Monitor", "Auriculares", "Webcam", "SSD 1TB"]
12clientes = [f"CLI-{i}" for i in range(1, 20)]
13
14def 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 }
31
32# Producir 5 pedidos con 2s de separación
33for 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)
39
40print("\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 boto3
2import json
3import time
4
5sqs = 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")
8
9def 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}...")
13
14 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 procesamiento
25 time.sleep(0.5)
26 sqs.delete_message(QueueUrl=queue_url, ReceiptHandle=msg["ReceiptHandle"])
27 print(f" [{handler_name}] OK - confirmado")
28
29# Ejecutar: python consumidor.py inventario
30# Ejecutar: python consumidor.py emails
31# Ejecutar: python consumidor.py almacen
32if __name__ == "__main__":
33 import sys
34 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_id
2import uuid
3
4# El productor incluye un correlation_id
5evento_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}
16
17# El consumidor (servicio-pagos) responde con OTRO evento
18evento_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)
2
3def 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}
10
11def 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"]) == 2
25 assert resultado["actualizaciones"][0]["stock_restado"] == 1

Testear handlers de eventos como funciones puras, sin infraestructura

## ejercicios

[01]

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.

Cargando editor...
[02]

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...
Cargando editor...
[03]

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"]
Cargando editor...
[04]

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.

Cargando editor...

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...