How would you backfill two years of history into a table that a nightly incremental load keeps current?
answer
- do not rewind the nightly job's bookmark
- cut the range into resumable pieces
- a manifest row per chunk
- the source is someone else's production
- newest first if consumers are waiting
basics
~20 sRun the backfill as a separate job with its own state, chunked into bounded half-open windows that each apply idempotently. Leave the nightly load running, throttle the source reads, and track chunk completion so failures resume instead of restarting.
solid answer
~50 sTreat the backfill as a **different job from the nightly load, with different state**. Rewinding the nightly watermark is the classic mistake: it turns one enormous unresumable query into the critical path for today's freshness. Instead, split the range into bounded half-open chunks along a natural partition — a month, a day — and record each chunk's status in a manifest so a failure resumes from the last incomplete chunk rather than the start. Each chunk writes through the same keyed, order-guarded merge the nightly job uses, so overlap between the two is harmless by construction and the backfill can run concurrently. Protect the source: read from a replica, cap concurrency, run off-peak, and watch replica lag rather than only your own job. Sequence chunks newest-first if consumers need recent history first. Finally, budget the target side — many small window writes on a columnar table create small files and merge debt that needs compaction afterwards.
code
sql · 12 lines-- one manifest row per chunk; the driver claims the next pending one
CREATE TABLE backfill_chunks (
source_table text NOT NULL,
window_lo date NOT NULL,
window_hi date NOT NULL, -- half-open: [lo, hi)
status text NOT NULL, -- pending | running | done | failed
attempts int NOT NULL DEFAULT 0,
rows_written bigint,
started_at timestamptz,
finished_at timestamptz,
PRIMARY KEY (source_table, window_lo)
);go deeper
Know that a large backfill is split into smaller date ranges rather than run as one query, and that each piece must be safe to re-run.
Explain why the backfill needs its own state and manifest, how half-open chunk bounds make each piece replayable, and why a keyed write lets it run alongside the nightly load.
Show the operating plan: chunk sizing against data density, concurrency you can throttle at runtime, replica reads, resume-from-manifest on failure, and compaction of the resulting small files.
Own the capacity and consumer-impact call: what the backfill may spend on source and target, what order serves the consumers, whether a shadow table and swap is warranted, and what completeness you publish while it runs.
## Frame the decision first A two-year backfill is a capacity and risk exercise more than a coding one. Before any query runs, get four numbers: how many rows and bytes the range holds, what the source can spare in read throughput without hurting production, what the target write will cost, and what the deadline actually is. Those decide chunk size, concurrency and whether the backfill runs for six hours or six days — and the honest answer is often that a slower backfill is strictly better, because nothing downstream is waiting on the oldest month. ## Rule one: separate job, separate state The most common failure is to rewind the nightly job's watermark to two years ago and let it run. That does three bad things at once: today's freshness now depends on a job that will take many hours, the whole range becomes one query with no resume point, and any failure restarts from the beginning. Run the backfill as its own job with its own state, and leave the nightly load untouched and on schedule. Freshness for recent data is preserved throughout, the two jobs can proceed at different speeds, and you can pause or kill the backfill at any moment without touching the production path. ## Rule two: bounded, replayable chunks Split the range along a natural boundary — usually the same partition column the target uses — into half-open windows: `[2024-01-01, 2024-02-01)`, and so on. Every chunk is a run of the ordinary extract job with explicit `lo` and `hi` parameters. That gives you three properties for free: any chunk can be re-run in isolation, progress is measurable in chunks completed, and the failure domain is one chunk rather than the whole range. Size chunks so a single one finishes in a comfortable unit of time — long enough that per-chunk overhead is negligible, short enough that losing one to a failure costs minutes rather than hours. Uneven data density matters here: a chunking scheme calibrated on a quiet month will produce one monstrous chunk over a peak period, so size by row count where you can rather than by calendar span alone. ## Rule three: a manifest, not a log line Keep a small table with one row per chunk: bounds, status, attempt count, rows read, rows written, start and end time. The job's driver picks the next chunk in `pending` state, marks it `running`, and marks it `done` only after the target write commits. This is what makes the backfill resumable and, just as importantly, auditable. When someone asks in three weeks whether March 2025 was actually loaded, the manifest answers. It also makes a partial backfill an honest object: you can tell consumers exactly which range is complete instead of "most of it, probably". ## Rule four: the write must tolerate overlap The backfill and the nightly job will touch the same keys — a row updated last week may also fall in a chunk you are replaying, and lookback windows widen the overlap further. Two properties make that a non-issue: - **Keyed writes.** Merge on identity, or replace whole partitions where the chunk maps onto partitions exactly. Never append. - **Order guards.** Update only when the incoming change value is at least the stored one, so an old chunk landing after the nightly job cannot overwrite fresher data. With both, the backfill needs no coordination with the nightly job beyond not saturating shared resources. Without them, you need locks, freeze windows and a change-management conversation — which is why the guards are worth building before the backfill, not during it. ## Rule five: protect the source The backfill's reads are the part that can cause an incident in someone else's system. Read from a replica rather than the primary where one exists. Cap concurrency to a fixed small number of chunks in flight, and make that number a knob you can turn down without redeploying. Prefer many bounded queries over one long-running scan, which on MVCC engines holds an old snapshot open and causes version bloat. Schedule the heavy phase outside the source's peak. And monitor the source, not just your job: replica lag, source CPU, lock waits and the source team's own alerts are the signals that tell you to throttle. A backfill that finishes in four hours and pages another team is a failure. ## Rule six: sequence for the consumers Order is a product decision. Newest-first gets the recent, most-queried history usable while the long tail is still loading, and it is the right default for most analytical consumers. Oldest-first suits a rebuild that must be strictly chronological, or where downstream models assume contiguous history from the beginning. Decide deliberately, publish the plan, and expose which range is complete so nobody builds a report on a half-loaded table and quietly reports wrong numbers. Where consumers cannot tolerate a partially loaded table at all, backfill into a shadow table and swap once complete — at the cost of double storage and a period where the shadow needs its own catch-up from the nightly stream. ## Rule seven: budget the target side and clean up after Many small window writes against a columnar or file-based target create many small files, fragmented partitions and merge debt. Plan the compaction or optimise step as part of the backfill, not as a surprise afterwards, and account for its cost in the estimate. On credit-billed warehouses, decide up front what the backfill is allowed to spend and put a ceiling on concurrency to enforce it. ## Close the loop When the last chunk is done, verify rather than declare: reconcile row counts per chunk against the source for a sample of ranges, check that no key exists twice, and confirm the nightly job's watermark was never touched. Then delete or archive the backfill job's state so the next person cannot accidentally resume it.
- Why is rewinding the nightly job's watermark the wrong way to trigger a backfill?It makes today's freshness depend on a multi-hour job, turns the whole range into one query with no resume point, and restarts from the beginning on any failure. It also destroys the record of where the incremental load actually was, so a partial failure leaves you unsure what is loaded.
- How do you decide chunk size for the backfill?Target a chunk that completes in a comfortable unit of time — small enough that a failure costs minutes, large enough that per-chunk overhead is noise. Size by row count rather than calendar span where data density is uneven, or a peak period produces one chunk many times larger than the rest.
- What signals tell you to throttle a running backfill?Signals from the source, not from your job: replica lag, source CPU and lock waits, and the source team's own alerts. Keep in-flight chunk concurrency as a runtime knob you can turn down without a deploy, and treat a backfill that pages another team as a failed backfill regardless of when it finished.
- When is backfilling into a shadow table and swapping worth the extra cost?When consumers cannot tolerate a partially loaded table — a published metric that must not be wrong mid-load, or a downstream contract on completeness. You pay double storage and must catch the shadow up from the ongoing stream before the swap, so reserve it for tables where a half-loaded state is genuinely unacceptable.
saying these in an interview costs you the question
- Rewinding the production watermark to start the backfill
- Running the whole range as one long unresumable query
- Pausing the nightly load for the duration of the backfill
- Ignoring source replica lag while reading at full speed
- Declaring completion without reconciling counts per chunk