aws.crafter.run

Fiabilidad

Outbox con DynamoDB Streams

También llamado Transactional outbox

Guardar el dato y el evento que lo anuncia en la misma transacción, de modo que sea imposible que uno exista sin el otro, y publicar después leyendo del cambio ya escrito.

DelicadoVerificado el

El problema

Tu función crea un pedido: lo guarda en DynamoDB y publica PedidoCreado para que el resto del sistema reaccione. Dos líneas de código, en ese orden.

Dos escrituras, dos destinoscrearPedidoDynamoDBel pedido se guardaSNSel evento no sale1 · PutItem2 · PublishEl pedido existe · nadie se ha enterado · el sistema cree que todo fue bien
Dos escrituras a dos sistemas no son una transacción. Da igual el orden: siempre hay un instante en el que una ha ocurrido y la otra no, y si el proceso muere justo ahí, nadie lo va a arreglar.

Entre esas dos líneas hay un instante en el que el pedido ya existe y el evento todavía no se ha publicado. Si el proceso muere ahí —la Lambda se queda sin tiempo, SNS devuelve un error, la red se cae— tienes un pedido que existe y del que nadie se ha enterado. El almacén no lo reservó, el email no salió, la analítica no lo cuenta.

Lo peor es que el sistema cree que todo fue bien: no hay excepción que capturar más tarde, ni cola donde el mensaje esté esperando. Simplemente falta.

Y no se arregla cambiando el orden. Si publicas primero y guardas después, el fallo produce lo contrario: un evento que anuncia un pedido que no existe, y consumidores que van a buscarlo y no lo encuentran. Has cambiado datos perdidos por datos fantasma.

La solución

Deja de enviar el evento y empieza a guardarlo. En la misma transacción que el dato.

Una escritura, un destinoTodo o nadacrearPedidoDynamoDBpedido + outboxpublicarOutboxEventBridgeTransactWriteItemsDynamoDB StreamsLos registros del stream viven 24 horas: ese es tu plazo para arreglar al consumidorStreams NO conserva la atomicidad de la transacción — ver trampas
El evento se guarda como un dato, no se envía como un mensaje. O se escriben los dos elementos o no se escribe ninguno, así que un pedido guardado siempre tiene su evento esperando. Publicar deja de ser algo que puede fallar a medias y pasa a ser un reintento.

La función escribe dos elementos con TransactWriteItems: el pedido y un elemento de outbox que contiene el evento a publicar. Es una operación de todo o nada, así que el estado “pedido sin evento” deja de existir.

A partir de ahí, publicar ya no es una operación que pueda fallar a medias: es un trabajo pendiente que está escrito en la base de datos. DynamoDB Streams entrega ese elemento a una segunda función, cuya única responsabilidad es publicarlo en EventBridge o SNS y marcarlo como enviado. Si falla, se reintenta; si la función tiene un bug, el registro sigue en el stream.

La ganancia real es de tipo, no de grado: has convertido un fallo silencioso en un reintento.

Cuándo usarlo

Cuándo NO usarlo

Cómo implementarlo

  1. Mismo elemento de tabla, prefijo propio. El evento se guarda como PK = PEDIDO#1001, SK = OUTBOX#<ulid>, así que viaja en la misma colección que el dato.
  2. Escribe los dos con TransactWriteItems. Hasta 100 acciones, hasta 4 MB agregados, misma cuenta y misma región.
  3. Activa el stream con NEW_IMAGE — necesitas el contenido del evento, no solo la clave.
  4. Filtra en el event source mapping para que la función solo despierte con los INSERT cuyo SK empieza por OUTBOX#. Sin filtro, cada cambio de cualquier elemento invoca la función.
  5. Arranca en TRIM_HORIZON, no en LATEST (ver trampas).
  6. Publica y marca. Tras publicar, borra el elemento de outbox o márcalo con un TTL corto.
  7. Activa ReportBatchItemFailures y pon un destino en caso de fallo.

El código

import { DynamoDBDocumentClient, TransactWriteCommand } from '@aws-sdk/lib-dynamodb';
import { EventBridgeClient, PutEventsCommand } from '@aws-sdk/client-eventbridge';

// ── 1. Escribir el dato y el evento como una sola operacion ────────
export const crearPedido = async (datos: Pedido) => {
  const evento = {
    tipo: 'PedidoCreado',
    id: datos.id,
    ocurrido: new Date().toISOString(),
    payload: { total: datos.total, cliente: datos.cliente },
  };

  await ddb.send(new TransactWriteCommand({
    TransactItems: [
      { Put: { TableName: 'app', Item: { PK: `PEDIDO#${datos.id}`, SK: 'META', ...datos } } },
      // El evento es un dato mas. Si esta escritura no ocurre,
      // la del pedido tampoco: es todo o nada.
      { Put: { TableName: 'app', Item: { PK: `PEDIDO#${datos.id}`, SK: `OUTBOX#${ulid()}`, evento } } },
    ],
  }));
};
// ── 2. Publicar lo que ya esta escrito ─────────────────────────────
export const publicarOutbox = async (e: DynamoDBStreamEvent): Promise<SQSBatchResponse> => {
  const fallidos: { itemIdentifier: string }[] = [];

  for (const r of e.Records) {
    // El filtro del event source mapping ya deberia dejar pasar solo
    // los INSERT de OUTBOX#, pero el codigo no se fia del despliegue.
    if (r.eventName !== 'INSERT') continue;
    const item = unmarshall(r.dynamodb!.NewImage!);
    if (!item.SK?.startsWith('OUTBOX#')) continue;

    try {
      await eb.send(new PutEventsCommand({
        Entries: [{
          EventBusName: 'app',
          Source: 'pedidos',
          DetailType: item.evento.tipo,
          Detail: JSON.stringify(item.evento),
        }],
      }));
      await borrarOutbox(item.PK, item.SK);
    } catch {
      fallidos.push({ itemIdentifier: r.eventID! });
    }
  }

  return { batchItemFailures: fallidos };
};

Te va a morder

Coste

Concepto Efecto
Escritura transaccional la capacidad de una escritura normal, por elemento
Capacidad de una transacción cancelada se consume igual
DynamoDB Streams por lectura de registros
Lambda de publicación invocaciones + duración
Publicación en EventBridge ~1,00 USD / millón

El coste dominante suele ser el doble de capacidad de escritura, no la Lambda extra. Escribir pedido y outbox cuesta cuatro unidades de escritura donde antes gastabas una; a cambio, dejas de perder eventos.

Un consuelo: habilitar transacciones no tiene coste adicional — pagas solo las lecturas y escrituras que las componen.

Fuentes

Límites y semántica de TransactWriteItems, consumo de capacidad doble, consumo en transacciones canceladas, y —lo más importante— que los registros de stream de una misma transacción pueden aparecer intercalados y sin garantías de orden: Amazon DynamoDB Transactions: How it works. Sondeo, ventana de agrupación, ParallelizationFactor, orden por elemento, entrega al menos una vez, lectores simultáneos por shard y el aviso sobre LATEST frente a TRIM_HORIZON: Using AWS Lambda with Amazon DynamoDB. Retención de 24 horas de los registros de stream: Core components of Amazon DynamoDB.

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