Backfills & Reconciliation — Idempotent Jobs, Checksums, Shadow Reads
Backfills copy historical rows while dual-write covers live traffic. Use keyset pages, a persisted resume token, and a lag gate. Prove the copy with chunk checksums and shadow reads. A matching row count with a failing checksum is not permission to cut over.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Overview
A backfill copies or transforms historical rows into the expanded schema or the new store. Dual-write covers live traffic. The job covers everything that already existed. Do it in chunks with keyset pagination, rate-limit it against replica lag and primary CPU, and prove it with checksums plus shadow reads. The job must be idempotent and resumable from a token stored in a control table. Observability is the product: rows per second, lag, and mismatch rate. A silent cron is not a backfill.
Exactly-once delivery of each page is the hard problem. Aim for at-least-once plus an idempotent upsert. Replaying a page must converge, not double-apply.
Flow
- 1
1. Dual-write covers live mutations
- next2. Keyset page from resume token
- 2
2. Keyset page from resume token
- next3. Transform and idempotent upsert
- 3
3. Transform and idempotent upsert
- next4. Sleep if replica lag or CPU high
- 4
4. Sleep if replica lag or CPU high
- next5. Advance token after page commit
- 5
5. Advance token after page commit
- next6. Chunk checksums, then shadow reads
- 6
6. Chunk checksums, then shadow reads
- next7. Repair mismatches before cutover
- 7
7. Repair mismatches before cutover
- next8. Read flag only under mismatch SLO
- 8
8. Read flag only under mismatch SLO
Lesson map
Backfills & Reconciliation — Idempotent Jobs, Checksums, Shadow Reads
Backfills copy historical rows while dual-write covers live traffic. Use keyset pages, a persisted resume token, and a lag gate. Prove the copy with chunk checksums and shadow reads. A matching row count with a failing checksum is not permission to cut over.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB a["1. Dual-write covers live mutations"] b["2. Keyset page from resume token"] c["3. Transform and idempotent upsert"] d["4. Sleep if replica lag or CPU high"] a -->|1. Dual-write covers live mutations| b b -->|2. Keyset page from resume token| c c -->|3. Transform and idempotent upsert| d
One page, then proof
The loop is: read a keyset page, upsert, maybe sleep, advance the token. Checksums start only after the token is at the end.
- 1
Dual-write is already on
Otherwise a row that changes after you copy it stays stale on the target. - 2
Keyset page
WHERE id is greater than the token, ORDER BY id, LIMIT N. No OFFSET. - 3
Idempotent upsert
Same page twice leaves the same target rows. Crash after the write and before the token update is safe to replay. - 4
Lag gate
If replica lag or primary CPU is over the line, return the same token and sleep. Do not advance. - 5
Checksum, then shadow
Per-chunk hash of sorted normalized payloads. Shadow reads sample the served path. Cut over only when the mismatch rate is inside the SLO.
Techniques that matter
Keyset pagination, not OFFSET. WHERE id > token ORDER BY id LIMIT 1000 stays stable under inserts. OFFSET scans and skips, and the skip grows. Concurrent deletes and inserts also shift which rows an OFFSET page means. You want an index on the cursor column.
Rate limit and lag. Pause when replica lag or primary CPU crosses the threshold. That is the same spirit as the online-DDL throttle: the copy is allowed to take calendar time so the primary stays inside its SLO.
Checksums. Hash each chunk of sorted, normalized payloads. Row counts plus a sampled deep compare are the cheap version. A Merkle-style tree of chunks scales when the table is huge. Counts can match while values differ: a boolean flipped, a null that should have been false, a timezone truncated. Do not cut over on counts alone.
Shadow reads. The production response still comes from the old path. An async compare, or a percent of traffic, reads the new path and increments migration_shadow_mismatch_total when they disagree. The mismatch rate is the soak signal for the cutover flag.
Resume tokens. Persist the last successful key in a control table after that page commits. Workers may crash. The next start reads the token and continues. Advancing the token before the upserts commit is how you skip a page forever.
When not to backfill in application code
- A pure SQL
UPDATEon a modest table whose locks you can afford, especially when online DDL already bounds the rewrite. - Historical data that only the warehouse needs.
- History that a CDC stream already copied. Do not build a second copier to feel busy. Connector lag and snapshot tooling stay on the CDC hub.
Prefer set-based SQL when the transform is simple and the table fits a maintenance budget. Prefer workers when you need a throttle, more than one store, or a transform that is not one statement.
A giant UPDATE holds locks, bloats the log, stalls replicas, and cannot pause. Chunked jobs trade calendar time for control. Skipping reconcile means you cut over to wrong values with a green row-count dashboard.
Keyset worker
The source list stands in for the primary. High lag returns the same token and writes nothing. Each upsert is idempotent, so a replayed page converges.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Page size 2 means the three-row table takes two successful calls to finish and a third call that finds an empty page and leaves the token at 3. The throttle call starts from token 0 and must not move.
Chunk checksum and shadow compare
Normalize before you hash: sort by id, then a stable encoding of the fields you claim migrated. A checksum of unsorted rows false-alarms. A checksum that ignores the column you transformed false-passes.
Press Run. Snippets must be self-contained — no network, files, or native modules.
The good chunk is intentionally unsorted. Sorting is part of the checksum, so order of arrival must not matter. The bad chunk differs by case. Counts would still match.
What you page on
| Signal | Why it exists |
|---|---|
backfill_rows_per_second | A job that "is running" at three rows a second will not finish in the window you promised |
| Resume token and percent complete | Crash recovery and a forecast. Percent without a token is a guess |
replica_lag_seconds | The gate. Over the SLO, sleep. Do not heroically catch up by hurting reads |
reconcile_chunk_mismatch_total | Checksums failed. Dig into that chunk. Do not cut over |
shadow_read_mismatch_rate | The served path still disagrees. This is the soak metric |
| Alert on mismatch rate | Job success only means the process exited. It does not mean the values match |
Quarantine the keys in a bad chunk, repair them, and re-hash. A global "rerun everything" is allowed. A global "close the ticket" is not.
Interview Q&A
Why keyset pagination instead of OFFSET?
Answer
OFFSET makes the database scan and skip a growing prefix. Under concurrent inserts and deletes, the same OFFSET is a different set of rows on the next run. A keyset predicate on an indexed cursor is proportional to the page, and the token means the same key after a crash.
How do you resume after a worker crash?
Answer
Persist the last successful token after the page's upserts commit. Idempotent upserts make a replay of that page safe. If you persist the token first, a crash skips the page. If you have no token table, you start at zero and hope the upserts are idempotent — they should be, but you also re-copy the whole table.
Counts match and the checksum fails. What do you do?
Answer
Values or nullability differ. Open the mismatching chunk, quarantine those keys, repair, and re-hash. Do not cut over. A boolean that flipped is invisible to COUNT(*).
When do you backfill in SQL, and when in a worker?
Answer
SQL for a simple single-table transform whose locks you can afford. A worker when you need a throttle, more than one store, a complicated transform, or a job that must run for days and pause. Both still need a way to prove the result.
Can a backfill replace dual-write?
Answer
No. Live writes during the copy race the job. Dual-write, or a CDC stream, covers the moving head. The backfill covers history. Running only the backfill leaves every row touched after its page was copied one version behind.
Why at-least-once plus upsert, instead of exactly-once pages?
Answer
Exactly-once across a crash, a retry, and two stores is a distributed transaction in disguise. An idempotent upsert makes a repeated page converge. That is the guarantee you can operate. The token tells you where to resume. The upsert tells you a resume is safe.
What does a shadow read prove that a checksum does not?
Answer
The checksum proves a chunk matched at hash time. The shadow read proves that the path you are about to serve still matches under live traffic, after dual-write has been running. Both can miss a bug the other catches. You want the rate, over a soak window, not one sample.
Why is a giant UPDATE the wrong default on a hot table?
Answer
It holds locks, writes a large amount of log, and cannot pause when replicas fall behind. Readers on those replicas see lag. You also cannot checkpoint halfway in a way an operator can resume without undoing the statement. Chunks are how you keep a cancel button.
Pitfalls
OFFSETpagination on a table that takes inserts during the job.- Advancing the resume token before the page commits.
- Cutting over because counts match.
- A checksum that does not sort, so identical sets look different.
- A checksum that hashes the wrong columns, so a bad transform looks clean.
- Backfilling history and leaving live writes on only the old store.
- Paging only on process crash, never on mismatch rate.
The job reports 18,402,991 source rows and 18,402,991 target rows. Chunk 44's checksum differs. Shadow mismatch rate is 0.4 percent and the SLO is 0.01 percent. Do you flip read_new? What do you inspect first, and what is still running while you inspect it?
Go deeper
Postgres UPDATE is the set-based tool when a chunked statement is enough. Use The Index, Luke's pagination chapter is the keyset argument with the index you need on the cursor. Martin Fowler's parallel change is why this job sits between expand and contract, not instead of them.
Backfill without reconcile is hope. Ship the mismatch metrics before the cutover flag. Next: rollback, feature flags, and cutover.