skip to content

You need to populate a newly added column for roughly 500 million existing rows while the system stays fully online. How do you design, run and supervise that backfill?

level: seniorimportance: should knowfreq 45%

answer

  1. Keyset pagination, never OFFSET
  2. One transaction per batch, cursor persisted
  3. Guard with IS NULL → idempotent, never clobbers live writes
  4. Throttle on replication lag, not a fixed sleep
  5. Verify remaining count before tightening constraints

basics

~20 s

Run bounded batches ordered by primary key, remembering the last key processed, committing each batch, throttling on replication lag and lock waits, and written so re-running is harmless. Never one statement across the whole table.

solid answer

~60 s

Never a single `UPDATE` over the table: one giant transaction holds locks for hours, accumulates enormous undo or row versions, floods replication, and if it fails at 95% you have nothing. The design I use: - **Key-range batches.** Process by primary key ranges — `WHERE id > :last AND id <= :last + N` — not `LIMIT/OFFSET`, whose cost grows as it advances. Each batch is its own transaction; a few thousand rows is a reasonable starting size. - **Resumable.** Persist the last processed key so a restart continues rather than starts over. - **Idempotent.** `WHERE new_col IS NULL` or an equivalent guard, so re-running a batch is a no-op and never clobbers a value dual write already produced. - **Throttled.** Between batches, check replication lag, lock waits and error rate; sleep or shrink the batch when they rise. This is what makes it invisible. - **Observable.** Rows processed, rate, ETA, last key, lag — plus a kill switch that stops it without a deploy. And it must coexist with live traffic already writing new rows correctly, so the backfill only fills gaps. When done, verify by counting the remaining unconverted rows before tightening any constraint.

code

sql · 7 lines
sql
UPDATE person
   SET display_name = concat_ws(' ', first_name, last_name)
 WHERE id > :last_id
   AND id <= :last_id + 5000
   AND display_name IS NULL;

-- caller persists :last_id + 5000, checks replication lag, then repeats

go deeper

for a junior

Know that large data changes are done in batches with commits, not as one enormous statement.

for a middle

Describe keyset batching, per-batch transactions, a persisted cursor, and an idempotent guard so re-runs are safe.

for a senior

Own the operations: lag-driven throttling, adaptive batch size, bloat and vacuum pressure, kill switch and progress metrics, and the race against live dual writes.

for a principal

Treat it as a scheduled load-shedding problem — days of deliberately slow work, thresholds agreed with whoever owns latency, and verification gating the constraint tightening that follows.

## Why a single statement is wrong `UPDATE t SET new_col = f(old_col)` over half a billion rows fails in several ways at once. It runs as one transaction, so it holds row locks for its entire duration, blocking conflicting writers. It generates undo or old row versions for every row, which bloats storage and, under MVCC, holds back garbage collection so unrelated tables grow too. It produces a single vast burst of replication traffic, pushing followers into lag. Its plan is a full table scan whose intermediate state is invisible. And if it fails or is cancelled near the end, the engine rolls the whole thing back — often taking as long again — leaving zero progress. ## Batch selection Order the work by the primary key and advance a cursor: ``` WHERE id > :last_id ORDER BY id LIMIT :batch ``` This is **keyset pagination**. Never use `OFFSET`, whose cost rises linearly as the job advances, so the last batches are the slowest. If the key is a UUID or otherwise unordered, iterate over its ranges anyway, or over a partition or time column that has an index. Batch size is a tuning knob, not a constant. Start small — 1,000 to 5,000 rows — measure the duration of a batch, and target something short, say under a second, so no batch holds locks long and cancelling is immediate. Adapt: shrink when latency or lag rises, grow when the system is idle. ## Idempotence and resumability Assume the job will be killed, redeployed, or run twice concurrently by accident. Two properties defend against that: - **Idempotent batches.** Guard with `AND new_col IS NULL`, or compute a value that is stable for the same input. Re-running then costs a scan and changes nothing. This guard is also what prevents the backfill from overwriting a fresher value written by live dual-write traffic — a real bug when the backfill reads a row, the application updates it, and the backfill then writes a stale derived value back. - **Persisted cursor.** Store the last completed key in a table so a restart resumes. Storing it in memory or in a shell variable means a restart begins again from zero. If the two could genuinely race, do the read and the write in one statement, or re-check the source value in the `WHERE` clause so a concurrent change causes the batch row to be skipped rather than overwritten. ## Throttling — the part that makes it invisible Between batches, sample the health signals that the backfill can damage: - **Replication lag.** The single most important one. Backfills write heavily and lag is the first casualty; pause above a threshold and resume below it. - **Lock waits and long transactions.** If the batch is queueing, back off. - **Application latency and error rate.** If the service degrades, stop. A fixed sleep is a crude version of this and better than nothing; feedback on measured lag is much better, because load varies through the day. ## Bloat and vacuum pressure An `UPDATE` in an MVCC engine writes a new row version and leaves the old one behind, so backfilling a whole table can double its physical size and every index on it before cleanup catches up. Batching with frequent commits lets garbage collection reclaim space as you go, but the job must run slowly enough for it to keep up. Watch table and index size during the run, not only afterwards. ## Operating it Run it as a supervised job — a script or a service task — not as a one-off statement in a console session that dies with the SSH connection. Give it: - structured progress logging and a metric for rows/second, remaining rows and estimated completion; - a kill switch (a flag or a row in a control table) so it can be stopped without a deploy; - a bounded runtime per invocation so it can be scheduled into quiet periods and picked up later. Expect it to take days on 500 million rows if it is throttled properly. That is the correct outcome, not a failure: the point is that nobody notices. ## Finishing Completion is a measured claim, not an assumption. Count remaining rows where the new column is still unset, and re-run the sweep once more after the count reaches zero to catch anything written by a path that was missed. Only then tighten constraints — `NOT NULL`, checks — since those assertions are only valid once every row and every writer complies. And keep the verification query around while dual write continues, so a straggler writer shows up as a non-zero count rather than as a constraint violation later. ## When it is not worth it On a table of a few million rows on a quiet system, a single throttled statement during a low-traffic period may be entirely fine. The batching machinery is proportional to risk: table size, write volume, and how much the business notices a latency spike.

  • Why keyset pagination rather than LIMIT with OFFSET?
    OFFSET makes the engine walk and discard every preceding row, so batch cost grows linearly with how far the job has advanced and the final batches become dramatically slower than the first. Keyset pagination seeks directly into the index at the remembered key, so every batch costs the same. It is also stable under concurrent inserts and deletes, which shift OFFSET-based windows and cause rows to be skipped or processed twice.
  • The backfill and live dual-write traffic can touch the same row. How do you make sure the backfill does not overwrite a fresher value?
    Make the batch statement conditional on the target still being unset — typically AND new_col IS NULL — so a row the application has already populated is skipped entirely. Doing the read and the write in a single statement also removes the read-then-write gap where a concurrent update could slip in. If the derived value depends on a source column that may change, re-check that source value in the WHERE clause so a changed row is left for the next pass.

saying these in an interview costs you the question

  • Running one UPDATE across the entire table and calling it a migration
  • Using LIMIT with OFFSET to page through the table
  • Keeping the progress cursor only in memory, so a restart begins from zero
  • Sleeping a fixed amount instead of reacting to measured replication lag
  • Assuming completion without counting the rows still unconverted
  • Ignoring the table and index bloat that a full-table update creates under MVCC

context