skip to content

An EventBridge Pipe reading from a Kinesis Data Streams source stops making progress because one record always fails downstream. Which Pipes source settings get the pipe moving again?

level: seniorimportance: nice to knowfreq 28%

answer

  1. ordered shard, blocked iterator
  2. queues route around it, streams do not
  3. cap by attempts and by age
  4. split the batch to isolate the pill
  5. discarding without a DLQ is data loss

basics

~10 s

For stream sources a pipe retries a failing batch in order, so one bad record blocks its shard. The fix is the source's failure controls: MaximumRetryAttempts, MaximumRecordAgeInSeconds, OnPartialBatchItemFailure set to AUTOMATIC_BISECT, and a DeadLetterConfig.

solid answer

~40 s

Kinesis and DynamoDB Streams are ordered per shard, so a pipe will not advance the iterator past a batch that keeps failing — that is the poison-pill stall. The pipe's source parameters carry the escape hatches: `MaximumRetryAttempts` caps how many times a batch is retried, `MaximumRecordAgeInSeconds` discards records older than a threshold, and `DeadLetterConfig` sends the discarded batch's metadata to an SQS queue or SNS topic so nothing is silently lost. `OnPartialBatchItemFailure` set to `AUTOMATIC_BISECT` makes the pipe split a failing batch in half and retry the halves, isolating the offending record instead of condemning the ninety-nine good ones with it. Also check `BatchSize` and `ParallelizationFactor`: a smaller batch narrows the blast radius, and parallelisation buys throughput but only within a shard's ordering guarantees.

code

json · 13 lines
json
{
  "KinesisStreamParameters": {
    "StartingPosition": "LATEST",
    "BatchSize": 100,
    "MaximumBatchingWindowInSeconds": 5,
    "MaximumRetryAttempts": 3,
    "MaximumRecordAgeInSeconds": 3600,
    "OnPartialBatchItemFailure": "AUTOMATIC_BISECT",
    "DeadLetterConfig": {
      "Arn": "arn:aws:sqs:us-east-1:111122223333:orders-pipe-dlq"
    }
  }
}

go deeper

for a junior

Know that stream sources are ordered, so a record that keeps failing blocks everything behind it in that shard, and that a pipe has settings to cap retries and divert failures.

for a middle

Name the source-level controls and what each bounds — retry attempts, record age, bisection on partial batch failure, and the dead-letter destination — and explain why an SQS source uses the queue's redrive policy instead.

for a senior

Diagnose before tuning: use iterator age to find the stuck shard and Pipes execution logging to locate the failing stage, then choose settings that match the data's value, and alert on the DLQ rather than letting it fill quietly.

for a principal

Own the drop-versus-block policy across pipelines. Decide which data classes may ever be discarded, mandate idempotent targets where they may not, and make the DLQ a monitored, redrivable path rather than a place records go to be forgotten.

## Why the stall happens at all Queue sources and stream sources fail differently, and this is the distinction the question is really testing. With an SQS source, each message is independent. A message that fails becomes visible again, gets retried, and after `maxReceiveCount` attempts the queue's own redrive policy moves it to the dead-letter queue — the queue keeps flowing around it. Notice where that logic lives: **on the queue**, not on the pipe. There is no dead-letter setting in a pipe's SQS source parameters, because SQS already owns that behaviour. With a Kinesis Data Streams or DynamoDB Streams source, records in a shard are strictly ordered and the consumer tracks a position. Advancing past a batch that failed would break ordering and lose data, so by default the pipe retries the batch — and keeps retrying. One record that always throws therefore halts *its entire shard*, while other shards continue. The symptom is a rising iterator age on one shard and a partial outage that looks maddeningly selective. ## The four controls that unblock it All of these live in the pipe's source parameters for stream sources: - **`MaximumRetryAttempts`** — the number of times a failing batch is retried before it is discarded. Without a cap, retries continue until the record ages out of the stream's retention. - **`MaximumRecordAgeInSeconds`** — discard records older than this. It is the time-based twin of the retry cap and is often the more meaningful bound, because what you usually care about is "do not let the pipeline fall more than an hour behind". - **`DeadLetterConfig`** — an SQS queue or SNS topic that receives information about the discarded batch. Set this whenever you set either of the two above, or discarding becomes silent data loss. Note the shape of what arrives: for stream sources the dead-letter record identifies the failed batch and its position rather than replaying the full payloads, so your recovery path usually re-reads the stream at that position. - **`OnPartialBatchItemFailure: AUTOMATIC_BISECT`** — on failure, split the batch in two and retry each half, recursing until the poison record is isolated to a batch of one. The good records in the batch get processed; only the genuinely bad one is discarded to the DLQ. This is the setting that turns "one bad record cost us a hundred good ones" into "one bad record cost us one record". ```json { "KinesisStreamParameters": { "StartingPosition": "LATEST", "BatchSize": 100, "MaximumBatchingWindowInSeconds": 5, "MaximumRetryAttempts": 3, "MaximumRecordAgeInSeconds": 3600, "OnPartialBatchItemFailure": "AUTOMATIC_BISECT", "DeadLetterConfig": { "Arn": "arn:aws:sqs:us-east-1:111122223333:orders-pipe-dlq" } } } ``` ## The tuning knobs around them **`BatchSize` and `MaximumBatchingWindowInSeconds`** trade latency against efficiency. They also decide the blast radius of a failure: with bisection off, batch size *is* the number of records a single poison pill takes down. A very large batch with no bisection is the configuration most likely to produce a dramatic incident. **`ParallelizationFactor`** lets more than one concurrent invocation process a single shard. It raises throughput when your target is the bottleneck, but concurrency within a shard weakens the strict end-to-end ordering you may have been relying on — raise it only when the workload tolerates that. **`StartingPosition`** matters during recovery, not steady state: `TRIM_HORIZON` reprocesses everything still in retention, `LATEST` skips the backlog. Choosing `LATEST` to "clear" a stall silently abandons the backlog. ## Diagnosing before you tune Setting retries and a DLQ stops the bleeding; it does not tell you what bled. Two things help. First, the stream's iterator age metric pinpoints which shard is stuck and when it started. Second, Pipes supports execution logging to CloudWatch Logs, S3 or Data Firehose with levels OFF, ERROR, INFO and TRACE — turning the level up gives per-execution detail on where in the source, enrichment or target sequence the failure occurred, which is usually the difference between blaming the target and discovering that the enrichment is the thing timing out. ## The judgment part Retries plus a DLQ are the right default, but they are a policy statement: you have decided that falling behind is worse than dropping a record. For a clickstream that is obviously correct. For a financial ledger it is not, and the right answer is instead to make the target idempotent and tolerant so nothing needs discarding — with alerting on iterator age so a human sees the stall within minutes. Say which of those two worlds you are in before you quote the settings; that is what separates a senior answer from a configuration recital.

  • Why is there no DeadLetterConfig in a pipe's SQS source parameters?
    Because SQS already owns that behaviour. A message that fails is retried until `maxReceiveCount` is reached, then the queue's redrive policy moves it to the queue's dead-letter queue. The pipe inherits that; adding a second mechanism would duplicate it. Stream sources have no such per-record concept, which is why the setting lives on the pipe there.
  • What is the downside of raising ParallelizationFactor to clear a backlog?
    It runs several concurrent processors over one shard, so records from that shard are no longer processed strictly in order end to end. If downstream logic depends on per-key ordering — applying updates to the same entity, for example — you can apply them out of sequence. Raise it only when the workload is order-tolerant or keys are partitioned another way.
  • How would you tell whether the failure is in the enrichment or the target?
    Enable Pipes execution logging to CloudWatch Logs and raise the level to INFO or TRACE. Each execution record shows how far through the source, enrichment and target sequence it got, so you can see whether the enrichment is timing out or the target is rejecting the payload — instead of inferring it from the target's own metrics.

A queue is a supermarket with many tills — one broken trolley moves to the side and everyone else pays. A stream shard is a single-track railway; the derailed wagon stops every train behind it until you cut it out.

saying these in an interview costs you the question

  • Assumes a bad stream record is skipped automatically
  • Says the queue's redrive policy handles a Kinesis source
  • Sets a retry cap with no dead-letter destination
  • Raises ParallelizationFactor without mentioning ordering
  • Restarts at LATEST to clear a stall, abandoning the backlog

context