skip to content

How do you run a ninety-day forecast backfill on the same batch capacity as tonight's retraining graph without starving it?

level: seniorimportance: nice to knowfreq 32%

answer

  1. same steps, historical partitions
  2. one task per date and region
  3. separate queue, lower priority
  4. cap sized from the nightly peak
  5. completeness contract on the window

basics

~20 s

Fan the backfill out into independent per-day, per-region tasks shaped like the nightly steps, run them under a concurrency cap sized from capacity the nightly graph does not need, and keep per-task completion durable so it resumes where it stopped.

solid answer

~40 s

Treat the backfill as the same graph run over historical partitions rather than as a special job. Each `(date, region)` is one task that replaces its own output, so the backfill can be paused, resumed and retried per partition. Submit it through a **separate queue with a lower priority and a concurrency cap**, and size that cap from what the nightly run needs at peak — if the engine has 300 slots and the nightly graph peaks at 200, the backfill gets 100. Record completion per task so an interrupted backfill resumes rather than restarting, and make sure backfill tasks never write the partitions the nightly run is writing at the same time. A corrupt source day then fails one task, not the campaign.

go deeper

for a junior

Understand that a backfill deliberately recomputes past days, and that splitting it into one task per day and region is what lets it be paused, resumed and retried instead of restarted.

for a middle

Explain the throttle: a separate lower-priority queue with a concurrency cap, sized from what the nightly graph needs at its peak rather than from whatever capacity looks free.

for a senior

Work the numbers — task count, average duration, cap, resulting wall clock — and protect the nightly fit's reservation from being starved by thousands of small deferrable tasks.

for a principal

Own the contract around it: which partition ranges each workload may write, what completeness threshold lets training proceed over a window with gaps, and who is told what was excluded from the model that results.

## Shape the backfill like the nightly graph A ninety-day backfill for grocery demand forecasts is not a different pipeline. It is the same steps run over historical partitions, and it should inherit every property the nightly graph has: declared inputs and outputs, an output that is **replaced** rather than appended to, and a durable completion record per task. The unit matters. One task per `(date, region)` gives four properties at once: - **Restart granularity** — a failure at day 88 costs one partition, not 88 days of work. - **Isolation** — a corrupt source file affects one partition's task. - **Throttleability** — a queue of thousands of small tasks can be metered; one enormous task cannot. - **Progress** — completion is observable as a count of finished partitions, not as a running process. The opposite shape — one long task spanning the window — holds a large resource reservation for hours, offers no throttle, and redoes everything on any failure. ## Throttle against shared capacity The nightly graph and the backfill compete for the same batch engine, and the nightly graph has a deadline: the refreshed forecasts must be published before stores open. So the backfill runs as the **deferrable** workload: 1. Submit backfill tasks through a **separate queue** at lower priority, so the scheduler drains nightly work first. 2. Cap backfill **concurrency** at capacity the nightly run does not need at its peak, not at capacity that happens to be idle right now. 3. Make the cap adjustable at runtime, so an on-call engineer can take it to zero without cancelling the campaign. A worked sizing: the engine has 300 task slots and the nightly graph peaks at 200, so the backfill cap is 100. The window is 90 days across 40 regions — 3,600 tasks — and a task averages 3 minutes, so 10,800 task-minutes at 100 concurrent slots is about **108 minutes** of wall clock. That is the number to take to the stakeholder asking when the corrected history will be ready, and it is also the number that tells you whether the campaign is worth running at full width or overnight in two halves. ## Resource-aware scheduling The two workloads do not have the same shape, and the scheduler has to know it. | Workload | Shape | Scheduling need | |---|---|---| | Nightly fit step | One long task, large reservation | Guaranteed slot at its start time; must not be preempted mid-fit | | Nightly aggregations | Wide, short, data-parallel | Burst capacity, finishes quickly | | Backfill tasks | Very wide, short, deferrable | A hard concurrency cap and low priority | A backfill of thousands of small tasks can starve a single large reservation simply by keeping the pool full when the fit step asks for its slots, so the fit's reservation has to be held rather than won in a race. This is why the cap is computed against the nightly peak and not against current idleness. ## Keep the two runs off each other's outputs Both workloads write partitions, and a backfill that reaches into today's partitions while the nightly run is writing them produces a result neither run intended. Two rules: - **Disjoint partition ranges.** The backfill owns historical partitions; the nightly run owns the current one. If the ranges must overlap, the overlap is scheduled, not concurrent. - **Atomic replacement per partition.** Because every task swaps its whole output partition, a reader always sees one complete version — never a half-backfilled day. ## Isolation and a completeness contract Fan-out gives isolation only if the graph is willing to use it. One region's source file for one historical day is malformed: - fail **that task**, not the campaign; - quarantine the partition and record the gap explicitly; - let the step that consumes the window enforce a **completeness contract** — proceed only if the missing fraction is below a stated threshold, and record which partitions were excluded so the resulting model's coverage is known; - retry the quarantined partition separately once the source is fixed, which is cheap because the task replaces exactly its own output. A graph that fails whole on one bad partition turns a fraction of a percent of missing data into a missed refresh. A graph that ignores the gap silently trains on a hole nobody recorded. The contract sits between the two: continue, but only within a stated tolerance, and always leave a record of what was left out.

  • One store-region's source file for a backfilled day is corrupt — should the graph fail?
    Fail that partition's task only; isolation is what per-partition fan-out was for. Quarantine the partition, keep going, and let the step consuming the window enforce a completeness contract: train only if the missing fraction is under a stated threshold, and record what was excluded so the model's coverage is known.
  • Why not run the ninety-day window as one large task?
    It has no restart granularity, so a failure near the end redoes the whole window; it has no throttle, so it cannot be metered against the nightly run; and it holds a large reservation for hours, which is exactly what starves the fit step waiting for its slots.
  • The backfill and the nightly run both want to write the same date's partition. What do you do?
    Keep the ranges disjoint by design — the backfill owns history, the nightly run owns the current partition. Where an overlap is genuinely needed, schedule it rather than run it concurrently, and rely on atomic whole-partition replacement so readers never see a half-written day.

saying these in an interview costs you the question

  • Runs the whole window as one long task
  • Sizes the cap from capacity that is idle right now
  • Raises the backfill's priority so it finishes sooner
  • Lets backfill tasks write partitions the nightly run is writing
  • Fails the whole campaign when one partition's source is corrupt
  • Skips missing partitions silently and trains on the gap