backenddrills

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.

InputContract to specify
Same entity and version againNo additional logical effect.
Older version after newer versionDoes not overwrite newer state.
Deletion followed by an old updateDeletion ordering remains protected.
Snapshot mixed with live changesTransition 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.

  1. Prepare duplicate, reordered and deleted-entity inputs.
  2. Interrupt a consumer between destination update and checkpoint advancement.
  3. Restart and replay the same bounded fixture.
  4. Compare authoritative source state with the rebuilt projection.
  5. 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 →

Explore 1008 scenarios · $28 one-time

Selected practice from the study paths

Free readings need no account. Full edition drills require verified access; opening a paid link does not expose its answers.

Recommended next readings