Data Consistency and Distributed Transactions Questions
Maintaining correctness of state across services and replicas: eventual consistency, conflict resolution (last-write-wins, CRDTs, vector clocks), the saga pattern, two-phase commit, and idempotency keys for exactly-once effects. Covers when to trade strict consistency for availability and how to reason about read-your-writes and monotonic guarantees. Focuses on the application/service layer rather than storage-engine internals.
Implement a PN-Counter (Positive-Negative counter) CRDT in Python with methods: increment(amount), decrement(amount), merge(other), and value(). Provide the data structure and explain why merge is commutative, associative, and idempotent.
Sample Answer
Direct answer: A PN-Counter (positive-negative counter) supports both increment and decrement by internally tracking two separate G-Counters, one for increments and one for decrements, and reporting the value as their difference; merging two PN-Counters is just merging each half independently (element-wise max per replica, same as a plain G-Counter), which keeps the whole structure commutative, associative, and idempotent.
Structured elaboration. The key design insight is that a PN-Counter is NOT a single vector you merge with subtraction, subtraction isn't idempotent or safely mergeable the way max is. Instead it's two independent G-Counters (grow-only, merge via max) combined only at READ time via subtraction.
class PNCounter:
def __init__(self, replica_id: str):
self.replica_id = replica_id
self._increments = {} # replica_id -> count of increments FROM that replica
self._decrements = {} # replica_id -> count of decrements FROM that replica
def increment(self, amount: int = 1) -> None:
if amount < 0:
raise ValueError("use decrement() for negative changes")
self._increments[self.replica_id] = self._increments.get(self.replica_id, 0) + amount
def decrement(self, amount: int = 1) -> None:
if amount < 0:
raise ValueError("amount must be non-negative; this reduces the counter")
self._decrements[self.replica_id] = self._decrements.get(self.replica_id, 0) + amount
def merge(self, other: "PNCounter") -> "PNCounter":
"""Returns a NEW PNCounter representing the merge of self and other.
Does not mutate either input."""
merged = PNCounter(self.replica_id)
all_inc_replicas = set(self._increments) | set(other._increments)
all_dec_replicas = set(self._decrements) | set(other._decrements)
merged._increments = {
r: max(self._increments.get(r, 0), other._increments.get(r, 0))
for r in all_inc_replicas
}
merged._decrements = {
r: max(self._decrements.get(r, 0), other._decrements.get(r, 0))
for r in all_dec_replicas
}
return merged
def value(self) -> int:
return sum(self._increments.values()) - sum(self._decrements.values())
Why merge is commutative, associative, and idempotent. Each half (_increments, _decrements) is merged via element-wise max, and max itself has all three properties (max(a,b) == max(b,a); max(max(a,b),c) == max(a,max(b,c)); max(a,a) == a), so applying the same max operation per-replica-per-half preserves those properties for the whole structure. value() is a pure, deterministic function of the merged state, it doesn't affect convergence itself, it's just how you read out the current answer, which is why subtraction (not itself idempotent) is safe to use HERE, it's read-time-only, never part of the merge.
Verification. Two replicas, A and B, both increment locally without communicating: A increments by 3, B increments by 2, then B decrements by 1. A's local state: _increments={A:3}, _decrements={}, value = 3. B's local state: _increments={B:2}, _decrements={B:1}, value = 1. Merging A into B (or B into A, order shouldn't matter): _increments = {A:3, B:2}, _decrements = {B:1}, value = (3+2) - 1 = 4. Merging the other direction gives the identical result (commutativity holds by construction, since it's built from max and set union, both order-independent), and merging twice (merged.merge(merged)) leaves the value unchanged (idempotent, since max(x,x) = x).
Trade-offs and pitfalls. A common bug when implementing this is trying to store a SINGLE net value per replica instead of separate increment/decrement counters and merging with subtraction directly, that breaks idempotence: subtracting twice would double-subtract on a re-merge, unlike max, which is safe to reapply. The two-G-Counter design is what makes re-merging (which WILL happen under retries or redundant gossip) safe.
Design a multi-region user profile service that must support 100M users, 50k profile updates per second globally, and 1M reads per second. Requirements: users see their own updates immediately (read-your-writes) within a region, other users see updates eventually (within a bounded window), and 99th-percentile read latency stays low per region. Sketch the high-level architecture, replication strategy, and how you provide the read-your-writes guarantee without strong global coordination.
Sample Answer
Direct answer: At 100M users and 50k updates/sec, you need per-user data partitioned by region with a "home region" per user for writes, asynchronous cross-region replication for eventual global visibility, and a session-scoped mechanism (sticky routing to the user's home region, or a causal token) so each user sees their own updates immediately without requiring global synchronous coordination.
Structured elaboration
Partitioning and write routing. Each user is assigned a home region (based on signup location or explicit preference), and writes for that user's profile are always directed there, giving each user's writes a single, consistent, low-latency write path with no cross-region coordination needed per write. This is what makes 50k writes/sec globally distributed and still individually cheap: it's 50k INDEPENDENT single-region writes, not 50k globally-coordinated ones.
Cross-region replication. Each region's writes replicate asynchronously to every other region (a standard multi-region replication topology, e.g. a change stream fanning out from each region's primary store to the others). This is what provides "eventually visible to other users" within the required bound (here, one minute), the replication pipeline's own throughput and lag characteristics need to comfortably clear that bound under peak load, with monitoring on actual observed lag, not just a theoretical budget.
Read-your-writes without global coordination. Since the user's OWN writes always land in their home region, and their own subsequent reads can be routed to that SAME home region (sticky, based on the user's identity, not their current network location) for some bound (or indefinitely), the user always sees their own latest update, at native single-region latency, no waiting for cross-region replication for their OWN view. Other users reading this profile from a different region see it once replication catches up, within the required window, satisfying "others see it eventually within 1 minute" without that requirement touching the write or same-user-read path at all.
Achieving sub-50ms P99 read latency per region. Reads (for other users viewing a profile, not the owner's own reads) are served from the LOCAL region's replica, so they never cross a region boundary, keeping latency to local-network/local-disk numbers rather than being bounded by cross-region round-trip time (which alone would often exceed 50ms). This is only possible because the read doesn't need to be perfectly fresh, it needs to be fresh within a minute, which the async replication path already provides.
Handling the home-region-unavailable case. If a user's home region has an outage, either the user experiences degraded write availability (a real, honest trade-off of this design, since the write path is intentionally NOT cross-region-coordinated) or the system fails over the user's home-region assignment to another region (a more complex, operationally significant decision that itself needs careful handling to avoid conflicting writes from the old and new home regions during the transition).
Worked example. User in Region A updates their bio. The write commits to Region A's primary (their home region) as a purely local operation, with no cross-region network hop on the critical path, so its latency is bounded by local disk/network characteristics rather than by inter-region round-trip time, the same reasoning that keeps the P99 read latency for other users' LOCAL reads low. The response confirms success, and if the user immediately reloads their own profile, that read is also routed to Region A (sticky by user identity), showing the update instantly, own-write visibility achieved with zero cross-region dependency. Meanwhile the change replicates asynchronously to Regions B and C; a different user in Region B viewing this profile sees the OLD bio for however long replication takes (comfortably under the 1-minute bound under normal load, monitored explicitly for tail cases during regional replication backlogs).
Trade-offs and pitfalls. A common design mistake at this scale is routing the OWNER's reads based on their CURRENT network location rather than their home region (e.g. "nearest region" routing applied uniformly to all reads), which breaks read-your-writes the moment a user travels or is routed to a different region than the one their write landed in, own-write reads need identity-based (not network-proximity-based) routing specifically for this reason.
Design a scalable approach to support atomic increments for a counter that's sharded across many keys (for example, a global 'likes' count). Compare a few approaches (per-shard counters with periodic aggregation, CRDT counters, a central counter service, optimistic CAS-based increments) on accuracy, throughput, read latency, and reconciliation cost.
Sample Answer
Direct answer: For a sharded counter that needs to keep incrementing correctly under high concurrency without coordination, a CRDT counter (a PN-Counter, sharded by key with one counter object per key) gives exact eventual convergence with no central bottleneck; per-shard counters with periodic aggregation trade a small aggregation delay for simpler infrastructure; a central counter service gives immediate global consistency at the cost of becoming a write bottleneck; and optimistic CAS (compare-and-swap)-based increments work for moderate contention but degrade under very high concurrent write rates on the same key.
Structured elaboration
CRDT counter (PN-Counter per key). Each region/replica maintains its own local increment count for a given counter key; merging (element-wise max plus sum) converges to the exact correct total with no coordination required at write time. Accuracy is eventually exact (once all increments have propagated and merged, the total is precisely correct, nothing is approximated or lost). Storage overhead is one entry per replica per counter key, and merge/read cost scales with the number of replicas, not the number of increments, since increments to the same replica's entry just add to that entry rather than each needing its own record.
Per-shard counters with periodic aggregation. Instead of a CRDT's continuous merge, each shard just accumulates its own local count independently, and a periodic batch job sums all shards' counts to produce the displayed total. Simpler to build (no CRDT library or merge logic needed), but the displayed total is only as fresh as the last aggregation run, real-time reads see a stale, under-counted total between aggregation cycles, an explicit trade-off, not a bug, if the use case tolerates it.
Central counter service. A single service (backed by its own strongly-consistent store) owns the counter and serializes all increments through it. Gives immediate, always-accurate reads with no aggregation delay, but throughput is capped by that single service's write capacity, and it becomes a single point of contention (and, without careful design, a single point of failure) as write volume grows, exactly the scaling problem sharding elsewhere in the system was meant to avoid.
Optimistic CAS-based increments. Each increment reads the current value, computes the new value, and writes it back with a compare-and-swap (only succeeding if the value hasn't changed since the read); a losing CAS retries. Works well at moderate contention (few concurrent writers per key), but retry rate grows sharply as concurrent writers to the SAME key increase, under very high contention (many writers hammering one popular key, a "hot key"), this can degrade into a large fraction of writes retrying repeatedly, hurting both latency and throughput.
Comparison
| Approach | Accuracy | Throughput under high concurrency | Read latency | Reconciliation needed |
|---|---|---|---|---|
| CRDT (PN-Counter) | Eventually exact | High, no coordination on write | Read = merge of current replica states, cheap | None; merge IS the reconciliation |
| Per-shard + periodic aggregation | Exact as of last aggregation, else stale | Very high, purely local writes | Fast but stale between cycles | The aggregation job itself |
| Central counter service | Always exact, immediately | Bounded by the service's own capacity | Fast, single source of truth | None needed, but scaling requires sharding the service itself eventually |
| Optimistic CAS | Always exact, immediately | Degrades under high same-key contention | Fast when uncontended | None; but retry storms under contention are a real operational risk |
Worked example. A "post likes" counter under a viral post might receive thousands of concurrent increments per second from a single popular key. A central counter service or CAS-based approach would see heavy contention specifically on that one hot key; a PN-Counter sharded across regions handles it gracefully since each region's writers only contend with OTHER writers in the SAME region (much lower local contention), converging the true global total asynchronously with no single bottleneck.
Trade-offs and pitfalls. Teams sometimes default to a central counter service for simplicity and only discover the bottleneck once a specific key goes viral or otherwise becomes hot, worth explicitly asking during design "what's the expected max increment rate on a SINGLE key," not just the aggregate system-wide rate, since that's what determines whether central/CAS approaches will hold up.
Design an algorithm to compact vector clocks or CRDT deltas in a long-running distributed system so their metadata doesn't grow unbounded, while preserving eventual convergence. What heuristics would you use, what trade-offs do they introduce in causality precision, and how would you recover or reconcile after a compaction?
Sample Answer
Direct answer: Compacting vector clocks or CRDT deltas without unbounded metadata growth means pruning entries for replicas or clients that haven't contributed recently, snapshotting the current merged state periodically so old deltas can be discarded, or replacing full history-tracking with a bounded summary (like a Merkle-tree digest or a periodic checkpoint) that trades some precision in causality tracking for a hard bound on metadata size.
Structured elaboration
Why unbounded growth happens. A vector clock's size grows with the number of DISTINCT node/client IDs that have ever contributed to a piece of data's history; in a system with many ephemeral clients (mobile devices that connect briefly, short-lived worker processes), that set never shrinks under a naive implementation, every client that ever touched the data leaves a permanent entry. CRDT deltas have an analogous problem: an op-based CRDT that never compacts its operation log, or a state-based CRDT whose metadata (tombstones for deleted elements, per-replica counters) accumulates without bound, grows linearly with total historical activity rather than with current, live state.
Time-based eviction. Prune vector-clock entries (or CRDT tombstones) for a node that hasn't contributed an update within some window (e.g. 90 days), on the reasoning that a node silent for that long is very unlikely to suddenly resurface with an old, unmerged update that still needs precise causal placement. The risk this introduces: if a pruned node DOES resurface with a genuinely old update, the system can no longer precisely place it in causal history relative to everything that happened since, it gets treated more conservatively (e.g. as concurrent with everything, or requiring a full reconciliation pass) rather than correctly ordered.
Snapshotting. Periodically (e.g. daily, or after N operations), compute the fully-merged current state and take that AS the new baseline, discarding the individual deltas/operations that led to it. Anything before the snapshot no longer needs individual per-operation metadata, only the snapshot itself, which is a fixed, bounded size regardless of how much history preceded it. The trade-off: you lose the ability to answer fine-grained "what was the exact sequence of changes" questions for anything before the snapshot, acceptable for most systems where only the CURRENT convergent state matters, not the full audit history.
Hash-based summaries (Merkle trees). Instead of tracking exact causal metadata per key, periodically compute a Merkle tree (or similar hierarchical hash structure) over ranges of keys; comparing two replicas' Merkle trees quickly identifies WHICH ranges differ (via hash mismatch) without needing per-key vector clocks for keys that HAVEN'T diverged, useful for anti-entropy reconciliation at scale, where most keys agree between replicas and only a small subset needs detailed, per-key metadata to resolve.
Trade-off in losing precision. Every one of these techniques trades exact, unbounded-precision causality tracking for a bounded metadata footprint; the shared risk is FALSE CONCURRENCY DETECTION, an update that's actually causally related to something in the pruned/snapshotted history gets treated as if it had no prior context, potentially triggering an unnecessary conflict-resolution step (harmless but wasteful) or, in a worse case, missing a genuine causal ordering the system should have preserved.
Worked example. A collaborative-document CRDT accumulates tombstones for every deleted character over the document's lifetime, after a year of edits, tombstone metadata can dwarf the actual live content. The system snapshots the document's current state (the surviving, non-deleted content) every night, discarding tombstones for anything deleted before the snapshot point (since there's no remaining live content that could still conflict with an already-deleted character from before the snapshot), while keeping fresh tombstones for anything deleted since, to still correctly handle a late-arriving concurrent edit from just before the snapshot.
Recovery and reconciliation after compaction. A replica that reconnects after being offline longer than the retention window (with updates the system can no longer precisely causally place) is handled by falling back to a full state comparison against the current snapshot rather than trying to replay individual deltas, effectively treating it as a fresh join with a reconciliation pass, rather than a normal incremental merge, a graceful but more expensive degraded path reserved for the rare case of a very-late reconnect.
Compare implementing a global transaction coordinator using distributed consensus (Raft/Paxos) versus relying on a centralized ACID database for coordinating cross-shard transactions. Analyze latency, throughput, operational complexity, availability, and developer ergonomics, and give recommendations for systems of different scale and reliability requirements.
Sample Answer
Direct answer: A global transaction coordinator can be built two ways: as a single logical service backed by a consensus protocol like Raft or Paxos (so its own state survives node failures), or by delegating coordination to a centralized ACID database that already provides durable, consistent state. Consensus-backed coordination scales better and avoids a single database becoming the bottleneck, but costs an extra replication round-trip and more operational complexity; a centralized ACID database is simpler to reason about and operate but becomes a throughput ceiling and a single point of contention as the system grows.
Structured elaboration
Consensus-backed coordinator. The coordinator's transaction log (which participant voted what, what was decided) is replicated across an odd number of nodes via Raft or Paxos. A write only counts once a majority of replicas have durably stored it, so losing a minority of nodes (including the current leader) doesn't lose the decision, a new leader is elected and continues from the replicated log. This directly attacks 2PC's core weakness: instead of "only the coordinator knows the decision," it's "only a MAJORITY of coordinator replicas needs to know it."
Centralized ACID database as the coordinator's store. Instead of building custom replication, you let a database (itself internally replicated, e.g. via its own consensus or primary-replica mechanism) hold the coordinator's transaction table. The coordinator process is then effectively stateless: it reads/writes transaction records to the DB, and if the coordinator process itself dies, a new instance can pick up any in-flight transaction by querying the DB for records that aren't yet marked done.
Latency, throughput, complexity comparison
| Dimension | Consensus-backed (Raft/Paxos) | Centralized ACID database |
|---|---|---|
| Write latency | One consensus round-trip (majority ack) per state transition, typically comparable to a DB commit in a well-run cluster, but tunable (e.g. quorum size) | One DB commit per state transition; usually well-optimized, but subject to that DB's own replication latency |
| Throughput ceiling | Scales with cluster size and can shard the log by transaction ID / key range | Bounded by the single database's write throughput, harder to shard without re-deriving your own version of the consensus problem |
| Operational complexity | You own consensus cluster operations: leader election, log compaction, membership changes | You inherit whatever operational model the database already has, usually more familiar to most teams |
| Failure model | Tolerates loss of a minority of nodes with no external dependency | A single logical database is often already the availability bottleneck for the rest of the system too, so this adds shared fate |
| Developer ergonomics | Requires understanding consensus semantics (majority writes, leader-only reads for linearizability) | Ordinary transactions and SQL; most engineers already know the mental model |
Worked example. A payments platform coordinating transfers across 20 account shards chose a Raft-backed coordinator cluster of 5 nodes instead of a single Postgres instance, specifically because a single Postgres instance capped them at roughly 3-4k coordinator writes/sec under their durability settings, while sharding coordinator state across a Raft-backed key range let them scale past that by adding more coordinator shards, each independently replicated. The cost was a dedicated on-call rotation for the consensus cluster that didn't exist before.
When consensus is preferred vs when a centralized DB (or plain 2PC) is still fine. Consensus-backed coordination earns its complexity when coordinator throughput or availability is itself the bottleneck, and when you already need consensus elsewhere in the system (e.g. for leader election), so you're not paying for a brand-new capability. A centralized ACID database remains the right default when transaction volume is modest, the team doesn't want to operate a consensus cluster, and the database is already a dependency the rest of the system tolerates. Plain 2PC without either (a single coordinator process with a local durable log) is fine only when the coordinator's own availability is not a differentiated requirement, i.e. a short outage of the coordinator is an acceptable cost.
Trade-offs and pitfalls. A common mistake is assuming consensus fixes 2PC's blocking problem for participants too, it only makes the COORDINATOR's decision durable and available; participants still block waiting to hear that decision, they just no longer wait for one specific fragile process to come back, they wait for a majority of a replicated cluster to answer, which is a much smaller and more bounded risk.
Unlock Full Question Bank
Get access to all Data Consistency and Distributed Transactions interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.