skip to content

In a CDC pipeline, what guarantees that two updates to the same row reach the sink in source order?

level: middleimportance: must knowfreq 62%

answer

  1. the guarantee is per row, not global
  2. one reader consumes the log in commit order
  3. how a change is routed decides what stays ordered
  4. round-robin distribution breaks it immediately
  5. hash the primary key; one writer per key

basics

~20 s

Ordering is per key, not global. One capture process reads the log in commit order, events are keyed by the row's primary key so a row's whole history travels one ordered route, and one worker applies that key at a time.

solid answer

~50 s

Three things have to line up. First, **capture**: a single reader consumes the source log in commit order, so the emitted stream is already correctly ordered. Second, **transport**: events are keyed by the source primary key and routed by a hash of that key, so all changes to one row stay on one ordered path — round-robin or size-based distribution destroys this immediately. Third, **apply**: only one worker may be writing a given key at a time, and if the sink loads micro-batches it must collapse each batch to the newest event per key before merging, since a merge with two source rows matching the same target row is nondeterministic or an error on most engines. What you get from this is per-key order only. There is no global ordering across keys or tables once you parallelize, and running two capture processes over the same table gives you no ordering at all.

code

text · 9 lines
text
source order for id=7:
  pos 100  UPDATE tier='silver'
  pos 101  UPDATE tier='gold'

round-robin distribution:
  worker A <- pos 100
  worker B <- pos 101
  B commits first  -> sink tier='gold'
  A commits second -> sink tier='silver'   (stale, and never corrected)

go deeper

for a junior

Remember that change events are keyed by the row's primary key and that this is what keeps one row's updates in order. Know that ordering is promised per row, not across the whole table.

for a middle

Walk the three layers — a single ordered log reader, routing by a hash of the key, one writer per key at the sink — and explain concretely what round-robin distribution does to two updates of the same row.

for a senior

Show the sink-side traps you have hit: several events for one key inside a micro-batch, a thread pool that does not partition by key, and repartitioning while the stream is live. Explain the stale-value-forever asymmetry versus duplicates.

for a principal

Own the scope of the guarantee you publish. Decide deliberately that consumers get per-key order and eventual convergence across tables, and be ready to justify refusing a request for global commit ordering on throughput grounds.

## Ordering is a per-key property The useful guarantee in a change-data-capture pipeline is not "events arrive in the order the database committed them" — that is unobtainable at any throughput worth having — but "for any single row, its changes are applied in the order the source made them." That is enough for a current-state sink to converge on the right value, and it is what every design decision below is protecting. ## Capture: one ordered reader Log-based capture reads the source's transaction log sequentially, so the events it emits for a table are already in commit order and each carries a monotonically increasing source position. This is why splitting capture for one table across two readers — two separate replication streams, or a range-partitioned "parallel CDC" scheme someone invented to go faster — is a correctness bug and not just an optimisation. Two readers have no shared position space, and a row whose key moves between the two ranges has its history split across two independent streams with no way to interleave them. ## Transport: key by the primary key Between capture and sink, the stream is almost always parallelised, and the parallelisation is what puts ordering at risk. The rule is to derive the routing decision from the row's primary key — hash the key, send it to the corresponding lane — so every change to that row follows the same lane and keeps its relative order. Any scheme that distributes by arrival, by batch size, or round-robin across workers will, sooner or later, let the newer of two changes to the same row be written first. That failure is worse than a duplicate: a duplicate converges, a lost ordering leaves a permanently stale value in the sink that nothing later repairs. ```text source order for id=7: UPDATE tier='silver' (pos 100) UPDATE tier='gold' (pos 101) round-robin: worker A takes 100, worker B takes 101 B commits first -> tier='gold' A commits second -> tier='silver' <- stale value wins, forever ``` A useful corollary: the number of lanes is a capacity decision, but changing it rehashes keys, so a key's old and new lanes can be drained concurrently during the change. Pausing the sink, or draining before repartitioning, is the usual mitigation. ## Apply: one writer per key, newest per batch The sink is the last place ordering can be lost. Two common mistakes: - **Concurrent writers on the same key.** If the sink loader hands rows to a thread pool without partitioning by key, two updates to the same row race in the database and the loser's value can land last. Partition sink work by key, or serialise per key. - **A micro-batch containing several events for one key.** A `MERGE` whose source side contains two rows matching the same target row is an error on some engines and nondeterministic on others. Collapse first: rank the batch's events per key by source position and keep the highest. ```sql WITH latest AS ( SELECT *, ROW_NUMBER() OVER (PARTITION BY id ORDER BY src_pos DESC) rn FROM cdc_batch ) MERGE INTO customers t USING (SELECT * FROM latest WHERE rn = 1) s ON t.id = s.id WHEN MATCHED THEN UPDATE SET email = s.email, src_pos = s.src_pos WHEN NOT MATCHED THEN INSERT (id, email, src_pos) VALUES (s.id, s.email, s.src_pos); ``` Collapsing is also a large performance win: a row updated a thousand times inside one batch is written once. ## What per-key ordering does not give you It is worth being explicit about the limits, because interviewers probe them: - **No cross-key ordering.** Two different customers' updates can be applied in either order. Usually fine. - **No cross-table atomicity.** A source transaction that touched two tables becomes independent events on independent lanes, so a consumer can briefly see one table updated and not the other. - **No protection against a *replay* landing out of order.** After a restart, re-delivered events can interleave with newly captured ones depending on how the sink batches. Where that is possible, the durable fix is to store the source position on the sink row and refuse to apply an older one. - **Nothing when the key itself changes.** If the routing key is a mutable business key, the change that alters it is routed by a different value than the row's earlier history, and ordering with that history is gone. Key on an immutable identifier. ## What a good answer sounds like Name the three layers — capture, routing, apply — say that keying by primary key is what buys the guarantee, and immediately state the scope: per key, not global. Then mention the two sink-side traps (concurrent writers on a key, multiple events for a key in one batch). Candidates who only say "CDC preserves order" have not operated one.

  • Your sink loads micro-batches and one batch contains three changes to the same row. What must it do before merging?
    Rank the batch's events per key by source position and keep only the newest, then merge that. A merge whose source side matches one target row more than once is an error on some engines and nondeterministic on others. Collapsing is also faster: a hot row updated a thousand times in a batch is written once.
  • Someone proposes running two capture processes over the same large table to double throughput. What do you say?
    That it removes the ordering guarantee. Two readers have independent positions and no way to interleave, so a row whose key falls near a boundary — or moves across one — can have its history split and applied in the wrong order. Scale by parallelising downstream of a single ordered capture, keyed by primary key, not by splitting the capture itself.

saying these in an interview costs you the question

  • Claims CDC preserves the source's global commit order end to end
  • Distributes a table's change events across workers round-robin
  • Uses the event timestamp or transaction id as the routing key
  • Runs two capture processes over one table to go faster
  • Merges a batch that holds several events for the same key

context