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.
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.
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
- Un cambio de estado tiene que anunciarse, y perder el aviso es un incidente: pedidos, pagos, altas, cambios de permisos.
- Hay varios consumidores que dependen de enterarse y no puedes ir preguntándoles si les llegó.
- Estás construyendo una saga coreografiada, donde cada servicio reacciona al evento del anterior.
Cuándo NO usarlo
Cómo implementarlo
- 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. - Escribe los dos con
TransactWriteItems. Hasta 100 acciones, hasta 4 MB agregados, misma cuenta y misma región. - Activa el stream con
NEW_IMAGE— necesitas el contenido del evento, no solo la clave. - Filtra en el event source mapping para que la función solo despierte con
los
INSERTcuyoSKempieza porOUTBOX#. Sin filtro, cada cambio de cualquier elemento invoca la función. - Arranca en
TRIM_HORIZON, no enLATEST(ver trampas). - Publica y marca. Tras publicar, borra el elemento de outbox o márcalo con un TTL corto.
- Activa
ReportBatchItemFailuresy 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 | 2× 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.