skip to content

In a Pipes and Filters pipeline built on message queues, a filter reads a message, starts processing it, and crashes before acknowledging completion to the queue. What failure modes does this create for the pipeline, and how do you design filters to survive it?

level: seniorimportance: must knowfreq 65%

answer

  1. at-least-once delivery -> redelivery on crash-before-ack
  2. idempotency: same result no matter how many times
  3. poison message -> retry limit + dead-letter queue
  4. dedup via message ID / idempotency key
  5. transactional outbox for state+publish together

basics

~20 s

The queue doesn't know the work finished, so it hands that message to someone else too — meaning it might get processed twice. Filters need to be written so doing the same work twice causes no harm (like setting a value, not adding to it).

solid answer

~40 s

Most queue-based pipes use at-least-once delivery: if a filter crashes before acknowledging, the message is redelivered, which means downstream effects (writes, external calls, emitted messages) can happen more than once. Filters must be idempotent — processing the same message twice produces the same end state as processing it once — typically via deduplication keys, idempotent writes (upserts, conditional writes), or tracking processed message IDs. Separately, a message that repeatedly crashes every filter that tries it ('poison message') needs a retry limit and a dead-letter queue so it doesn't infinitely block the pipeline; and if a filter emits partial side effects before crashing, those need to be safe to repeat or explicitly compensated.

go deeper

for a junior

Understands in plain terms that a crash can cause the same message to be tried again, and that doing the same work twice can cause a problem (like double-charging), without needing precise delivery-semantics vocabulary.

for a middle

Names at-least-once delivery and idempotency explicitly, and can describe deduplication by message ID as a fix.

for a senior

Distinguishes the duplicate-processing failure mode from the poison-message failure mode, proposes retry limits plus dead-letter queues, and recognizes partial-side-effect scenarios (state change vs. downstream publish) as a distinct problem.

for a principal

Brings in patterns like transactional outbox for atomicity across state and messaging, discusses pushing idempotency guarantees into external systems (payment gateway idempotency keys) where that's more robust than filter-local dedup, and reasons about ordered-pipe stalling from poison messages.

## Two steps that can be interrupted This scenario is the central reliability challenge of building Pipes and Filters on real message queues, and it stems from the fact that 'processing a message' and 'acknowledging a message' are two separate steps that can be interrupted between each other. Almost all cloud message queues and brokers (Amazon SQS, Azure Service Bus, RabbitMQ, Kafka consumer groups) default to **at-least-once delivery semantics**: a consumer pulls a message, and the broker only removes it (or advances the consumer's offset) once the consumer explicitly acknowledges completion. If the filter crashes, gets killed, or times out after pulling the message but before acknowledging, the broker has no way to know whether work actually completed, so its only safe option is to make the message visible again for another attempt — either by the same filter instance after restart or by a different instance in the same consumer group. **The result:** the same message gets processed more than once, and any side effect that filter performs can happen more than once too: - writing to a database - calling a downstream API - publishing a message to the next pipe ## Designing every filter to be idempotent The standard, non-negotiable fix is designing every filter to be **idempotent**: processing the same message N times must produce the same observable end state as processing it exactly once. There are several concrete techniques. 1. The cleanest is making the underlying operation naturally idempotent — an upsert (insert-or-replace by message ID) instead of an insert, a conditional write ('set balance to X only if version = Y') instead of a blind increment, or a PUT-style 'set the final state' operation instead of a PATCH-style delta. 2. Where the operation isn't naturally idempotent (e.g., 'charge a credit card' or 'send an email' — calling twice really does have two effects), the filter needs an explicit deduplication layer: extract a unique identifier from the message (or generate one deterministically, such as a hash of its content), check a durable store (a database table, a cache with TTL) for whether that ID was already processed before doing the real work, and record it as processed atomically with the side effect. 3. Many managed queues also offer built-in deduplication windows (e.g., SQS FIFO queue content-based dedup) that catch exact-duplicate redeliveries within a time window, but that's a narrower guarantee than application-level idempotency and doesn't cover all failure timing. ## The poison message A second, distinct failure mode from the same scenario is the **poison message**: a message that, for some reason (malformed data, a bug triggered only by this specific payload, a downstream dependency that always rejects it), causes the filter to crash or fail every single time it's attempted. - Without a limit, at-least-once redelivery means this message gets retried forever, consuming worker capacity and — worse — typically blocking ordered processing behind it if the pipe preserves order (as with a single partition in Kafka or a FIFO queue), stalling every message that arrived after it. - The mitigation is a bounded retry count per message plus a **dead-letter queue (DLQ)**: after N failed attempts, the broker (or the filter's own retry wrapper) moves the message off the main pipe into a separate DLQ where it stops blocking the pipeline, and an operator or an automated remediation process can inspect and reprocess it later. - Nearly every managed queue (SQS redrive policies, Service Bus dead-lettering, Kafka via a manual DLQ topic pattern) supports this natively. ## Partial side effects A third, more subtle failure is **partial side effects**: if a filter's job is 'write to database, then publish to the next pipe' and it crashes between those two steps, a naive retry either re-does the DB write (fine if idempotent) or skips straight to publishing without the DB write ever having happened reliably, depending on how the crash-and-retry logic is structured. This is where patterns like **transactional outbox** (write the DB change and the 'to-publish' record in one local transaction, then a separate relay process publishes from the outbox reliably) become relevant for filters that need to guarantee both a state change and a downstream message happen together, or neither does. ## Pushing the guarantee into the system that matters A concrete real-world instance: an order-fulfillment pipeline's `charge-payment` filter reads an order-created message, calls a payment gateway, and publishes a payment-completed message. If it crashes after the gateway call succeeds but before publishing, at-least-once redelivery causes it to run again — and without idempotency, the customer would be charged twice. Real systems solve this by passing an idempotency key (often the order ID or message ID) to the payment gateway itself, since most major payment processors (Stripe, for instance) support idempotency keys precisely so that a retried charge request with the same key returns the original result instead of creating a second charge — pushing the idempotency guarantee down into the external system that actually matters, rather than trying to solve it only at the filter's own boundary.

  • Why can't you just rely on the queue's built-in deduplication instead of making the filter itself idempotent?
    Built-in dedup (like a FIFO queue's content-based dedup window) typically only catches exact duplicate messages within a limited time window and doesn't help if the filter itself calls an external system twice for what the queue considers one logical delivery, or if the duplicate arrives outside that window. Application-level idempotency (deduplication keys checked against durable state, or naturally idempotent operations) is the only guarantee that covers all the ways redelivery can actually happen.
  • How does a dead-letter queue prevent a poison message from stalling an ordered pipeline?
    Once a message exceeds its retry limit, the broker (or filter's retry wrapper) removes it from the main pipe entirely and routes it to a separate DLQ, which frees the main pipe to continue delivering subsequent messages instead of endlessly retrying the same blocking one. This matters most on ordered pipes (a single Kafka partition, an SQS FIFO queue) where one stuck message would otherwise stall everything behind it.
  • What's the transactional outbox pattern solving here, specifically?
    It solves the 'write to DB and publish a message' atomicity problem: since a local DB transaction and a separate message publish can't be done as one atomic operation across two different systems, the outbox pattern writes both the state change and a record of the message-to-publish in the same local DB transaction, then a separate relay process reliably publishes from that outbox table — guaranteeing the state change and the message either both happen or neither does, closing the crash-between-steps gap.

Like a warehouse worker who picks an item off the shelf but the system only marks it 'shipped' after they scan the shipping label. If they get pulled away after picking but before scanning, the system assumes it's still there and sends someone else to pick the same item — so the process only works safely if 'picking it twice' can't result in shipping two items to the customer.

saying these in an interview costs you the question

  • Assumes at-least-once delivery means messages are never duplicated
  • Proposes 'just don't crash' instead of designing for idempotency
  • No concept of a dead-letter queue or retry limit for permanently-failing messages
  • Thinks idempotency only matters for database writes, not external API calls
  • Confuses message ordering guarantees with delivery guarantees

context