skip to content

A retraining graph's aggregate step is retried after a timeout and one region's daily units double — what went wrong?

level: middleimportance: must knowfreq 68%

answer

  1. at-least-once, not exactly-once
  2. a timeout describes the caller
  3. append adds, replace overwrites
  4. staging per attempt, then atomic swap
  5. side effects need a natural key

basics

~20 s

The step appends. Scheduled steps run at least once, so a retry re-executed work whose rows were already written, and the partition now holds two copies. A step must replace its declared output, not add to it.

solid answer

~40 s

A timeout is a statement about the *observer*, not the worker: the first attempt may have been slow, or may have finished writing and died before its completion record landed. Either way the scheduler offers **at-least-once execution**, so the aggregate step ran twice. Because it appended rows for that region, the second run added a second copy and the daily units doubled. The fix is not fewer retries but a different write shape: each attempt computes into its own staging location and then **publishes atomically over the whole declared output partition**, so two attempts leave the same partition rather than two contributions. Rerunning is then a function of the declared inputs, and the partial output of a killed attempt is an orphan nobody reads.

code

pseudocode · 10 lines
pseudocode
step build_daily_units(run_id, date, region):
    target  = declared_output(date, region)
    staging = temp_location(target, attempt_id)

    discard_if_exists(staging)
    for each sale in declared_input_sales(date, region):
        accumulate(staging, key = (sale.store, sale.item), units = sale.units)

    publish_atomic(staging, target)   # replaces the whole partition
    record_complete(run_id, step = "build_daily_units", partition = (date, region))

go deeper

for a junior

Remember that a retried step runs its body again, so a step that adds rows to what is already there will count the same sales twice. Replacing the output is what makes a repeat harmless.

for a middle

Explain at-least-once execution and why a timeout says nothing about whether the worker finished, then describe the staging-plus-atomic-swap shape and why the swap must cover exactly the partition the task owns.

for a senior

Diagnose the doubled region from run history and partition counts, and extend the argument to the side effects outside the swap — emitted messages, audit rows, external calls — which need a natural key derived from the work.

for a principal

Set it as a platform rule rather than a per-step habit: every step's output is a function of its declared inputs, published atomically, and any step that cannot meet that has to justify the exception before it is scheduled.

## Why a retry is not a clean rerun The nightly graph for grocery demand forecasts aggregates raw sales into daily units per store-item, one task per region. A task exceeds its time limit, the scheduler retries it, and that region's units come back doubled. Nothing malfunctioned. Two facts collided: 1. **A timeout is an observation about the caller.** It means the scheduler stopped waiting, not that the worker stopped working. The first attempt may still be running, or may have written its rows in full and then died before its completion record was durable. 2. **Scheduled execution is at-least-once.** Because any acknowledgement can be lost after the work is done, a scheduler that never loses work must be willing to repeat it. **Exactly-once execution** is not on offer; the *effect* of exactly-once is something the step has to provide. An appending step breaks under exactly that. Attempt one wrote 12 million rows; attempt two wrote them again; the partition holds 24 million and every downstream aggregate for that region is twice what it should be. Worse, the model does not error — it trains on a region whose demand looks to have doubled overnight, and the forecast is quietly wrong for that region only. ## The two write shapes | | Append to the output | Replace the declared output | |---|---|---| | Effect of a second attempt | Adds a second contribution | Leaves the same content | | Effect of a partial attempt | Half a contribution, silently kept | Orphan staging nobody reads | | Safe to retry on timeout | No | Yes, for the declared output | | What a reader sees mid-write | Growing, inconsistent data | Old content, then new, never both | | Recovery after a crash | Manual deletion of the bad slice | Rerun the step | The second shape is what makes a graph restartable at all. The step's output is defined as a **function of its declared inputs** rather than as an accumulation of whatever ran, so repeating the step is a no-op on the final state. ## Building the replacement 1. **Compute into attempt-private staging.** Every attempt writes to a location unique to that attempt, so two concurrent attempts never interleave rows. 2. **Publish atomically over the whole declared output.** A single swap makes the partition go from old content to new content with nothing in between. Concurrent attempts become last-writer-wins on identical content. 3. **Make the scope of the swap equal the scope of the task.** If the task owns `(date, region)`, the swap replaces that partition — never a wider one, or a retry of one region wipes its neighbours. 4. **Treat leftover staging as garbage, not data.** Orphans from killed attempts are cleaned on a schedule; nothing downstream may read them. 5. **Keep the computation deterministic given the declared inputs.** A step that reads the wall clock or samples without a fixed seed produces a different output on the second attempt, and the swap then silently changes results rather than reproducing them. ## What the swap does not cover The atomic publish protects the step's **declared output**. Anything the step does outside it repeats on every attempt: - rows appended to a shared audit or metrics table; - a message emitted to a downstream consumer; - a counter incremented, an alert raised, a notification sent; - a call to an external system that has its own state. Those need their own protection — a **natural key** derived from the work itself, such as the run, the step and the partition, so a repeat collides with the first write instead of adding to it. A key derived from the attempt is worthless: each attempt invents a new one, so nothing ever collides and every retry duplicates. ## How this shows up and how to confirm it Double counting rarely announces itself. The forecast for one region jumps, the evaluation metric for that region degrades, and everything else is green. The confirmation is cheap: count rows or sum units in the suspect partition against a neighbouring region of similar size, then check whether the step's task has more than one attempt recorded for that run. If a doubled partition lines up with a retried task, you have both the symptom and the mechanism. The remediation is the same as the design: make the step replace its partition, then rerun it. Because the rerun is now a function of the declared inputs, it repairs the damaged partition instead of adding a third copy to it.

  • The retry starts while the first attempt is still writing — does the atomic swap still save you?
    For the declared output, yes: each attempt writes its own staging and the publish is one swap, so the loser's work is simply discarded and the partition holds one complete result. It does not save side effects outside the swap — emitted messages or appended audit rows — which still need a natural key.
  • Would exactly-once execution from the scheduler be a better fix?
    It is not available. A worker can finish its write and die before acknowledging, so a scheduler that never loses work must be willing to repeat it. You get the exactly-once *effect* by making the step's output a function of its declared inputs, published atomically, which is cheaper and survives the scheduler being replaced.
  • How do you repair a partition that already holds two copies?
    Rerun the step for that partition once the write shape is a replacement. The rerun recomputes from the declared inputs and swaps over the bad content, so it repairs rather than adds. Then re-run the descendants that consumed the doubled data, since their outputs are wrong too.

A kitchen that re-plates a spoiled order makes a fresh plate and swaps it for the one on the pass. A kitchen that adds another scoop to the plate already sitting there has just served a double portion, and nobody in the dining room can tell which it was.

saying these in an interview costs you the question

  • Thinks the scheduler guarantees each step body executes exactly once
  • Fixes it by raising the timeout so the retry never fires
  • Appends with an attempt id and leaves readers to deduplicate
  • Believes a timed-out attempt cannot have written anything
  • Swaps a partition wider than the task owns
  • Keys deduplication on the attempt instead of the work