Data engineering
Data pipelines: source versions, replay and reconciliation
By the BackendDrills editorial team · Published and checked October 6, 2026 · 7-minute read
Before you start: Transactions and at-least-once event delivery. Find this reading in a study path →
A healthy consumer does not demonstrate that a projection represents the source correctly. Define identity, source ordering, replay behavior and reconciliation before using derived data for an important decision. Connector details depend on configuration and version.
Arrival time is not the business version
An original customer-directory fixture emits version 17 with an updated address, followed by version 18 removing that address. During retries, a downstream stage receives 18 before 17. Applying each event in arrival order resurrects stale data. Decide which source version orders changes for the same entity and make an obsolete update unable to replace a newer one.
| Input | Contract to specify |
|---|---|
| Same entity and version again | No additional logical effect. |
| Older version after newer version | Does not overwrite newer state. |
| Deletion followed by an old update | Deletion ordering remains protected. |
| Snapshot mixed with live changes | Transition follows the connector's documented model. |
Debezium's PostgreSQL connector exposes change records and source metadata; snapshot, delete and restart behavior require reading its configuration-specific contract. Do not assume a business timestamp is a sufficient total order or that one stream guarantees ordering after arbitrary downstream repartitioning. Write which metadata the consumer uses and what happens when it is missing.
Replay needs an acceptance rule
A pipeline may replay after restart or checkpoint recovery. Store the applied version or operation identity at the boundary that changes the projection. If processing also sends an email or calls another system, that external effect needs its own duplicate and recovery contract. Updating a checkpoint and updating a destination are distinct operations unless your design explicitly coordinates them.
Measure completeness and correctness
Queue lag is useful but incomplete. A pipeline can report low lag after discarding malformed records or applying stale updates. For a fixture of 20 entities, compare source and projection IDs, versions, deletion state and one business total. Intentionally remove one record and confirm that reconciliation detects it even when the queue is empty.
- Prepare duplicate, reordered and deleted-entity inputs.
- Interrupt a consumer between destination update and checkpoint advancement.
- Restart and replay the same bounded fixture.
- Compare authoritative source state with the rebuilt projection.
- Record missing, extra, stale and conflicting results separately.
Finish with a recovery note naming the replay range, source authority, stop condition and post-replay comparison. Define freshness and completeness targets independently. A lead can use this note to assign repair ownership; an architect can use it to challenge whether a derived store is appropriate for a correctness-sensitive operation.
Check the underlying behavior
Original illustrative examples, prepared with AI assistance and checked against the linked primary documentation. No customer incident or vendor endorsement is claimed. Our editorial approach.
Practice the next decision
Try a complete free backend drill: inspect evidence, make three decisions and review the reasoning. No account or card required.
Try a free incident drill →Selected practice from the study paths
- CDC updates arrive out of order · Full edition drill
- The inventory cannot explain its balance · Full edition drill
Free readings need no account. Full edition drills require verified access; opening a paid link does not expose its answers.