Las cuatro lecciones anteriores han dado a MercadoFresco las piezas: colas para desacoplar, temas para repartir, un bus para enrutar por contenido y flujos para orquestar. Pero han ido dejando preguntas aplazadas con un «lo veremos en 07-05», y todas apuntan al mismo sitio. ¿Qué significa exactamente «al menos una vez» cuando el que se duplica es un cobro de 48,20 €? ¿Cómo se escribe un consumidor que puede recibir el mismo mensaje tres veces sin cobrar tres veces? ¿Cuántas veces hay que reintentar, y cuándo parar para no tumbar al que se está recuperando? ¿Qué se hace exactamente con 340 mensajes en una DLQ un lunes por la mañana? ¿Y cómo se garantiza que un pedido escrito en Aurora se publica siempre, si no hay transacción que abarque la base de datos y el bus de eventos?
Esta lección no presenta servicios nuevos. Convierte los cuatro que ya conoces en criterio de diseño: cuándo usar cada uno, qué garantías dan de verdad y qué tienes que poner tú porque el servicio no lo pone. Es la lección que separa una arquitectura que funciona en la demo de una que sobrevive a un viernes de octubre con el proveedor de pagos a medio gas.
Aviso de coste. La tabla
mercadofresco-idempotenciabajo demanda con TTL de 24 h cuesta unos 0,90 USD al mes con el volumen de MercadoFresco. Lo caro no es esto: es un cobro duplicado, que además de los 48 € cuesta una llamada a atención al cliente y una devolución. Datos ficticios.
Contenido
- Garantías de entrega: las tres, y por qué una es casi mentira
- Idempotencia: la propiedad que lo arregla todo
- La tabla
mercadofresco-idempotencia - Un consumidor idempotente completo
- Reintentos: retroceso exponencial y jitter
- Qué errores merecen reintento y cuáles no
- Dónde reintentar: SDK, servicio o flujo
- La tormenta de reintentos
- Interruptor de circuito para el proveedor de pagos
- Colas de mensajes fallidos: clasificar y reprocesar
- El runbook de la DLQ de Marta
- Orden y agrupación: cuándo importa de verdad
- El patrón outbox y el problema de la doble escritura
- Contrapresión, amortiguación y estrangulamiento
- Tabla decisoria: SQS, SNS, EventBridge, Step Functions o síncrono
- La arquitectura de integración completa de MercadoFresco
- Errores comunes y consejos
- Ejercicios
- Conclusión
Garantías de entrega: las tres, y por qué una es casi mentira
| Garantía | Qué promete | Qué falla | Dónde aparece |
|---|---|---|---|
| Como mucho una vez | Nunca duplica | Puede perder | SNS a correo/SMS, UDP, «dispara y olvida» |
| Al menos una vez | Nunca pierde | Puede duplicar | SQS estándar, SNS a SQS, EventBridge, Lambda asíncrono |
| Exactamente una vez | Ni pierde ni duplica | Cuesta y casi nunca es de extremo a extremo | SQS FIFO (en su ventana), Step Functions estándar |
Prácticamente toda la mensajería seria elige al menos una vez, y por una razón sencilla: entre perder un pedido y procesarlo dos veces, lo segundo tiene arreglo y lo primero no.
El «exactamente una vez» casi nunca existe de extremo a extremo, y conviene entender por qué. Imagina el consumidor que cobra la tarjeta:
- Recibe el mensaje de
cola-mercadofresco-pedidos. - Llama a la pasarela. La pasarela cobra.
- La red se corta antes de que llegue la respuesta.
- El proceso muere, o lanza una excepción, y no borra el mensaje.
- El tiempo de visibilidad expira. Otro consumidor recibe el mismo mensaje.
- Llama a la pasarela. La pasarela cobra otra vez.
Ningún servicio de mensajería puede evitarlo, porque el problema no está en la entrega: está en que cobrar y confirmar que has cobrado son dos operaciones distintas y entre ellas cabe un fallo. SQS FIFO deduplica el envío del productor durante 5 minutos, no el efecto del consumidor. Step Functions estándar garantiza que su propia máquina no repite un paso, pero si la Lambda del paso cobró y falló al responder, Step Functions reintentará el paso.
La conclusión operativa es corta y hay que aceptarla sin resistencia: el sistema entrega al menos una vez; el efecto exactamente una vez lo pone tu consumidor. Y eso se llama idempotencia.
Idempotencia: la propiedad que lo arregla todo
Una operación es idempotente si ejecutarla varias veces con la misma entrada produce el mismo resultado que ejecutarla una vez. No significa «no hace nada la segunda vez»: significa que el estado final es el mismo.
| Operación | ¿Idempotente? | Por qué |
|---|---|---|
SET stock:FRUT-011 = 40 |
Sí | Escribir un valor absoluto |
DECR stock:FRUT-011 BY 2 |
No | Cada ejecución resta otra vez |
PutItem con la misma clave y datos |
Sí | Sobrescribe con lo mismo |
UpdateItem ADD contador 1 |
No | Incremento relativo |
INSERT con clave primaria pedido_id |
Sí (falla la segunda) | La restricción lo impide |
| Cobrar 48,20 € en la pasarela | No | Dos cobros, dos apuntes |
| Enviar un correo | No | Dos correos en la bandeja |
| Generar la miniatura de una foto | Sí | Sobrescribe el mismo objeto de S3 |
s3:PutObject con la misma clave |
Sí | Sobrescribe |
| Publicar un evento en el bus | No (los consumidores lo ven dos veces) | — |
La regla que se deduce: las operaciones que fijan un valor son idempotentes; las que lo modifican en relativo o producen efectos externos, no. Y cuando la operación no es idempotente por naturaleza, hay que hacerla idempotente añadiendo una clave de idempotencia.
Una clave de idempotencia es un identificador estable de la operación lógica, no del mensaje. La diferencia es crucial:
- ❌
MessageIdde SQS. Cambia si el productor reenvía el mismo trabajo. No sirve. - ❌ Un UUID generado en el consumidor. Cambia en cada intento. No sirve para nada.
- ✅
f"cobro:{pedido_id}"— deriva del dominio, es estable entre reintentos y entre productores. - ✅
f"correo-confirmacion:{pedido_id}"— un correo por pedido, sea cual sea el camino. - ✅
f"reserva:{pedido_id}:{sku}"— un movimiento de stock por línea de pedido. - ✅ Un hash del contenido canónico, cuando no hay identificador natural.
Fíjate en que la clave incluye qué operación además de sobre qué. PED-084417 a secas no vale:
el mismo pedido genera un cobro, un correo y varias reservas, y son operaciones distintas que deben
poder ejecutarse cada una una vez.
Muchas APIs de terceros aceptan la clave directamente. La pasarela de pago de MercadoFresco admite una
cabecera Idempotency-Key: enviando cobro:PED-2026-084417, un segundo intento devuelve el mismo
cobro en lugar de crear uno nuevo. Cuando el proveedor lo soporta, esa es siempre la mejor opción,
porque la garantía la da quien tiene el estado. Cuando no lo soporta, hay que llevar el registro
nosotros.
La tabla mercadofresco-idempotencia
aws dynamodb create-table \
--table-name mercadofresco-idempotencia \
--attribute-definitions AttributeName=clave_idempotencia,AttributeType=S \
--key-schema AttributeName=clave_idempotencia,KeyType=HASH \
--billing-mode PAY_PER_REQUEST \
--sse-specification Enabled=true,SSEType=KMS,KMSMasterKeyId=alias/mercadofresco-datos \
--tags Key=Proyecto,Value=mercadofresco Key=Entorno,Value=produccion \
Key=Componente,Value=integracion Key=Propietario,Value=marta \
Key=CentroCoste,Value=tecnologia \
--profile mercadofresco-dev --region eu-west-1
aws dynamodb update-time-to-live --table-name mercadofresco-idempotencia \
--time-to-live-specification 'Enabled=true,AttributeName=expira_en' \
--profile mercadofresco-dev --region eu-west-1Un elemento tiene esta forma:
{
"clave_idempotencia": "cobro:PED-2026-084417",
"estado": "COMPLETADO",
"resultado": { "referencia_cobro": "PAY-77321", "importe_eur": 48.20 },
"iniciado_en": "2026-08-02T18:41:07Z",
"completado_en": "2026-08-02T18:41:08Z",
"expira_en": 1785350467,
"ejecucion": "arn:aws:states:eu-west-1:111122223333:execution:...:PED-2026-084417"
}Cuatro decisiones de diseño, todas con motivo:
El estado tiene tres valores, no dos. EN_CURSO, COMPLETADO y FALLIDO. EN_CURSO es el que
resuelve el caso difícil: dos consumidores que reciben el mismo mensaje a la vez. El primero marca
EN_CURSO con una escritura condicional; el segundo ve que existe y se retira. Con solo dos estados,
ambos verían «no está» y cobrarían los dos.
Se guarda el resultado. Si el trabajo ya se hizo, el segundo intento no debe simplemente ignorar
el mensaje: debe devolver el mismo resultado que la primera vez. En una máquina de estados, eso
significa que el paso continúa con la misma referencia_cobro y el proceso sigue adelante en lugar de
romperse.
TTL de 24 horas (expira_en en segundos epoch, el mismo mecanismo de mercadofresco-carritos en
06-02). El plazo debe cubrir con holgura la vida máxima de un mensaje en el sistema: retención de la
cola (4 días) es demasiado, y una hora es demasiado poco si algo se queda en la DLQ y se reprocesa por
la tarde. 24 horas es el equilibrio de MercadoFresco. Sin TTL, la tabla crece indefinidamente y pagas
almacenamiento por registros que nadie va a consultar jamás.
Cifrada con alias/mercadofresco-datos, porque el resultado puede contener referencias de cobro.
Un consumidor idempotente completo
import json, os, time
from datetime import datetime, timezone
import boto3
from botocore.exceptions import ClientError
ddb = boto3.resource("dynamodb", region_name="eu-west-1")
tabla = ddb.Table("mercadofresco-idempotencia")
TTL_SEGUNDOS = 24 * 3600
PLAZO_EN_CURSO = 900 # 15 min: mas alla, se asume que el otro consumidor murio
class TrabajoEnCurso(Exception):
"""Otro consumidor esta procesando esta misma operacion ahora mismo."""
def reservar_ejecucion(clave):
"""Marca la operacion como EN_CURSO. Devuelve None si somos los primeros,
o el elemento existente si ya estaba (completado, fallido o en curso)."""
ahora = int(time.time())
try:
tabla.put_item(
Item={
"clave_idempotencia": clave,
"estado": "EN_CURSO",
"iniciado_en": datetime.now(timezone.utc).isoformat(),
"expira_en": ahora + TTL_SEGUNDOS,
},
# La escritura condicional es el corazon del patron: solo escribe si no
# existe, o si existe pero es un EN_CURSO abandonado hace mas de 15 min.
ConditionExpression="attribute_not_exists(clave_idempotencia) "
"OR (estado = :en_curso AND caduca_bloqueo < :ahora)",
ExpressionAttributeValues={":en_curso": "EN_CURSO", ":ahora": ahora},
)
return None # somos los primeros
except ClientError as e:
if e.response["Error"]["Code"] != "ConditionalCheckFailedException":
raise
return tabla.get_item(Key={"clave_idempotencia": clave})["Item"]
def completar(clave, resultado):
tabla.update_item(
Key={"clave_idempotencia": clave},
UpdateExpression="SET estado = :c, resultado = :r, completado_en = :t",
ExpressionAttributeValues={
":c": "COMPLETADO", ":r": resultado,
":t": datetime.now(timezone.utc).isoformat(),
},
)
def liberar(clave):
"""Fallo transitorio: se borra el registro para que el reintento pueda entrar."""
tabla.delete_item(Key={"clave_idempotencia": clave})
def procesar_pedido_idempotente(mensaje):
"""Consumidor de cola-mercadofresco-pedidos. Puede recibir el mismo mensaje N veces."""
pedido = json.loads(mensaje["Body"])
clave = f"cobro:{pedido['pedido_id']}" # clave de dominio, no MessageId
existente = reservar_ejecucion(clave)
if existente and existente["estado"] == "COMPLETADO":
# Ya se hizo. Devolvemos el MISMO resultado y borramos el mensaje.
log.info("duplicado ignorado", extra={"clave": clave})
return existente["resultado"]
if existente and existente["estado"] == "EN_CURSO":
# Otro consumidor lo tiene ahora mismo. NO borramos el mensaje: que vuelva.
raise TrabajoEnCurso(clave)
try:
resultado = pasarela.cobrar(
pedido["cliente_id"], pedido["importe_eur"],
idempotency_key=clave, # doble red: tambien en el proveedor
)
completar(clave, {"referencia_cobro": resultado.referencia,
"importe_eur": pedido["importe_eur"]})
return {"referencia_cobro": resultado.referencia}
except TarjetaRechazada:
# Fallo permanente: se registra como FALLIDO para no reintentar en bucle.
tabla.update_item(
Key={"clave_idempotencia": clave},
UpdateExpression="SET estado = :f", ExpressionAttributeValues={":f": "FALLIDO"},
)
raise
except Exception:
# Fallo transitorio: se libera el bloqueo para que el reintento pueda entrar.
liberar(clave)
raiseLos cuatro puntos que hacen que esto funcione de verdad:
- La escritura condicional es atómica. DynamoDB garantiza que solo un consumidor gana la carrera. No hace falta ningún bloqueo distribuido; es una única llamada.
- El
EN_CURSOcaduca. Sin la condición sobrecaduca_bloqueo, un consumidor que muera entre elput_itemy elcompletardejaría la clave bloqueada 24 horas y ese pedido no se cobraría nunca. - Fallo transitorio libera; fallo permanente no. Si la red falló, hay que permitir el reintento. Si la tarjeta se rechazó, reintentar es tirar peticiones a la basura.
- La clave viaja también a la pasarela. Si el proveedor soporta
Idempotency-Key, se usan las dos redes: la nuestra y la suya. Si la nuestra falla justo en el hueco entre cobrar y registrar, la suya nos cubre.
Powertools for AWS Lambda implementa este patrón completo con un decorador, incluyendo la tabla, el
TTL, el bloqueo y la caché del resultado. En Python basta con @idempotent sobre el manejador,
configurando DynamoDBPersistenceLayer y el event_key_jmespath que extrae la clave del evento. En
producción es lo razonable; escribirlo a mano una vez, como aquí, es lo que permite entender qué hace y
depurarlo cuando falle.
Reintentos: retroceso exponencial y jitter
Reintentar sin espera es contraproducente: si el servicio falla por saturación, los reintentos inmediatos lo saturan más. El retroceso exponencial (exponential backoff) multiplica la espera en cada intento; el jitter la aleatoriza.
Sin jitter (base 1 s, factor 2) Con jitter completo
intento 1 → 1,00 s intento 1 → 0,42 s
intento 2 → 2,00 s intento 2 → 1,17 s
intento 3 → 4,00 s intento 3 → 0,88 s
intento 4 → 8,00 s intento 4 → 6,31 s
intento 5 → 16,00 s intento 5 → 3,05 s
450 clientes reintentando 450 clientes reintentando
| | | | .. . ... . .. .. . ...
↑ ↑ ↑ ↑ distribuidos en el intervalo
picos sincronizados sin picosEl jitter no es un detalle de afinado: es lo que evita que la recuperación provoque la siguiente caída. Sin él, si la pasarela se cae 30 segundos, las 450 ejecuciones en curso reintentan exactamente en el mismo instante y la tumban otra vez justo cuando se estaba levantando.
import random, time
def con_reintentos(fn, intentos=5, base=1.0, techo=30.0):
"""Retroceso exponencial con jitter completo. Solo reintenta lo reintentable."""
for n in range(intentos):
try:
return fn()
except ErrorPermanente:
raise # 4xx de negocio: no se reintenta
except ErrorTransitorio as e:
if n == intentos - 1:
raise # se agotaron: que suba
espera = random.uniform(0, min(techo, base * (2 ** n)))
log.warning("reintento", extra={"intento": n + 1, "espera_s": round(espera, 2),
"error": str(e)})
time.sleep(espera)random.uniform(0, ...) es el jitter completo, la variante que mejor dispersa la carga según los
análisis de AWS. Existen alternativas —jitter «igualado» (mitad + aleatorio(0, mitad)) o
«descorrelacionado»— pero el completo es el más simple y el que mejor funciona en la mayoría de los
casos. El techo evita que el cuarto reintento espere ocho minutos.
Qué errores merecen reintento y cuáles no
| Error | ¿Reintentar? | Por qué |
|---|---|---|
500, 502, 503, 504 |
Sí | Fallo del servidor, probablemente transitorio |
429 / ThrottlingException |
Sí, con más espera | Estás yendo demasiado rápido |
ProvisionedThroughputExceededException |
Sí | Igual, en DynamoDB |
| Tiempo de espera de conexión o lectura | Sí, con cuidado | Puede haberse ejecutado: exige idempotencia |
400 / ValidationException |
No | El mensaje está mal; reintentar no lo arregla |
401 / 403 / AccessDenied |
No | Faltan permisos; hay que corregir la política |
404 |
Depende | Si es una escritura reciente, puede ser consistencia eventual |
409 / conflicto |
Depende | Con escritura condicional, suele significar «ya está hecho» |
Error de negocio (TarjetaRechazada) |
No | Reintentar no cambiará la decisión del banco |
La distinción práctica —5xx sí, 4xx no— tiene dos matices importantes. El 429 es un 4xx que sí
se reintenta, porque significa exactamente «vuelve más tarde», y hay que hacerlo con más espera de lo
normal (respeta la cabecera Retry-After si viene). Y los tiempos de espera agotados son el caso
peligroso: no sabes si la operación se ejecutó. Reintentarlos es correcto solo si tu operación es
idempotente; si no lo es, un reintento tras un tiempo de espera es exactamente cómo se cobra dos veces
a un cliente.
Un error muy común es tratar ConditionalCheckFailedException de DynamoDB como transitorio. No lo es:
significa que la condición no se cumplió —normalmente, que ya existe—, y reintentar dará el mismo
resultado. En el patrón de idempotencia, ese error es la respuesta correcta, no un fallo.
Dónde reintentar: SDK, servicio o flujo
Hay tres capas de reintento y se acumulan, lo que sorprende a mucha gente.
| Capa | Quién | Configuración | Alcance |
|---|---|---|---|
| SDK | boto3 / botocore | retries={"max_attempts": 5, "mode": "standard"} |
Llamadas a APIs de AWS |
| Servicio | SQS, Lambda, EventBridge, SNS | maxReceiveCount, MaximumRetryAttempts |
Entrega del mensaje |
| Flujo | Step Functions | Retry con BackoffRate y JitterStrategy |
Un paso del proceso |
El peligro es el producto. Si el SDK reintenta 5 veces, el consumidor recibe el mensaje 4 veces
(maxReceiveCount=4) y el Retry del flujo lo intenta 4 más, una caída de 30 segundos puede generar
80 llamadas para un solo pedido. Con 450 pedidos en curso son 36.000 llamadas contra un servicio que
ya estaba mal.
La disciplina de MercadoFresco: una capa manda y las demás son mínimas. En el proceso de pedido, la
capa que manda es el Retry de Step Functions, que es donde está la política de negocio; el SDK se deja
en max_attempts=2 para absorber fallos de red instantáneos, y el maxReceiveCount de las colas actúa
solo como red de seguridad frente a caídas del propio consumidor. Escribe en algún sitio cuántos
intentos totales puede recibir una dependencia en el peor caso: si el número supera 15, hay algo mal.
La tormenta de reintentos
La tormenta de reintentos (retry storm) es el fallo en cascada más habitual en sistemas distribuidos, y merece describirse entero porque casi siempre se diagnostica al revés:
- La pasarela de pago se degrada: la latencia sube de 240 ms a 4 s.
- Los consumidores esperan más, así que el ritmo de proceso cae y la cola crece.
- Lambda ve la cola crecer y escala: de 20 a 60 sondeadores.
- La pasarela recibe el triple de peticiones justo cuando peor está, y empieza a devolver 503.
- Cada 503 dispara reintentos. La carga se multiplica de nuevo.
- La pasarela cae del todo. Todos los reintentos fallan y consumen
maxReceiveCount. - Miles de pedidos legítimos acaban en la DLQ.
- La pasarela se recupera. Todos los reintentos pendientes salen a la vez y la tumban otra vez.
Cuatro defensas, y hacen falta las cuatro:
- Jitter en todos los reintentos. Rompe la sincronización de los pasos 5 y 8.
- Tope de concurrencia (
MaximumConcurrencyen el mapeo de origen de eventos). Evita el paso 3: la cola crece, pero la presión sobre la pasarela no. - Presupuesto de reintentos. Reintentar como máximo un porcentaje del tráfico (un 10 % es habitual). Si más del 10 % de las llamadas son reintentos, se dejan de reintentar hasta que la proporción baje.
- Interruptor de circuito. Deja de llamar del todo mientras el servicio esté caído.
Interruptor de circuito para el proveedor de pagos
El interruptor de circuito (circuit breaker) es un componente que cuenta fallos y, cuando pasan de un umbral, deja de llamar al servicio y falla de inmediato. Suena drástico y es exactamente lo que hace falta: si la pasarela lleva 20 fallos seguidos, la petición 21 también va a fallar, y lo único que consigue es gastar 30 segundos de tiempo de espera y mantener la presión sobre un servicio que intenta recuperarse.
stateDiagram-v2
[*] --> Cerrado
Cerrado --> Abierto: 20 fallos en 60 s
Abierto --> SemiAbierto: pasan 30 s
SemiAbierto --> Cerrado: 3 pruebas correctas
SemiAbierto --> Abierto: 1 prueba falla
note right of Cerrado
Pasa todo.
Se cuentan los fallos.
end note
note right of Abierto
No se llama al servicio.
Falla al instante (fail fast).
end note
note right of SemiAbierto
Deja pasar unas pocas
peticiones de prueba.
end note
El estado se guarda en mercadofresco-catalogo (ElastiCache/Valkey, 06-05), porque tiene que ser
compartido entre todos los consumidores: un interruptor en memoria de proceso no sirve de nada
cuando hay 20 Lambdas concurrentes, cada una con su propio contador.
r = redis.Redis(host="mercadofresco-catalogo.xxxxx.cache.amazonaws.com",
port=6379, ssl=True, decode_responses=True)
UMBRAL_FALLOS, VENTANA_S, ESPERA_ABIERTO_S, PRUEBAS_SEMI = 20, 60, 30, 3
class CircuitoAbierto(Exception):
"""El servicio esta marcado como caido: no se le llama."""
def llamar_con_interruptor(nombre, fn):
k_estado, k_fallos = f"cb:{nombre}:estado", f"cb:{nombre}:fallos"
k_exitos = f"cb:{nombre}:exitos_semi"
estado = r.get(k_estado) or "cerrado"
if estado == "abierto":
# SET NX: solo el primer consumidor que llega tras la espera pasa a semiabierto.
if r.set(f"cb:{nombre}:sonda", "1", nx=True, ex=ESPERA_ABIERTO_S):
r.set(k_estado, "semiabierto", ex=ESPERA_ABIERTO_S * 4)
r.delete(k_exitos)
else:
raise CircuitoAbierto(nombre) # falla al instante, sin llamar
try:
resultado = fn()
except Exception:
fallos = r.incr(k_fallos)
r.expire(k_fallos, VENTANA_S) # ventana deslizante aproximada
if estado == "semiabierto" or fallos >= UMBRAL_FALLOS:
r.set(k_estado, "abierto", ex=ESPERA_ABIERTO_S * 10)
r.delete(k_fallos, k_exitos)
raise
if estado == "semiabierto":
if r.incr(k_exitos) >= PRUEBAS_SEMI:
r.set(k_estado, "cerrado") # el servicio ha vuelto
r.delete(k_fallos, k_exitos, f"cb:{nombre}:sonda")
else:
r.delete(k_fallos) # racha limpia
return resultadoLo importante no es el código, es qué se hace cuando el circuito está abierto. Fallar rápido solo tiene valor si hay un plan:
- El cobro no puede degradarse: si la pasarela está caída, el pedido no se confirma y se le dice al cliente que lo intente en unos minutos. Es mejor que 30 segundos de espera y un error genérico.
- El aviso al ERP sí puede esperar: el mensaje se devuelve a la cola con
ChangeMessageVisibility(300)y se reintenta en cinco minutos, sin gastar recepciones. - El correo de confirmación puede degradarse a un proveedor secundario, o simplemente retrasarse.
Y el circuito debe ser observable: cada transición a abierto publica una métrica en
MercadoFresco/Tienda y una alarma avisa a Marta por alertas-mercadofresco. Un interruptor abierto que
nadie ve es una funcionalidad apagada en silencio.
Colas de mensajes fallidos: clasificar y reprocesar
Un mensaje en una DLQ no es un error: es una pregunta sin responder. Lo primero es clasificarlo, porque los tres tipos exigen acciones completamente distintas.
| Tipo | Síntoma | Causa | Acción |
|---|---|---|---|
| Envenenado | Todos los intentos fallan igual, error de parseo | Formato inválido, campo faltante, versión desconocida | Corregir el consumidor o descartar; nunca redrive tal cual |
| Transitorio | Ráfaga de mensajes en la misma franja horaria | Dependencia caída, límite superado | Redrive cuando el servicio vuelva |
| De datos | Un mensaje suelto, error de negocio | SKU descatalogado, cliente borrado, precio nulo | Corregir el dato o el mensaje, luego redrive |
La forma de distinguirlos en 30 segundos: mira la distribución temporal. Si los 340 mensajes entraron entre las 18:12 y las 18:41, es transitorio y el redrive lo arregla. Si entraron goteando a lo largo de tres días, es envenenado o de datos, y hacer redrive solo los devolverá a la DLQ.
def inspeccionar_dlq(url_dlq, muestra=10):
"""Lee sin borrar: devuelve el mensaje a los 5 s para no consumir recepciones."""
r = sqs.receive_message(
QueueUrl=url_dlq, MaxNumberOfMessages=muestra, WaitTimeSeconds=5,
VisibilityTimeout=5, MessageAttributeNames=["All"],
AttributeNames=["ApproximateReceiveCount", "SentTimestamp"],
)
resumen = collections.Counter()
for m in r.get("Messages", []):
try:
cuerpo = json.loads(m["Body"])
resumen[str((cuerpo.get("version", "?"), sorted(cuerpo.keys())[:3]))] += 1
except json.JSONDecodeError:
resumen["JSON_INVALIDO"] += 1
print(json.dumps({"message_id": m["MessageId"], "cuerpo": m["Body"][:300],
"recepciones": m["Attributes"]["ApproximateReceiveCount"],
"enviado_en": m["Attributes"]["SentTimestamp"]}, ensure_ascii=False))
return resumenVisibilityTimeout=5 es el detalle que convierte esto en una herramienta segura: los mensajes vuelven a
la DLQ enseguida y no se pierden si el script muere. Nunca inspecciones una DLQ borrando mensajes.
Y la alarma, que ya montamos en 07-01 pero que ahora se entiende del todo:
ApproximateNumberOfMessagesVisible > 0 sobre la DLQ es la señal de salud más valiosa de todo el
módulo, porque es la única que dice «hay trabajo pagado que no se ha hecho». Va a
alertas-mercadofresco con umbral 0 y periodo de 5 minutos.
El runbook de la DLQ de Marta
Un runbook es un procedimiento escrito que alguien puede seguir a las 3 de la mañana sin pensar. Este es el de MercadoFresco, y vive en el repositorio junto a la definición de las colas.
1. Contener. ¿Sigue creciendo? Mira ApproximateNumberOfMessagesVisible de la DLQ y
ApproximateAgeOfOldestMessage de la cola de origen. Si crece, el problema está vivo: el redrive
puede esperar; primero se para la hemorragia. Si es la dependencia la que está caída, considera
desactivar temporalmente el mapeo de origen de eventos para que los mensajes se acumulen en la cola —que
los guarda 4 días— en lugar de agotar reintentos y caer en la DLQ.
2. Clasificar. Ejecuta inspeccionar_dlq sobre 10 mensajes. Anota: ¿misma franja horaria o
goteo? ¿mismo error? ¿ApproximateReceiveCount uniforme? Decide envenenado, transitorio o de datos.
3. Diagnosticar. Busca en CloudWatch Logs por el MessageId de dos o tres mensajes para ver la
excepción real. Cruza la franja horaria con las métricas de las dependencias y con trail-mercadofresco
por si hubo un cambio de configuración.
4. Corregir. Según el tipo: desplegar el consumidor arreglado, esperar a que el servicio vuelva, o corregir el dato de origen. No hay redrive sin corrección previa; devolver mensajes a una cola cuyo consumidor sigue roto solo duplica el trabajo y llena los registros.
5. Reprocesar. Con límite de velocidad, siempre:
aws sqs start-message-move-task \
--source-arn arn:aws:sqs:eu-west-1:111122223333:mercadofresco-pedidos-fallidos \
--max-number-of-messages-per-second 20 \
--profile mercadofresco-dev --region eu-west-1
aws sqs list-message-move-tasks \
--source-arn arn:aws:sqs:eu-west-1:111122223333:mercadofresco-pedidos-fallidos \
--profile mercadofresco-dev --region eu-west-16. Verificar. La DLQ debe quedar a cero y NumberOfMessagesDeleted de la cola de origen debe subir
en la cantidad esperada. Si los mensajes vuelven a la DLQ, para y regresa al paso 2: la corrección no
era la correcta.
7. Registrar. Cuántos mensajes, qué causa, qué se cambió y qué habría avisado antes. Este paso es el que evita que el mismo incidente ocurra tres veces.
Una advertencia sobre el reprocesado. Todos los mensajes que vuelven pasarán otra vez por el consumidor, y algunos pudieron procesarse parcialmente antes de fallar. Si el consumidor no es idempotente, el redrive es una máquina de generar duplicados. Aquí es donde la primera mitad de esta lección deja de ser teoría.
Orden y agrupación: cuándo importa de verdad
El orden se pide mucho más de lo que se necesita, y cuesta caro. Antes de exigirlo, hazte tres preguntas: ¿los mensajes afectan al mismo dato? ¿la operación es relativa (sumar, restar) o absoluta (fijar)? ¿el resultado final cambia si se aplican al revés?
| Caso | ¿Orden? | Por qué |
|---|---|---|
| Reservar y liberar 2 unidades del mismo SKU | Sí | Operaciones relativas sobre el mismo contador |
| Correos de dos pedidos distintos | No | Independientes |
| «Pedido creado» y «pedido cancelado» del mismo pedido | Sí | Cancelar antes de crear no tiene sentido |
| Actualizaciones de precio del mismo SKU | Depende | Si llevan marca de tiempo, gana la última: no hace falta |
| Eventos de analítica | No | Se agregan; el orden es irrelevante |
El orden global es carísimo y casi nunca necesario. Exigirlo en una cola FIFO con un solo
MessageGroupId limita a un mensaje en vuelo: da igual que tengas 50 consumidores, procesas en serie.
Con el sku como grupo, MercadoFresco tiene 3.400 grupos, orden garantizado donde importa y paralelismo
de 3.400 donde no.
Elegir la clave de grupo es la misma decisión que la clave de partición de DynamoDB (06-02): tan
granular como el orden lo permita. Para los movimientos de stock, el sku. Para el ciclo de vida de un
pedido —creado, pagado, preparado, enviado—, el pedido_id, porque el orden importa dentro de un pedido
y no entre pedidos.
Hay una alternativa que evita FIFO por completo y conviene conocer: hacer que el orden no importe.
Si cada mensaje lleva una marca de tiempo o un número de versión y el consumidor descarta lo que sea más
antiguo que el estado actual —una escritura condicional del tipo if version > version_actual—, los
mensajes desordenados se resuelven solos. Es más trabajo en el consumidor y muchísimo mejor en
rendimiento. Es la misma idea que la idempotencia: el consumidor robusto vale más que la garantía cara
del transporte.
El patrón outbox y el problema de la doble escritura
En 07-01 dejamos pendiente este agujero. El código era:
pedido_id = aurora.insertar_pedido(carrito, cliente, cobro.referencia) # 1
sqs.send_message(QueueUrl=COLA_PEDIDOS, MessageBody=json.dumps(cuerpo)) # 2No hay transacción que abarque Aurora y SQS. Si el paso 2 falla —red, estrangulamiento, el proceso muere— el pedido existe en la base de datos y nadie se entera: no se prepara, no se envía correo, no se reparte. Invertir el orden no ayuda: entonces el riesgo es anunciar un pedido que no existe, que es peor. Es el problema de la doble escritura (dual write), y no se resuelve con reintentos, porque el proceso puede morir entre las dos operaciones.
Solución 1: tabla outbox. Se escribe el evento en la misma transacción que el pedido, en una tabla de la propia base de datos. Un proceso aparte lee esa tabla y publica.
BEGIN;
INSERT INTO pedidos (pedido_id, cliente_id, importe_eur, estado)
VALUES ('PED-2026-084417', 'CLI-30912', 48.20, 'confirmado');
INSERT INTO outbox (id, agregado, tipo_evento, carga, publicado)
VALUES (gen_random_uuid(), 'PED-2026-084417', 'PedidoConfirmado',
'{"pedido_id":"PED-2026-084417","importe_eur":48.20}'::jsonb, false);
COMMIT;O el pedido y su evento existen, o no existe ninguno de los dos: eso es lo que da la transacción.
Después, un publicador lee las filas con publicado = false, las envía y las marca. Puede fallar y
reintentar sin problema: publicará el evento dos veces como mucho, que es «al menos una vez», que es
justo lo que los consumidores idempotentes de esta lección saben tolerar.
def publicar_outbox():
"""Se ejecuta cada segundo. Idempotente y reanudable."""
filas = aurora.query(
"SELECT id, agregado, tipo_evento, carga FROM outbox "
"WHERE publicado = false ORDER BY creado_en LIMIT 100 FOR UPDATE SKIP LOCKED"
)
for f in filas:
eb.put_events(Entries=[{
"EventBusName": "bus-mercadofresco",
"Source": "mercadofresco.tienda",
"DetailType": f["tipo_evento"],
"Detail": f["carga"],
}])
aurora.execute("UPDATE outbox SET publicado = true, publicado_en = now() "
"WHERE id = %s", f["id"])FOR UPDATE SKIP LOCKED permite varios publicadores en paralelo sin que se pisen, y ORDER BY creado_en
conserva el orden por si importa.
Solución 2: captura de cambios sobre Streams. Si el almacén de escritura es DynamoDB, el patrón es
aún más limpio: se escribe solo en la tabla, y DynamoDB Streams genera el evento
automáticamente. No hay doble escritura porque solo hay una escritura. Un pipe-mf-... de EventBridge
Pipes (07-03) lee el stream, filtra y publica en el bus. Para Aurora existe el equivalente con
Debezium/DMS leyendo el registro de transacciones, aunque es bastante más pesado de operar.
| Tabla outbox | Streams + Pipes | |
|---|---|---|
| Dónde vive el estado | Base de datos relacional | DynamoDB |
| Piezas que mantener | Tabla + publicador + limpieza | Ninguna: es gestionado |
| Orden | Controlable con ORDER BY |
Por clave de partición |
| Latencia | La del sondeo (1–5 s) | Menos de 1 s |
| Carga extra en la BD | Escritura + sondeo | Ninguna |
| Cuándo usarlo | Aurora es la fuente de la verdad | DynamoDB es la fuente de la verdad |
MercadoFresco usa las dos: outbox en Aurora para los eventos de pedido, y Pipes sobre los Streams de
mercadofresco-carritos para CarritoAbandonado. Y la tabla outbox necesita su propia limpieza —borrar
lo publicado hace más de 7 días— o crecerá sin control, exactamente como los carritos zombis de 06-02.
Contrapresión, amortiguación y estrangulamiento
La cola como amortiguador. El viernes entran 900 pedidos/hora en ráfagas; los consumidores procesan 600. Sin cola, las 300 diferencias serían errores 503 en la cara del cliente. Con cola, la profundidad sube a 1.200 mensajes hacia las 20:00 y vuelve a cero a las 22:00. El pico se convierte en tiempo, que es un recurso mucho más barato que la capacidad.
La regla para dimensionar es sencilla: el sistema no necesita aguantar el pico, necesita aguantar la media más un margen, y tener cola suficiente para el área bajo la curva del pico. Si el pico dura 4 horas con 300 mensajes/hora de exceso, la cola llegará a unos 1.200 mensajes: perfectamente normal.
Límites de concurrencia. Aurora aurora-mf-escritor aguanta unas 200 conexiones. Si Lambda escala a
400 invocaciones concurrentes, cada una con su conexión, la base de datos rechaza conexiones y todo
falla, incluida la tienda. Tres capas de defensa:
| Capa | Mecanismo | Valor en MercadoFresco |
|---|---|---|
| Cola → Lambda | --scaling-config MaximumConcurrency |
20 en correo, 12 en ERP |
| Lambda (cuenta) | Concurrencia reservada por función | 50 para las que tocan Aurora |
| Aurora | RDS Proxy o pool en la aplicación | Multiplexa 400 clientes sobre 100 conexiones |
Estrangulamiento controlado. Cuando el que se satura es un tercero, hay que limitar el ritmo desde
nuestro lado: invocation-rate-limit-per-second en los destinos de API de EventBridge (07-03),
throttlePolicy.maxReceivesPerSecond en las políticas de entrega de SNS (07-02), MaxConcurrency en un
Map de Step Functions (07-04). Todos son la misma idea: es mejor ir despacio a propósito que ir
deprisa y provocar una caída.
Y una regla de higiene: contrapresión hacia arriba, nunca hacia abajo. Cuando el sistema va saturado, la respuesta correcta es dejar que la cola crezca y avisar, no aumentar el paralelismo contra una dependencia que ya está sufriendo.
Tabla decisoria: SQS, SNS, EventBridge, Step Functions o síncrono
| Servicio | La pregunta que responde | Úsalo cuando | Evítalo cuando |
|---|---|---|---|
| Llamada síncrona | «¿Cuál es el resultado, ahora?» | El usuario necesita el dato para continuar | El resultado no forma parte de la respuesta |
| SQS | «¿Quién hace este trabajo, cuando pueda?» | Un consumidor, trabajo duradero, amortiguación | Hay varios interesados |
| SNS | «¿A quiénes hay que avisar de esto?» | Fan-out de un hecho, latencia mínima, SMS/correo | Hace falta enrutar por contenido complejo |
| EventBridge | «¿A dónde va esto según lo que dice?» | Varios tipos de evento, eventos de AWS/SaaS, archivo | Volumen enorme con latencia crítica |
| Step Functions | «¿En qué punto va el proceso y qué hay que deshacer?» | Pasos dependientes, esperas, compensaciones | Hechos independientes sin estado |
Tres combinaciones que resuelven la mayoría de los casos reales: SNS→SQS (fan-out con durabilidad), EventBridge→SQS→Lambda (enrutamiento por contenido con amortiguación) y EventBridge→Step Functions (un hecho arranca un proceso). Y una que casi nunca es buena idea: SNS→Lambda directo para trabajo que importa, por todo lo dicho en 07-02.
La arquitectura de integración completa de MercadoFresco
flowchart TD
CLI([Cliente pulsa Confirmar pedido]) --> APP[Tienda en asg-mercadofresco-tienda]
APP -->|1. cobrar 240 ms| PAY[/Pasarela de pago/]
APP -->|2. INSERT pedido + outbox<br/>misma transaccion 38 ms| AUR[(aurora-mercadofresco-pedidos)]
APP -->|3. responde ~400 ms p95| CLI
AUR -.->|publicador de outbox| BUS{{bus-mercadofresco}}
DDB[(mercadofresco-carritos<br/>Streams)] -->|pipe-mf-carritos-abandonados| BUS
BUS -->|regla: PedidoConfirmado| SFN[[mercadofresco-procesar-pedido]]
BUS -->|regla: fan-out| TEMA{{mercadofresco-pedido-confirmado}}
BUS -->|regla: StockBajo| ALERT{{alertas-mercadofresco}}
TEMA --> QA[(cola-mercadofresco-almacen)]
TEMA --> QB[(cola-mercadofresco-correo)]
TEMA --> QC[(cola-mercadofresco-analitica)]
SFN -->|waitForTaskToken| QA
QA --> ERP[/ERP del almacen/]
QB --> LC[Lambda correo<br/>idempotente, max 20] --> SES[/SES/]
QC --> RS[(wg-mercadofresco-analitica)]
QA -.->|4 fallos| DLQ[(mercadofresco-pedidos-fallidos)]
QB -.->|4 fallos| DLQ
DLQ -.->|alarma| ALERT
LC -.->|clave de idempotencia| IDEM[(mercadofresco-idempotencia<br/>TTL 24 h)]
SFN -.->|compensacion| ALERT
style APP fill:#cfe2ff
style BUS fill:#cfe2ff
style SFN fill:#d1e7dd
style DLQ fill:#f8d7da
style IDEM fill:#fff3cd
Lo que se ha ganado, medido:
| Antes del módulo 7 | Después | |
|---|---|---|
| Confirmación del pedido (p50 / p95) | 2.893 ms / 9.100 ms | 312 ms / 400 ms |
| Puntos de fallo síncronos | 8 | 2 (pasarela y Aurora) |
| Pedidos rotos tras cobrar, por hora en pico | ~7 | 0 en tres meses |
| Caída de 2 h del ERP | ~1.800 ventas perdidas | 0: se procesan al volver |
| Añadir un consumidor nuevo | Despliegue de la tienda | Una suscripción |
| Saber en qué punto va un pedido | Imposible | Historial de la ejecución |
| Coste mensual de integración | 0 | ~62 USD |
Sesenta y dos dólares al mes —SQS casi gratis, SNS 2,60, EventBridge 0,95, Step Functions 54, DynamoDB 0,90— a cambio de que confirmar un pedido vuelva a ser una sola cosa rápida y fiable, y de que el resto del mundo se entere a su ritmo.
Errores Comunes y Consejos
Usar el MessageId como clave de idempotencia. Cambia si el productor reenvía el mismo trabajo, así
que no protege del caso más frecuente. La clave se deriva del dominio: cobro:<pedido_id>.
Tabla de idempotencia sin TTL. Crece indefinidamente y pagas almacenamiento por registros que nadie consultará. TTL de 24 h.
Idempotencia con solo dos estados. Sin EN_CURSO, dos consumidores simultáneos ven «no está» y
ejecutan los dos. La escritura condicional con EN_CURSO caducable es lo que cierra la carrera.
Bloquear con EN_CURSO sin caducidad. Un consumidor que muere a mitad deja esa operación bloqueada
hasta que expire el TTL, y ese pedido no se procesa nunca.
Reintentar sin jitter. La recuperación provoca la siguiente caída.
Reintentar errores 4xx de negocio. Retrasa la respuesta y no arregla nada. Excepción: el 429, que
sí se reintenta y con más espera.
Acumular reintentos en tres capas. SDK × servicio × flujo puede multiplicar por 80 la carga sobre una dependencia degradada. Una capa manda; las demás, mínimas.
Redrive sin corregir antes. Los mensajes vuelven a la DLQ, y si el consumidor no es idempotente, además duplican efectos.
Inspeccionar una DLQ borrando mensajes. Usa VisibilityTimeout corto y no borres nunca durante el
diagnóstico.
Pedir orden global «por si acaso». Un solo MessageGroupId convierte una cola distribuida en un
proceso serie.
Escribir en la base de datos y publicar sin transacción. El problema de la doble escritura. Outbox o Streams; no hay tercera opción que funcione.
Consejo: hazlo idempotente antes de hacerlo rápido. Un consumidor idempotente permite reintentar sin miedo, reprocesar DLQ, reproducir archivos de EventBridge y desplegar dos veces por error. Es la propiedad que más problemas evita por línea de código.
Consejo: escribe el runbook antes del incidente. A las 3 de la mañana no se diseña un procedimiento.
Consejo: mide el peor caso de intentos por dependencia. Si el número supera 15, revisa la configuración: alguna capa está reintentando de más.
Ejercicios
Ejercicio 1: el cobro duplicado del viernes
Un cliente reclama que se le cobraron dos veces 48,20 €. Los datos: cola-mercadofresco-pedidos con
VisibilityTimeout 120 s; el consumidor llama a la pasarela con read_timeout de 180 s; el SDK está en
max_attempts=5; maxReceiveCount=4; no hay tabla de idempotencia; la pasarela soporta
Idempotency-Key pero no se usa; los registros muestran un pico de latencia de la pasarela a las 19:14.
Responde: (a) los dos mecanismos independientes que pudieron duplicar el cobro; (b) cuál es más probable dados los números; (c) cinco correcciones ordenadas por eficacia; (d) cuál es la única que elimina el problema de raíz; (e) qué habría cambiado si se hubiera usado una cola FIFO.
Ejercicio 2: 340 mensajes en la DLQ un lunes
mercadofresco-pedidos-fallidos tiene 340 mensajes. ApproximateAgeOfOldestMessage de la cola de
origen es 12 s (normal). Al inspeccionar 10 mensajes: todos con ApproximateReceiveCount = 5,
SentTimestamp repartido entre el domingo a las 22:04 y las 23:51, y todos con "version": 2 en el
cuerpo mientras el consumidor espera "version": 1. Ese domingo por la noche hubo un despliegue.
Aplica el runbook: (a) los pasos 1 y 2 con tu conclusión; (b) qué tipo de fallo es y por qué la distribución temporal despista; (c) las dos correcciones posibles con sus pros y contras; (d) el comando de reprocesado y por qué el límite de velocidad importa aquí; (e) qué salvaguarda de diseño lo habría evitado.
Ejercicio 3: rediseñar la publicación del pedido
La tienda hace INSERT en Aurora y luego put_events en bus-mercadofresco. Marta detecta que 1 de
cada 900 pedidos existe en Aurora pero no generó evento: no se preparó, no se avisó al cliente y nadie
lo supo hasta que reclamó.
Diseña la solución: (a) por qué reintentar el put_events no basta; (b) el esquema de la tabla outbox y
la transacción SQL; (c) el publicador, indicando cada cuánto corre y cómo evita duplicar y pisarse con
otros; (d) qué garantía de entrega ofrece el conjunto y qué exige a los consumidores; (e) qué dos cosas
más hay que operar que antes no existían.
Soluciones
Solución 1
(a) Los dos mecanismos. Uno: el tiempo de visibilidad es menor que el tiempo de proceso. La cola
tiene 120 s y el consumidor puede esperar 180 s a la pasarela. Si a las 19:14 la pasarela tardó 150 s,
el mensaje se hizo visible a los 120 s, otro consumidor lo cogió y cobró en paralelo. Dos: los
reintentos del SDK sobre una llamada no idempotente. Con max_attempts=5, si la pasarela cobró y
falló al responder —o tardó más que el tiempo de espera—, botocore reintenta y cobra otra vez, todo
dentro de la misma recepción del mensaje.
(b) Cuál es más probable. El primero. La coincidencia entre el pico de latencia y la relación 120 < 180 es demasiado exacta: cualquier llamada que superara los 120 s producía una doble entrega garantizada. El segundo mecanismo también actuó probablemente, pero requiere que la llamada falle de una forma concreta, mientras que el primero solo requiere que tarde.
(c) Cinco correcciones por eficacia.
- Tabla de idempotencia con clave
cobro:<pedido_id>en el consumidor. Es la única que hace inofensivo el duplicado, venga de donde venga. - Usar
Idempotency-Keyen la pasarela. Coste casi nulo y la garantía la da quien tiene el estado del cobro. Debería haberse hecho desde el primer día. VisibilityTimeouta 400 s, holgadamente por encima delread_timeoutde 180 s.read_timeouta 30 s ymax_attemptsdel SDK a 2. Ciento ochenta segundos es demasiado para una pasarela; si tarda tanto, es mejor fallar y reintentar de forma controlada.- Alarma sobre la latencia p99 de la pasarela y un interruptor de circuito, para dejar de llamarla cuando se degrada en lugar de acumular llamadas de 150 s.
(d) La única que elimina el problema de raíz es la 1 (con la 2 como refuerzo). Las correcciones 3, 4 y 5 reducen la probabilidad, pero SQS entrega al menos una vez por diseño: puede duplicar aunque todo esté perfectamente configurado. Solo un consumidor idempotente convierte el duplicado en un no-evento.
(e) Con cola FIFO. Habría ayudado poco. FIFO deduplica el envío del productor dentro de una
ventana de 5 minutos, y aquí el problema estaba en el consumidor: el mensaje se entregó una vez y se
procesó dos por expiración de visibilidad. La única aportación de FIFO sería que, con
MessageGroupId = pedido_id, no habría un segundo mensaje del mismo pedido en vuelo —pero el mismo
mensaje reentregado sigue siendo el mismo mensaje, y el efecto duplicado ocurre igual—. Es un buen
ejemplo de FIFO usado como sustituto de la idempotencia, que es un error caro.
Solución 2
(a) Pasos 1 y 2. Contener: la cola de origen tiene 12 s de antigüedad, es decir, el problema no
está vivo. El consumidor actual funciona con normalidad y no hay hemorragia; se puede diagnosticar sin
prisa. Clasificar: los 340 mensajes entraron en una ventana de 1 h 47 min del domingo por la noche,
coincidiendo con un despliegue, y todos tienen ApproximateReceiveCount = 5 y "version": 2.
Conclusión: son mensajes envenenados, generados por un productor desplegado antes que su consumidor.
(b) Tipo de fallo y por qué despista. Son envenenados —fallo de formato, no de disponibilidad—
pero su distribución temporal es la de un fallo transitorio: concentrados en una ventana. Lo que
despista es que la ventana coincide con el intervalo en que el productor nuevo convivió con el consumidor
antiguo, no con una caída. Lo que desambigua es el contenido: "version": 2 frente a un consumidor que
espera 1. Regla: la distribución temporal es una pista, el contenido es la prueba.
(c) Dos correcciones. Desplegar el consumidor que entiende la versión 2 y hacer redrive. Pro: no
se pierde ni un pedido y el sistema queda en su estado deseado. Contra: hay que desplegar de urgencia y
verificar que la versión 2 se procesa bien. Escribir un consumidor de compatibilidad que traduzca
versión 2 a versión 1. Pro: no toca el consumidor principal. Contra: añade una pieza permanente para un
problema temporal. La primera es claramente mejor; la segunda solo tiene sentido si el consumidor nuevo
no está listo y los pedidos no pueden esperar. Y el fondo del asunto: el productor no debía haberse
desplegado antes que el consumidor; con version en el cuerpo, el consumidor antiguo podría haber
ignorado los campos nuevos en lugar de fallar, que es justamente por qué la evolución compatible solo
añade campos (07-03).
(d) Reprocesado.
aws sqs start-message-move-task \
--source-arn arn:aws:sqs:eu-west-1:111122223333:mercadofresco-pedidos-fallidos \
--max-number-of-messages-per-second 10 --profile mercadofresco-dev --region eu-west-1El límite importa porque estos 340 pedidos llevan más de 24 horas parados: al procesarlos golpean a la vez la pasarela, Aurora y el ERP, además del tráfico normal del lunes por la mañana. A 10 por segundo tardan 34 segundos y no se nota; de golpe, podrían provocar exactamente la tormenta de reintentos que describe esta lección.
(e) La salvaguarda. Dos, en realidad. La regla de evolución compatible: nunca desplegar un productor con un contrato nuevo antes que sus consumidores, y hacer que los cambios sean solo aditivos para que un consumidor antiguo pueda ignorar lo que no conoce. Y la alarma de la DLQ, que habría avisado el domingo a las 22:09 en lugar del lunes a las 09:00; con ella, alguien habría podido revertir el despliegue en minutos y ningún pedido habría esperado 11 horas.
Solución 3
(a) Por qué no basta reintentar. Los reintentos cubren el fallo de la llamada, no el fallo del
proceso. Si la instancia muere entre el COMMIT de Aurora y el put_events —despliegue, reducción
del ASG, OutOfMemory, un fallo de hardware— no queda nadie que reintente y no hay rastro de que
faltara publicar algo. La probabilidad es baja, pero 1 entre 900 con 240.000 pedidos al mes son 266
pedidos perdidos al mes.
(b) Tabla y transacción.
CREATE TABLE outbox (
id uuid PRIMARY KEY,
agregado text NOT NULL,
tipo_evento text NOT NULL,
carga jsonb NOT NULL,
publicado boolean NOT NULL DEFAULT false,
creado_en timestamptz NOT NULL DEFAULT now(),
publicado_en timestamptz
);
CREATE INDEX idx_outbox_pendientes ON outbox (creado_en) WHERE publicado = false;El índice parcial es importante: la consulta del publicador solo mira filas no publicadas, y un índice
completo crecería con toda la historia. La transacción es la del apartado correspondiente: INSERT en
pedidos e INSERT en outbox dentro del mismo BEGIN/COMMIT.
(c) El publicador. Corre cada segundo, como proceso de fondo en las instancias de trabajadores o
como Lambda disparada por EventBridge Scheduler. Lee 100 filas pendientes con FOR UPDATE SKIP LOCKED,
que permite varios publicadores en paralelo sin que dos cojan la misma fila. Publica y marca
publicado = true. Si muere entre publicar y marcar, la fila se republicará: es un duplicado aceptable.
Conviene usar el id de la fila outbox como identificador de deduplicación aguas abajo, de modo que el
duplicado sea detectable.
(d) Garantía. Al menos una vez, de extremo a extremo. Nunca se pierde un evento —porque está en la transacción— y puede duplicarse —porque el publicador puede morir entre publicar y marcar—. Eso exige que todos los consumidores sean idempotentes, que es precisamente lo que esta lección ha construido. No se puede tener outbox y consumidores no idempotentes: sería cambiar un problema de pérdida por uno de duplicación.
(e) Dos cosas nuevas que operar. Primero, el proceso publicador: hay que monitorizarlo (métrica de filas pendientes, alarma si superan 500 o si la más antigua lleva más de 60 s) porque si se para, los pedidos vuelven a no publicarse y ahora el fallo es silencioso y global en lugar de esporádico. Segundo, la limpieza de la tabla: un borrado diario de las filas publicadas hace más de 7 días, o crecerá sin control como los carritos zombis de 06-02. Y como coste indirecto, la latencia del evento sube de milisegundos al segundo del sondeo, algo irrelevante en este caso porque todo lo que hay detrás es asíncrono.
Conclusión
Este módulo empezó con una tienda que hacía ocho cosas seguidas antes de responder y termina con una
arquitectura donde confirmar un pedido son dos operaciones —cobrar y escribir— y 400 milisegundos.
Entre medias han aparecido las piezas: cola-mercadofresco-pedidos y sus hermanas para desacoplar,
mercadofresco-pedido-confirmado para que un hecho llegue a todos los interesados,
bus-mercadofresco para enrutar por contenido y recibir lo que publica el propio AWS, y
mercadofresco-procesar-pedido para gobernar el proceso largo con su camino de compensación.
Pero lo que de verdad sostiene todo eso no es ninguno de los cuatro servicios: es lo de esta lección.
Que el sistema entrega al menos una vez y que el efecto exactamente una vez lo pone tu consumidor,
con una clave de dominio y una escritura condicional en mercadofresco-idempotencia. Que los reintentos
necesitan jitter o la recuperación provoca la siguiente caída, y que acumularlos en tres capas
multiplica la carga sobre lo que ya está sufriendo. Que un interruptor de circuito compartido en
ElastiCache vale más que treinta segundos de espera contra un servicio caído. Que una DLQ es una
pregunta sin responder y necesita un runbook que alguien pueda seguir a las 3 de la mañana, empezando
por contener y clasificar, y sin redrive antes de corregir. Que el orden global se pide mucho más de
lo que se necesita y cuesta el paralelismo entero. Que escribir en la base de datos y publicar sin
transacción pierde un evento de cada novecientos, y que el outbox o los Streams son las dos únicas
respuestas reales. Y que la contrapresión —dejar que la cola crezca en lugar de aumentar la presión—
es lo que convierte un pico en tiempo en lugar de en una caída.
La arquitectura de MercadoFresco es ya sólida. Y sin embargo, hay un problema del curso que sigue
exactamente igual que el primer día: todo esto se ha desplegado a mano. Las colas se crearon con
aws sqs create-queue escrito en un terminal, las reglas con put-rule, la máquina de estados subiendo
un JSON. Nadie sabe con certeza si el entorno de pruebas tiene la misma configuración que producción.
Luis sigue subiendo cambios por SSH un viernes por la tarde, con el código en el portátil y los dedos
cruzados. No hay pruebas automáticas que digan si un cambio en el consumidor rompe la idempotencia. No
hay forma de volver atrás salvo copiar el fichero anterior, si es que alguien lo guardó. Y el incidente
del ejercicio 2 —un productor desplegado antes que su consumidor— es exactamente el tipo de error que un
proceso de despliegue serio hace imposible.
En el módulo 8, «Herramientas para desarrolladores», empezando por 08-01, «AWS CodeCommit», atacamos el cuarto problema del curso: los despliegues arriesgados. Veremos dónde vive el código y cómo se gobiernan los cambios, cómo se construye y se prueba automáticamente con CodeBuild, cómo se despliega sin cortar el servicio y se revierte solo con CodeDeploy, y cómo se encadena todo en un flujo continuo con CodePipeline, hasta montar un pipeline de extremo a extremo que lleve un cambio de Luis desde su portátil hasta producción sin que nadie tenga que escribir un comando un viernes a las siete de la tarde.
Curso de AWS
Módulo 1: Introducción a AWS
- ¿Qué es AWS?
- Configuración de tu cuenta de AWS
- Infraestructura global de AWS
- Consola de administración de AWS
- AWS CLI y SDKs
Módulo 2: Servicios principales de AWS
Módulo 3: Redes y entrega de contenido
- Amazon VPC
- Grupos de seguridad y listas de control de acceso
- Elastic Load Balancing
- Amazon CloudFront
- Route 53
Módulo 4: Seguridad e identidad
- AWS Identity and Access Management (IAM)
- AWS Key Management Service (KMS)
- Secrets Manager y Parameter Store
- AWS Shield
- AWS WAF
Módulo 5: Monitorización y gestión
- Amazon CloudWatch
- AWS X-Ray y trazabilidad distribuida
- AWS CloudTrail
- AWS Config
- AWS Trusted Advisor
Módulo 6: Bases de datos
- Cómo elegir la base de datos adecuada
- Amazon DynamoDB
- Amazon Aurora
- Amazon Redshift
- Amazon ElastiCache
Módulo 7: Integración de aplicaciones
- Amazon SQS
- Amazon SNS
- Amazon EventBridge
- AWS Step Functions
- Patrones de integración: idempotencia, reintentos y colas de mensajes fallidos
