El problema
Tienes un endpoint que crea pedidos. Al principio hacía una cosa: guardar el pedido en DynamoDB. Luego el equipo de producto pidió un email de confirmación. Después, descontar del inventario. Después, mandar un evento a analítica.
Cada petición se añadió donde era más fácil: dentro de la misma función Lambda, una detrás de otra.
Hay tres cosas rotas aquí, y solo una es obvia.
La obvia es la latencia: el cliente espera la suma de los tres pasos, no el más lento. Si guardar tarda 120 ms, el email 800 ms y el inventario 300 ms, tu API responde en 1,22 s para hacer un trabajo que en paralelo costaría 800 ms.
La segunda es el acoplamiento de fallos. El servicio de email es de un tercero
y se cae los martes. Cuando se cae, tu función lanza una excepción, el cliente ve
un 500… y el pedido ya está guardado en DynamoDB. Ahora tienes un pedido que
existe, un cliente que cree que falló, y ningún email. Esto no es un bug: es la
consecuencia directa de meter tres operaciones con modos de fallo distintos dentro
de una transacción que no es una transacción.
La tercera es la fricción para cambiar. Añadir un cuarto efecto —avisar al almacén— significa tocar, probar y desplegar la función que crea pedidos. La función que crea pedidos se convierte en el sitio por donde pasa todo el mundo, y en el sitio donde todo el mundo rompe cosas.
La solución
Invierte quién conoce a quién. El productor deja de llamar a nadie: publica un hecho —“se ha creado el pedido 4711”— y termina. Los consumidores se suscriben a ese hecho por su cuenta.
La pieza que hace el reparto es un topic de SNS. Publicas un mensaje y SNS entrega una copia a cada suscripción. La pieza que hace que esto sea robusto es que cada suscripción es una cola SQS, no una Lambda directa.
Esa segunda decisión es la que separa un fan-out que aguanta de uno que parece que funciona:
- La cola absorbe el pico. Si entran 5 000 pedidos en diez segundos, la cola los guarda; las Lambdas los consumen al ritmo que puedan.
- La cola guarda el mensaje mientras el consumidor está caído. Con SNS a Lambda directo, si tu función falla y agota los reintentos, el mensaje se pierde. Con una cola en medio, el mensaje sigue ahí cuando vuelves.
- Cada cola tiene su propia DLQ y su propia concurrencia. El email puede ir a 10 en paralelo y el inventario a 2, sin negociarlo con nadie.
Cuándo usarlo
- Un hecho del dominio dispara varios efectos que no se necesitan entre sí: crear pedido, y de ahí email, inventario, analítica, antifraude.
- Los consumidores tienen modos de fallo o ritmos distintos y no quieres que el lento contagie al rápido.
- Esperas añadir consumidores que hoy no existen, y quieres hacerlo sin tocar al productor.
- Al cliente no le hace falta el resultado de esos efectos para continuar.
Cómo implementarlo
- Nombra el evento como un hecho pasado, no como una orden:
PedidoCreado, noEnviarEmail. Si el nombre es una orden, has escrito una llamada RPC con pasos extra. - Crea el topic y publica desde el productor. El productor no sabe —y no debe poder averiguar— quién escucha.
- Una cola SQS por consumidor. Nunca dos consumidores compartiendo cola: en una cola cada mensaje lo recibe uno solo, así que compartirla convierte tu fan-out en un reparto de carga.
- Suscribe cada cola al topic y activa raw message delivery (ver trampas).
- Da permiso a SNS para escribir en cada cola mediante la política de recursos de la cola, con condición sobre el ARN del topic.
- Añade una DLQ a cada cola con
maxReceiveCountentre 3 y 5. - Conecta cada Lambda a su cola con
reportBatchItemFailuresactivado. - Haz idempotente a cada consumidor. No es opcional: ver trampas.
El código
// Productor: publica el hecho y termina.
import { SNSClient, PublishCommand } from '@aws-sdk/client-sns';
const sns = new SNSClient({});
export const crearPedido = async (evento: APIGatewayProxyEventV2) => {
const pedido = await guardarEnDynamo(JSON.parse(evento.body!));
await sns.send(new PublishCommand({
TopicArn: process.env.TOPIC_PEDIDOS,
Message: JSON.stringify(pedido),
// Los atributos son lo que las filter policies pueden mirar.
// El cuerpo del mensaje no es filtrable salvo que declares
// MessageBody como ambito de filtrado en la suscripcion.
MessageAttributes: {
tipo: { DataType: 'String', StringValue: 'PedidoCreado' },
pais: { DataType: 'String', StringValue: pedido.pais },
},
}));
return { statusCode: 202, body: JSON.stringify({ id: pedido.id }) };
};
// Consumidor: idempotente y con fallo parcial.
import type { SQSEvent, SQSBatchResponse } from 'aws-lambda';
export const enviarEmail = async (evento: SQSEvent): Promise<SQSBatchResponse> => {
const fallidos: { itemIdentifier: string }[] = [];
for (const registro of evento.Records) {
try {
// Con raw message delivery activado, el body ES el mensaje.
// Sin el, seria JSON.parse(registro.body).Message: un string
// dentro de un sobre. Ver trampas.
const pedido = JSON.parse(registro.body);
// Idempotencia: SNS entrega al menos una vez y SQS estandar
// tambien. Los duplicados no son una posibilidad remota,
// son una certeza a largo plazo.
if (await yaProcesado(pedido.id)) continue;
await mandarCorreo(pedido);
await marcarProcesado(pedido.id);
} catch {
// Solo este mensaje vuelve a la cola. Sin esto, un unico
// mensaje venenoso hace reprocesar el lote entero.
fallidos.push({ itemIdentifier: registro.messageId });
}
}
return { batchItemFailures: fallidos };
};
Te va a morder
Coste
Los precios se mueven; los órdenes de magnitud no. Para us-east-1, aproximadamente:
| Concepto | Precio |
|---|---|
| Publicaciones en SNS | ~0,50 USD / millón (primer millón gratis) |
| Entregas de SNS a SQS | sin coste de entrega |
| Peticiones de SQS estándar | ~0,40 USD / millón (primer millón gratis) |
| Invocaciones Lambda | ~0,20 USD / millón, más duración |
La consecuencia práctica es que el reparto es casi gratis y los consumidores no: publicar un millón de eventos cuesta céntimos, pero cada rama que añades multiplica las peticiones de SQS y las invocaciones de Lambda. El coste de un fan-out crece con el número de consumidores, no con el de eventos.
Fuentes
Cuotas de SNS y distinción hard/soft: Amazon SNS endpoints and quotas. Garantías de entrega: Amazon SQS standard queues. Regla de facturación por porciones de 64 KB: SQS pricing. Los importes concretos proceden de las páginas de precios de AWS, que se renderizan dinámicamente y no pudieron citarse literalmente: trátalos como orden de magnitud, no como cifra contractual.