One malformed record. 61,000 duplicate loyalty credits.
A retail platform streams purchase events into a Kinesis Data Stream (streaming log — consumers track their own position via a checkpoint; the stream never deletes on read). A Lambda function (serverless compute) consumes the stream through an event source mapping (polls the shard and invokes Lambda with batches — advances the checkpoint only when the whole batch succeeds) with a batch size of 100. For each record, the function calls the loyalty service and credits points to the customer's account.
The pipeline has run cleanly for a year at 2 million events per day.
On Tuesday at 9:40 AM, a producer deployment ships a bug: one purchase event is published with a null user_id. It lands at position 73 of a batch. The function processes records 1 through 72 — credits 72 customers — then throws on record 73. The invocation fails.
The event source mapping retries. Not record 73 — the batch. Records 1 through 72 are credited again. Record 73 throws again. Retry. Again. The mapping's retry setting is the default: keep retrying until the record ages out of the stream's 24-hour retention.
Support notices at 1:50 PM: four hours, roughly 850 retry cycles, 61,000 duplicate credits, and every purchase event behind the bad record on that shard is stuck waiting. No Lambda alarm fired — the function's errors look like a routine blip at first glance, because the error count is one per retry, not 61,000.
Why did already-processed records run again, and what is the correct fix?
A) Kinesis (streaming log) delivered the same records multiple times during the incident — but a stream is not a delivery service that can duplicate; the consumer reads from a position it controls, and the stream served exactly what was asked for
B) The checkpoint advances per batch, not per record — a thrown error rewinds the consumer to the last committed checkpoint, so every record after it, including the 72 that succeeded, re-executes on every retry until the batch finally succeeds or the poison record expires
C) Two concurrent Lambda invocations processed the same shard and raced — but the event source mapping invokes one batch at a time per shard precisely to preserve ordering; there is no concurrent second reader on the shard
D) The loyalty service retried the credits internally — but it received separate, fully-formed requests carrying no shared idempotency key; from its side these were 850 distinct instructions to credit points, and it executed them correctly
Answer in the comments.
Top comments (5)
The answer is B.
A Kinesis consumer's progress is a single number: the sequence checkpoint. The event source mapping (polls the shard, invokes Lambda with batches, commits the checkpoint on success) reads a batch starting at the checkpoint, invokes the function, and only if the invocation returns successfully does it advance the checkpoint past that batch. Success is all-or-nothing at batch granularity.
So when record 73 throws, the mapping does not know that 72 records succeeded — it only knows the batch failed. It rewinds to the checkpoint and reads the same 100 records again. The 72 loyalty credits were side effects fired into the world; the checkpoint has no memory of them. And because the default retry policy for stream sources is to retry until the data expires from retention, a single poison record turns into a four-hour duplication engine — while also blocking the shard, since ordering guarantees forbid skipping ahead.
This is not an AWS quirk. It is the fundamental contract of checkpoint-based consumption: your commit granularity defines your duplicate window. A Kafka consumer committing offsets per batch has the identical failure mode.
Two fixes, both needed:
Fix 1 — Shrink the duplicate window: enable partial batch responses (
ReportBatchItemFailures). The function reports the sequence number where processing failed; the mapping commits the checkpoint past the 72 successes and retries only from record 73. Pair it with a bounded retry policy — setmaximumRetryAttempts, enablebisectBatchOnFunctionErrorso the mapping halves the batch to isolate the poison record, and configure an on-failure destination (SQS or SNS target that receives the failed record's shard and sequence metadata) so the bad record is skipped, captured, and the shard unblocks.Fix 2 — Make the side effect safe to repeat: send an idempotency key with every credit call —
event_idis sitting right there in the record. Even with perfect checkpointing, at-least-once processing is the contract of every stream consumer. Fix 1 makes duplicates rare; Fix 2 makes them harmless. Only Fix 2 protects the customer's balance.Day 3 of this series was the push side of at-least-once: a platform retry re-running a non-idempotent function. This is the pull side: a checkpoint rewind re-running non-idempotent work. Different mechanism, same law — anything that fires a side effect must survive being executed twice.
A — A Kinesis stream is storage with a cursor, not a delivery pipeline. It cannot "send duplicates" because it does not send — the consumer reads from a position the consumer owns. Every repeated record in this incident was re-requested by the consumer after its own checkpoint rewound.
C — The event source mapping serializes processing per shard to preserve ordering: one batch in flight at a time. (Even with a parallelization factor, records with the same partition key stay in order on one processor.) The re-processing here is sequential retry, not a race.
D — The loyalty service behaved correctly. It received hundreds of independent, valid requests with no idempotency key linking them. Deduplicating them without a key would require it to guess — the platform sent the duplicates, and the platform owns the fix.
Also, it would mean a lot to me if you could support my content and stay in touch 🙏
YouTube: youtube.com/@system-design-lab
LinkedIn: linkedin.com/in/joud-awad/
Medium Blog: joudwawad.medium.com/
Substack: joudawad.substack.com/