skip to content

Retrying your dispatcher pipeline duplicates the attempt rows written before the carrier call, so how do you restructure it so only the carrier call repeats?

level: seniorimportance: should knowfreq 46%

answer

  1. retry re-runs its own upstream
  2. scope follows placement
  3. wrap the risky call as its own source
  4. run-once effects stay outside
  5. inner retry keeps the outer subscription alive

basics

~20 s

Make the carrier call its own per-item source and attach the retry to that inner source, so resubscription re-enters only the call. The recording stays in the outer pipeline, outside the retried scope, and runs once.

solid answer

~50 s

A retry re-runs whatever sits above it **inside the stream it is attached to** — so the retried scope is a placement decision, not a fixed property of the pipeline. Wrap the carrier call as a source of its own for one message, attach the retry to that inner source, and compose it into the outer pipeline through the stage that subscribes to a per-item source. Now a carrier failure resubscribes only to the carrier call. The recipient load and the attempt record sit in the outer pipeline and run once. The narrower scope buys a second thing that matters more in production: the failure never reaches the outer subscription, so it is contained to one message instead of tearing down the pipeline and everything else in flight. If an effect genuinely has to live inside the retried scope, key it so repeats converge on one record instead of appending another.

code

pseudocode · 10 lines
pseudocode
// retry attached to the outer pipeline
pipeline = messages
    .map(compose)
    .effect(record_attempt)        // writes one record
    .flat_map(carrier_call)        // may fail
    .retry(max_attempts = 3)

// a carrier failure resubscribes to `messages`:
//   record_attempt writes a SECOND record,
//   and every other message in flight is cancelled

go deeper

for a junior

The takeaway is that where a retry is attached decides what it repeats. A retry around one call repeats that call; a retry around the whole pipeline repeats the whole pipeline.

for a middle

Be able to restructure the pipeline: put the risky call in its own per-item source, attach the retry there, and compose it back through the stage that subscribes to such a source. Explain which stages then run once.

for a senior

Bring the containment argument, not just the duplicate one: an outer failure cancels everything in flight and rebuilds the subscription, while an inner retry keeps the pipeline alive. Name the memory and concurrency cost you took on.

for a principal

Set the placement rule for the codebase and the condition for an effect to sit inside a retried scope at all, then decide where the record of contained per-item failures goes so that quiet containment does not become invisible loss.

## The retried scope is a placement decision The rule is narrow and exact: a retry re-subscribes to the stream it is attached to. Attach it to the outer pipeline and the retried scope is the whole pipeline from its source. Attach it to a small inner source that handles one item, and the retried scope is just that source. Nothing else about the pipeline changes. This is why `retry duplicates my writes` is almost never a reason to abandon retrying — it is a reason to move the retry down to the smallest stream that contains the thing you actually want repeated. ## Two placements, two very different pipelines | | Retry on the outer pipeline | Retry on the inner per-item source | |---|---|---| | What re-runs | recipient load, composition, attempt record, carrier call | the carrier call only | | Duplicate effects | one extra attempt record per retry | none | | Other items in flight | cancelled with the subscription | unaffected, still flowing | | The outer subscription | released and re-established | never released | | Retry state | one scope for the whole pipeline | one per item currently retrying | The third and fourth rows are the ones engineers underestimate. A failure that reaches the outer level is terminal for that subscription: the pipeline is torn down, every other message being processed at that moment is cancelled, and the retry then builds a new pipeline from scratch. Containing the failure to the inner source turns a pipeline-wide event into a per-message one. ## Classify, then place Walk the stages and sort them before deciding where the retry sits: - **Must run exactly once** — the attempt record, anything that notifies, publishes or charges. These belong outside the retried scope, or inside it only once keyed. - **Safe to repeat as-is** — pure composition, reads whose extra load you can afford. - **The stage you are retrying** — the carrier call. The scope should be built around this one and as little else as possible. The retry then goes at the smallest boundary that encloses the third group and none of the first. ## When the effect cannot move out Sometimes the effect genuinely has to be inside — you want a record of each attempt, including the failed ones, written by the thing that made the attempt. Then make repetition converge instead of accumulate: 1. Give the write a **key** that is stable for this logical send and this attempt number, so the same attempt written twice addresses the same record. 2. Make the write **converge on that key** — write-or-update the record rather than appending a new one. 3. Make sure the key is derived from the value, not generated fresh inside the retried scope; a key generated on each subscription is a new key every attempt, which is the same duplicate with extra steps. ## The tempting wrong fix Moving every effect below the retry point looks like a fix and quietly changes the behaviour: the attempt record now only exists for sends that eventually succeeded, and the failed attempts — the ones the record was for — leave no trace. Decide whether you are moving an effect because it does not belong in the retried scope, or because you have not worked out how to key it. Those are different decisions with different consequences. ## What the narrower scope costs The inner placement is not free, and a senior answer says so: - **Assembly is more complex.** There is now a per-item source with its own decision inside the outer pipeline, and two retry-shaped things to reason about instead of one. - **Retry state multiplies.** Each item currently retrying holds its own attempt count, its own timing and the item itself in memory. Many items failing at once means many live retry scopes, not one. - **Concurrency and retry load interact.** The concurrency limit of the stage that runs those per-item sources now also caps how many retries can be in flight at once, which makes that limit a load-control knob it was not before. - **Failures stop being loud.** Contained per item, a failure no longer takes the pipeline down — which is the goal, and which also means nothing forces anyone to notice it. The containment has to be paired with the failure being recorded somewhere. The summary an interviewer is listening for: the retry did not duplicate anything the author did not put inside its scope, and the scope is chosen by where the retry is attached.

  • What happens to the other messages in flight when a failure reaches the outer pipeline instead of an inner retried source?
    The failure is terminal for that subscription, so the pipeline is torn down: every message being processed at that moment is cancelled, and the retry then starts a new subscription from the source. An inner retry contains the failure to the one message that hit it.
  • The attempt record must stay inside the retried scope. How do you stop resubscription multiplying rows?
    Key the write by something stable for this logical send and attempt number, and make it converge on that key — update rather than append. The key must come from the value itself; one generated inside the retried scope is regenerated on every attempt and defeats the point.
  • Why not simply move every effect below the retry point?
    Because it changes what is recorded. Effects moved downstream only happen for work that eventually succeeded, so the failed attempts leave no trace — and the attempt record usually existed precisely to capture those.
  • What new cost does the per-item retry introduce under a burst of failures?
    Each retrying item holds its own attempt count, timing and payload while it waits, so many simultaneous failures mean many live retry scopes rather than one. The concurrency limit of the stage running those inner sources becomes the cap on concurrent retries.

saying these in an interview costs you the question

  • Calls duplicate writes an unavoidable cost of retrying at all.
  • Thinks a retry only repeats the stage immediately before it.
  • Moves every effect downstream and loses the record of failed attempts.
  • Expects other in-flight items to survive a failure at the outer level.
  • Assumes per-item retry costs nothing in memory or concurrency.
  • Keys the converging write with a value generated inside the retried scope.