aws.crafter.run

Comunicación

Fan-out

También llamado Fan-out / publish-subscribe

Un productor emite un hecho una sola vez y varios consumidores independientes reaccionan a él en paralelo, sin que el productor sepa cuántos son ni qué hacen.

DirectoVerificado el

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.

AntesClienteAPI GatewayprocesarPedidouna sola funciónDynamoDBguardar pedidoSESemail de confirmaciónAPI de inventarioservicio de un tercero120 ms800 ms300 ms123El cliente espera 1 220 ms · si (2) falla, el pedido guardado en (1) queda huérfano
Un solo procesador, tres responsabilidades. El cliente espera 1,22 s (la suma, no el máximo) y el fallo de cualquiera de los tres pasos tira toda la operación — incluido el pedido que ya estaba guardado.

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.

DespuésAPI GatewaycrearPedidopublica y respondeSNStopic pedidos40 msConsumidores independientesSQScola-emailenviarEmailSQScola-inventarioactualizarStockSQScola-analyticsregistrarMétrica123Una copia del mensaje por cola · cada cola con su propia DLQ y su propia concurrencia
Un productor, tres consumidores que no se conocen. El cliente recibe respuesta en 40 ms. Cada rama reintenta, falla y escala por su cuenta: si el email se cae, el inventario ni se entera. Añadir una cuarta rama no toca el código del productor.

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:

Cuándo usarlo

Cómo implementarlo

  1. Nombra el evento como un hecho pasado, no como una orden: PedidoCreado, no EnviarEmail. Si el nombre es una orden, has escrito una llamada RPC con pasos extra.
  2. Crea el topic y publica desde el productor. El productor no sabe —y no debe poder averiguar— quién escucha.
  3. 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.
  4. Suscribe cada cola al topic y activa raw message delivery (ver trampas).
  5. 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.
  6. Añade una DLQ a cada cola con maxReceiveCount entre 3 y 5.
  7. Conecta cada Lambda a su cola con reportBatchItemFailures activado.
  8. 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.

Patrones relacionados

Un enlace sin la relación nombrada es un "ver también". Aquí cada uno dice qué relación tiene y por qué.