A serverless function is triggered by a Kinesis (or Kafka) stream with several shards/partitions. One shard has a record that always throws an exception when processed. What happens to that shard's throughput while the bad record isn't skipped or fixed, and why doesn't it affect the stream's other shards?
answer
- one reader+checkpoint per shard
- ordering only within a shard, not across
- checkpoint can't skip a failed record
- poison pill = stuck shard, others unaffected
- bisect-on-error / max retries / on-failure destination
basics
~20 sThat one shard gets stuck: every batch keeps re-trying starting from the same bad record, so new records behind it never get processed — like a traffic jam on one lane. Other shards each have their own separate reader, so they keep moving fine; the jam doesn't spread.
solid answer
~50 sStream-based triggers assign one poller/consumer per shard (or per partition, in Kafka terms), and each shard's records are delivered and must be checkpointed strictly in order — the function can't skip ahead to a later record until the current one is acknowledged. If a specific record repeatedly throws, the batch containing it keeps failing, the checkpoint never advances past it, and the poller keeps re-delivering the same batch starting at that record, effectively freezing that one shard's processing while everything queued behind it backs up — this is the classic 'poison pill' problem. Because each shard has an independent reader and checkpoint, and shards are the unit of parallelism and ordering, a stuck shard has zero effect on other shards' throughput; they continue advancing on their own separate checkpoints. Mitigations include enabling bisect-on-error (splitting the failing batch to isolate the bad record), setting a maximum retry count with a destination for discarded records, or skipping past unprocessable records after configured retries.
go deeper
Should grasp that one bad record can block processing behind it on the same shard while giving a basic reason (retries keep hitting the same record).
Should explain that each shard has its own reader/checkpoint and that this is why other shards are unaffected, plus name at least one mitigation like retry limits.
Should describe the full mechanism — per-shard ordering guarantee forcing sequential checkpointing, iterator age as the diagnostic signal, and concrete mitigations like bisect-on-error and on-failure destinations.
Should reason about partition-key design and hot-shard risk as a root cause of concentrated blast radius, weigh batch-size and retry-limit trade-offs system-wide, and design defensive per-record error handling as the primary prevention layer rather than relying solely on platform retry mechanics.
## Shards, readers and checkpoints Stream-based event sources like Kinesis Data Streams and Kafka partition their data into multiple independent, ordered logs — 'shards' in Kinesis terminology, 'partitions' in Kafka's. Ordering is guaranteed only within a single shard or partition: records written to the same shard are delivered to consumers in the exact order they were written, but there is no ordering guarantee at all across different shards. - To consume a stream, the serverless platform's poller creates **one independent reader per shard**, each tracking its own position via a **checkpoint** (a sequence number or offset) that records how far that specific shard has been successfully processed. - A function invocation for a given shard is handed a batch of records read starting from that shard's last checkpoint, and — critically — the checkpoint only advances once the invocation reports success for that batch. ## Why ordering forces sequential checkpointing This shard-per-reader, strictly-ordered design exists because many stream use cases genuinely depend on ordering — for example, a stream of 'account balance changed' events for a given user must be applied in the order they happened, or the final balance will be wrong. Kinesis and Kafka guarantee this by routing all records with the same partition key to the same shard and processing that shard sequentially, and the trade-off for that guarantee is that a shard cannot skip ahead: a reader can't process record 105 before record 104 has been durably marked complete, because doing so would violate the ordering contract downstream consumers depend on. ## The poison-pill failure mode This sequential, per-shard checkpointing directly causes the **poison-pill** failure mode. If record 104 has a malformed payload, an unhandled type the deserializer can't parse, or triggers a bug that always throws, the batch containing it fails every time it's retried, because retrying delivers the exact same record again. Since the checkpoint can't advance past an unacknowledged batch, the poller keeps re-delivering starting at that same position indefinitely (subject to whatever retry/backoff policy is configured), and every record behind it in that shard — 105, 106, 107, and onward — sits unprocessed and unread, growing the shard's **iterator age** (the gap between the current time and the timestamp of the oldest unprocessed record) without bound. From an operational standpoint, this shows up as one shard's iterator-age or consumer-lag metric climbing steadily while every other shard's stays flat, which is usually the first diagnostic signal an on-call engineer sees. ## Why the other shards keep moving The isolation between shards is a direct consequence of the same architecture that causes the poison-pill problem: since each shard has its own independent reader and checkpoint, a stall on one shard has no coupling to any other shard's reader, checkpoint, or invocation cadence. A stream with 20 shards experiencing a poison pill on shard 7 continues processing shards 1 through 6 and 8 through 20 at full throughput; only whatever fraction of the total record volume was routed to shard 7 (which depends on the stream's partition-key distribution) is affected. This is a meaningful trade-off in itself: - **High shard-cardinality** limits the blast radius of any single stuck record. - **Hot-key partitioning** (many records routed to the same shard, for example all events for one very active user) means it can concentrate both throughput and risk on that shard disproportionately. ## Mitigations Mitigating the poison-pill problem generally takes one of a few forms, and serious stream-processing deployments typically use more than one. - **Bisect-on-function-error** splits a failing batch in half and retries each half separately, repeating the bisection until the exact offending record is isolated to a batch of one, which limits how many good records get needlessly retried alongside the bad one. - **Setting a maximum retry attempts or maximum record age** lets the platform give up on an unprocessable record after a bounded number of tries and either skip it (advancing the checkpoint past it, accepting data loss for that record) or route it to an on-failure destination such as an SQS queue or S3 bucket for offline inspection, similar in spirit to a dead-letter queue. - **Defensive handler design** — validating and safely handling malformed payloads inside a try/catch per-record within the batch rather than letting one bad record throw an exception that fails the whole batch — is the most effective prevention, since it stops the poison pill from ever fully forming. ## A concrete scenario A concrete scenario: a clickstream-analytics pipeline ingests user-activity events into a 50-shard Kinesis stream partitioned by user ID. A client SDK bug starts emitting a malformed JSON payload for one specific misconfigured mobile app version, and every event from users on that app version routes (via partition-key hashing) to a small handful of shards. Those shards' consumer lag climbs into the hours while the other 45+ shards stay near-zero lag, and downstream dashboards for those specific users go stale — exactly the isolated, per-shard blast radius the architecture produces, and exactly the scenario bisect-on-error and an on-failure destination are designed to contain.
- How does bisect-on-error help isolate a poison-pill record, and what's the cost of using it?It splits a failing batch into two smaller batches and retries each half independently, repeating recursively until the offending record ends up alone in a batch of one, at which point only that single record keeps failing and retrying while its neighbors in the original batch get processed successfully. The cost is extra invocations and latency during the bisection process, since isolating one bad record out of, say, a batch of 100 can take several rounds of halving.
- Why can't you just increase the batch size to make a stream-triggered function more efficient without considering the poison-pill risk?A larger batch means more good records get held hostage behind a single bad one before bisect-on-error or retry limits kick in, and it also means more work gets redone on every retry attempt of a failing batch. It trades average-case throughput efficiency for a larger worst-case blast radius when a poison pill does occur.
- If a team decides to skip a permanently failing record after 3 retries, what's the actual data-loss trade-off they're accepting?They're accepting that the specific event represented by that record is permanently dropped from processing — it will never be retried again once the checkpoint advances past it — in exchange for unblocking every subsequent record on that shard. This is usually paired with routing the skipped record to an on-failure destination so it isn't silently lost, just deferred to manual or offline reprocessing.
Each shard is its own single-lane road with its own toll booth (checkpoint) that won't let traffic through until the car in front is cleared. A stalled car (bad record) backs up only that one lane — every other road (shard) keeps flowing freely because they have completely separate lanes and toll booths.
saying these in an interview costs you the question
- Thinks a stuck shard will eventually block or slow down other shards in the same stream
- Believes skipping a bad record is automatic with no configuration required
- Doesn't know ordering is only guaranteed within a shard/partition, not across the whole stream
- Assumes retrying a failed batch will process a different record than the one that originally failed
- Confuses stream checkpointing with SQS's delete-based acknowledgment model