Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Eventual vs strong consistency

Data platform · Distributed Systems Basics

Eventual vs strong consistency

Mediumdata-platform-27
eventual-consistencyreplication-lagdistributed-systemscdc

Question

What is eventual consistency, and where does it show up in data engineering?

Solution

Eventual consistency means that storage replicas receive updates asynchronously and converge to the same state over time if no new mutations occur, meaning read requests may temporarily return stale values. In contrast, strong consistency guarantees that any read executed after a completed write always reflects that newest mutation, but it imposes significant latency and coordination overhead. Data engineering pipelines frequently interact with eventual consistency when reading from database replicas, stream buffers, and distributed key-value layers.

Convergence across distributed nodes

In modern architectures, eventual consistency appears wherever data is replicated without synchronous distributed locking:

  • Database read replicas: Change data capture engines often tail read replicas instead of the primary transactional node to avoid production lock contention. If the replica lags five seconds behind the primary, CDC readers and downstream Kafka topics broadcast events with built-in temporal lag.
  • Distributed NoSQL databases: Systems such as Amazon DynamoDB and Apache Cassandra use tunable quorum reads and writes. Their default read configurations query a single replica node to minimize latency, occasionally returning outdated records until background gossip and hinted handoffs complete.
  • Distributed caching tiers: Redis caches and CDN edges hold cached payloads with time-to-live expiration policies, returning stale snapshots while the underlying warehouse or operational store has updated.
Write -> Primary DB --(async replication lag: 2-5s)--> Read Replica
                                                              |
Pipeline Consumer reads stale record until replica catches up <+

Pipeline designs must be constructed to handle this replica drift gracefully:

Where pipelines absorb replication lag

Attempting to force strong consistency across high-volume analytical workloads causes excessive network round-trips, lock queues, and increased compute bills. Battle-tested data platforms design pipelines to tolerate lag instead of fighting it.

When extracting CDC streams from read replicas, pipelines track source commit timestamps and log sequence numbers (LSN) rather than wall-clock ingestion time. Downstream aggregation stages apply watermark windows and deduplication logic, buffering out-of-order records until watermarks advance. When exposing data marts to dashboards, batch jobs write to staging tables and swap partition pointers atomically, preventing analytical users from querying intermediate states during prolonged convergence cycles.

PreviousNext