El problema
Tu función consume una cola y procesa diez mensajes por invocación. Es lo sensato: una invocación para diez mensajes cuesta la décima parte que diez invocaciones.
Hasta que uno de los diez lleva un campo que tu código no espera y lanza una excepción.
batchItemFailures es una lista de mensajes que vuelven. En Kinesis y DynamoDB Streams es un checkpoint: Lambda toma el número de secuencia más bajo de la lista y reintenta desde ahí en adelante.Lo que ocurre por defecto está documentado sin ambigüedad: “si tu función encuentra un error al procesar un lote, todos los mensajes de ese lote vuelven a ser visibles en la cola”. Los nueve que funcionaron perfectamente vuelven a procesarse.
Y no una vez: en cada vuelta. El mensaje venenoso sigue ahí, así que el ciclo
se repite hasta agotar maxReceiveCount o la retención. Cada vuelta son nueve
efectos duplicados —nueve emails, nueve cobros— más una invocación que pagas.
En streams —Kinesis y DynamoDB Streams— es peor, porque ahí el orden importa y AWS lo protege: “para asegurar el procesamiento en orden, el event source mapping pausa el procesamiento del shard afectado hasta que el error se resuelve”. Un solo registro malo bloquea el shard entero.
La solución
Decirle a Lambda qué falló exactamente, en vez de dejar que lo deduzca de una excepción.
Se activa incluyendo ReportBatchItemFailures en la lista FunctionResponseTypes
del event source mapping, y devolviendo desde el handler:
{ "batchItemFailures": [ { "itemIdentifier": "<id>" } ] }
La documentación es tajante sobre el primer paso, y es donde falla mucha gente:
“aunque tu código devuelva respuestas de fallo parcial, Lambda no las procesa
salvo que ReportBatchItemFailures esté explícitamente activado”. Sin esa
casilla, tu return es un objeto que nadie lee.
Y ahora lo importante: no significa lo mismo en los dos sitios
En SQS, batchItemFailures es lo que parece: la lista de mensajes que vuelven
a la cola. El resto se borra.
En streams, es un punto de control. Cita literal: “si el array
batchItemFailures contiene varios elementos, Lambda usa el registro con el
número de secuencia más bajo como checkpoint. Después reintenta todos los
registros a partir de ese punto”.
Es decir: en un stream no puedes saltarte el registro 4 y quedarte con el 5 al 10. Devolver el 4 significa reintentar del 4 al 10. Tiene sentido —el orden es la razón de ser de un stream— pero contradice frontalmente la intuición que uno trae de SQS, y es el origen de duplicados que nadie se explica.
Cuándo usarlo
- Siempre que uses lotes. No hay un caso razonable en el que quieras que nueve mensajes buenos se reprocesen por culpa de uno.
- En cuanto el consumidor tenga efectos observables: cobrar, enviar, escribir.
- En streams, además, para que un registro venenoso no bloquee el shard.
Cuándo NO usarlo
Cómo implementarlo
- Activa
ReportBatchItemFailuresen el event source mapping. Sin esto, lo demás no sirve de nada. - Captura por registro, no por lote. Un
tryalrededor del bucle entero es exactamente lo que estás intentando evitar. - Devuelve el identificador correcto:
messageIden SQS, número de secuencia en streams. - En streams, devuelve el primero que falló y corta. Como el resto se va a reintentar igualmente, seguir procesando es trabajo que se tirará.
- Ajusta el tamaño de lote a la duración de tu función, y recuerda que hay un techo de payload de 6 MB que no se puede tocar.
- Idempotencia, siempre: la entrega es al menos una vez y los reintentos son parte del diseño.
- DLQ detrás, para que el venenoso acabe saliendo del bucle.
El código
// ── SQS: batchItemFailures es una lista de mensajes que vuelven ─────
export const consumir = async (e: SQSEvent): Promise<SQSBatchResponse> => {
const fallidos: { itemIdentifier: string }[] = [];
for (const r of e.Records) {
try {
await procesar(JSON.parse(r.body));
} catch {
// Solo este vuelve a la cola. Los demas se borran.
fallidos.push({ itemIdentifier: r.messageId });
}
}
return { batchItemFailures: fallidos };
};
// ── Streams: batchItemFailures es un CHECKPOINT ─────────────────────
// Devolver el registro 4 significa reintentar del 4 al 10. Por eso, en
// cuanto uno falla, se devuelve y se corta: seguir procesando los
// posteriores es trabajo que se va a tirar.
export const consumirStream = async (
e: DynamoDBStreamEvent,
): Promise<DynamoDBBatchResponse> => {
for (const r of e.Records) {
try {
await procesar(r);
} catch {
return { batchItemFailures: [{ itemIdentifier: r.dynamodb!.SequenceNumber! }] };
}
}
// Lista vacia = lote completo correcto.
return { batchItemFailures: [] };
};
Te va a morder
Coste
| Configuración | Efecto |
|---|---|
BatchSize: 1 |
una invocación por mensaje — el más caro |
BatchSize: 10 (por defecto en SQS) |
una invocación por cada diez |
| Sin fallo parcial, con un venenoso | × reintentos, sobre el lote entero |
| Con fallo parcial | el lote bueno se paga una vez |
El batching es de las optimizaciones más rentables que existen en Lambda: pasar de 1 a 10 divide entre diez las invocaciones y buena parte del arranque amortizado. Lo que casi nadie calcula es el coste del camino de fallo: sin fallo parcial, un solo mensaje malformado en un lote de diez multiplica por el número de reintentos el coste de procesar los otros nueve — y añade sus efectos duplicados, que suelen costar más que la factura.
Dicho de otro modo: el batching ahorra en el camino feliz y el fallo parcial evita que el camino infeliz se lo coma.
Fuentes
Comportamiento por defecto ante un error, pausa del shard para preservar el
orden, entrega al menos una vez, ventana de agrupación y su irreversibilidad en
las fuentes de 500 ms, techo de payload de 6 MB, y que los throttles no cuentan
como reintentos:
How Lambda processes records from stream and queue-based event sources.
Activación mediante FunctionResponseTypes, sintaxis de la respuesta, semántica
de checkpoint por número de secuencia más bajo, condiciones exactas de éxito y
de fallo completo, y bisección del lote:
Configuring partial batch response with DynamoDB and Lambda.
Comportamiento de los mensajes y la espera de hasta 20 segundos con colas de
poco tráfico:
Using Lambda with Amazon SQS.