CACS352 Distributed System

Distributed SystemUnit 1013 min read

Consistency Models & Replication in Distributed Systems

Unit 10 of Distributed System: Explores how data consistency is maintained across replicated copies in distributed systems, compares strong vs. eventual consistency models, and explains replication strategies (primary-backup, multi-leader, peer-to-peer) with real-world examples like eSewa transactions and Daraz invento

TAKEAWAYS:

  • Replication improves availability and fault tolerance but introduces consistency challenges.
  • Consistency models range from strong (linearizability) to eventual (weak consistency).
  • Primary-backup replication is simple but creates bottlenecks; multi-leader replication enables offline edits but risks conflicts.
  • Conflict resolution uses timestamps, vectors, or application-specific rules (e.g., eSewa’s last-write-wins for transaction IDs).
  • CAP theorem shows you can only prioritize two of three: Consistency, Availability, or Partition tolerance.
  • Quorum-based replication (e.g., Cassandra) trades consistency for performance by requiring only a subset of replicas to agree.

1. Why Replicate Data?

Distributed systems replicate data to:

  • Improve availability (failover to backup nodes).
  • Reduce latency (serve users closer to them).
  • Handle high throughput (distribute read/write load).
Node FailuresNetwork PartitionsFault ToleranceLoad BalancingReduced LatencyPerformanceHigh UptimeGeographic RedundancyAvailabilityDistributed System
Primary reasons for data replication in distributed systems, illustrated with a decision tree.

Replication is a scaling technique because it:

  1. Offloads reads from the primary (e.g., Daraz’s inventory replicas in multiple data centers).
  2. Enables partition tolerance (CAP theorem: Availability over Consistency during network splits).
  3. Supports geographic distribution (eSewa’s payment gateways replicate across Kathmandu and Pokhara).

2. Consistency Models: The Spectrum of Trade-offs

Consistency defines how quickly and reliably replicas agree on data. The spectrum ranges from strong to eventual:

Model Definition Example Use Case Trade-off
Linearizability Operations appear instantaneous at a single logical time. Bank transactions (NMB’s real-time balance updates). High latency, strict ordering.
Sequential Consistency Operations execute in a total order across all nodes. WhatsApp message delivery (messages arrive in sent order). Requires global clocks.
Causal Consistency Operations respect causality (if A → B, B must not precede A). Pathao ride requests (driver assignments). Complex to implement.
Eventual Consistency Replicas converge eventually (no guarantees on speed). Social media feeds (Facebook/Twitter posts). Accepts stale data temporarily.

FIGURE: Consistency Spectrum

```mermaid
flowchart TD
  A["Linearizability"] -->|"Stronger"| B["Sequential Consistency"]
  B -->|"Stronger"| C["Causal Consistency"]
  C -->|"Stronger"| D["Eventual Consistency"]
  A -->|"Requires"| E["Global Clock"]
  D -->|"Tolerates"| F["Stale Reads"]
  E -->|"Example"| G["Pathao Ride Requests"]
  F -->|"Example"| H["Social Media Feeds"]
  G -->|"Complexity"| I["Complex to Implement"]
  F -->|"Trade-off"| J["Accepts Stale Data Temporarily"]

Worked Example: Daraz’s Inventory

  • Model: Eventual Consistency
  • Scenario: A user buys a phone (inventory drops from 10 to 9).
  • Trace:
    1. Primary node decrements stock to 9.
    2. Replica in Pokhara updates eventually (e.g., 5 seconds later).
    3. If a user checks stock in Pokhara during this gap, they see 10 (temporarily inconsistent).
  • Why? Daraz prioritizes availability (users see items) over strong consistency (immediate updates).

3. Replication Strategies

Replication strategies define who writes to whom and how replicas sync. Three dominant models:

A. Primary-Backup Replication

  • How it works:
    • A primary node handles all writes; backups replicate asynchronously or synchronously.
    • Reads can go to any replica (but primary is authoritative).
  • Advantages:
    • Simple to implement (e.g., MySQL master-slave).
    • Strong consistency if synchronous.
  • Disadvantages:
    • Bottleneck: Primary is a single point of failure.
    • Latency: Backups lag behind (e.g., NTC’s DNS servers replicate from Kathmandu to remote offices).

FIGURE: Primary-Backup Replication

```figure
{"type":"layers","layers":["Client","Primary Node","Backup 1","Backup 2"],"right":["Write/Read","Async/Sync Write","Async/Sync Write","Read"],"caption":"Primary-Backup Replication: Client interactions with primary and backups. NTC’s DNS replication from Kathmandu to remote offices is an example."}

B. Multi-Leader Replication

  • How it works:
    • Multiple nodes can accept writes (e.g., GitHub’s branch merges).
    • Conflicts resolved via last-write-wins (LWW), timestamps, or application logic.
  • Advantages:
    • Offline support: Edits on mobile apps sync later (e.g., WhatsApp messages).
    • Lower latency: Users write to the nearest leader.
  • Disadvantages:
    • Conflict resolution is complex (e.g., two users editing the same Daraz order).
    • Data divergence: Temporarily inconsistent replicas.

Conflict Example: eSewa Transaction

  • Scenario: User A transfers ₹1000 to User B in Kathmandu; User C does the same in Pokhara.
  • Conflict: Both leaders process the transfer, but the total balance in the database may temporarily show ₹2000 (double-spend risk).
  • Solution: eSewa uses transaction IDs to detect and rollback duplicates.

C. Peer-to-Peer (P2P) Replication

  • How it works:
    • All nodes are equal; writes propagate via gossip protocols (e.g., Cassandra).
    • No designated primary.
  • Advantages:
    • Scalable: No single bottleneck (e.g., Bitcoin’s blockchain).
    • Fault-tolerant: Nodes can join/leave dynamically.
  • Disadvantages:
    • Weak consistency: Eventual by default.
    • Complex coordination: Requires conflict-free replicated data types (CRDTs).

4. Conflict Resolution Techniques

When replicas diverge, systems use one of these:

Technique How It Works Example Pros/Cons
Last-Write-Wins (LWW) The latest timestamped write wins. eSewa’s transaction IDs. Simple but loses data on conflicts.
Vector Clocks Tracks causality (e.g., "A → B") to resolve dependencies. Git’s merge conflicts. Complex but preserves causality.
Application Logic Custom rules (e.g., "highest bidder wins"). Daraz auctions. Flexible but requires domain knowledge.
CRDTs (Conflict-Free Replicated Data Types) Data structures that auto-resolve conflicts (e.g., sets, counters). Riak’s distributed counters. Strong consistency without coordination.

FIGURE: Vector Clock Example

```mermaid
stateDiagram-v2
    [*] --> Node1: {A:1}
    Node1 --> Node2: {A:1, B:2}
    Node2 --> Node3: {A:1, B:2, C:3}
    Node3 --> Node1: {A:1, B:2, C:3, D:4}

Nodes track causality to merge changes without conflicts.


5. CAP Theorem: The Impossible Triangle

The CAP theorem states that in a partitioned network, you can only guarantee two of:

  1. Consistency (all reads return the most recent write).
  2. Availability (every request gets a response, even during partitions).
  3. Partition Tolerance (system works despite network failures).
111ConsistencyAvailabilityPartition Tolerance
CAP Theorem: Trade-offs between consistency, availability, and partition tolerance in distributed systems. No system can satisfy all three simultaneously under
System Type Prioritizes Example
Strongly Consistent C + P Google Spanner (global databases).
Available A + P Cassandra (eventual consistency).
Partition-Tolerant A + C (but not P) Local databases (no partitions).

Real-World Example: NTC vs. Ncell

  • NTC (Consistency): Uses synchronous replication for call records (C + P).
  • Ncell (Availability): Relies on asynchronous replication to handle rural outages (A + P).

6. Quorum-Based Replication

To balance consistency and performance, systems use quorums:

  • Read Quorum (R): Minimum replicas to read from.
  • Write Quorum (W): Minimum replicas to write to.
  • Condition: R + W > N (where N = total replicas) ensures consistency.

Example: Cassandra’s Tunable Consistency

  • N = 6 replicas, R = 4, W = 4:
    • A write must succeed on 4 nodes.
    • A read must check 4 nodes (but can return stale data if <4 replicas are up).
  • Trade-off: Higher W = stronger consistency but lower availability.

FIGURE: Quorum Overlap


7. Replication Protocols

Protocols define how replicas sync. Key examples:

A. Primary-Backup Protocols

  1. Synchronous Replication:

    • Primary waits for acknowledgments from all backups before committing.
    • Example: Oracle’s Data Guard.
    • Downside: High latency (blocks writes).
  2. Asynchronous Replication:

    • Primary commits immediately; backups catch up later.
    • Example: MySQL master-slave.
    • Downside: Data loss if primary fails before syncing.

B. State Transfer Protocols

  • Snapshot Transfer: Replicas periodically exchange full state dumps (e.g., ZooKeeper).
  • Log Transfer: Replicas replay write-ahead logs (WAL) (e.g., PostgreSQL’s streaming replication).

8. Replication in Practice: eSewa’s Payment System

Scenario: User A transfers ₹5000 to User B via eSewa. Replication Model: Multi-Leader with Conflict Resolution

  1. Write Path:

    • Request → Primary Node (Kathmandu) → Deducts ₹5000 from A’s balance.
    • Conflict Detection: If Pokhara’s leader also processes the transfer, eSewa’s transaction ID ensures no double-spend.
    • Sync: Changes propagate to Pokhara via asynchronous replication.
  2. Read Path:

    • User B checks balance → Reads from Pokhara replica (lower latency).
    • Stale Read Risk: Pokhara may show ₹5000 temporarily if Kathmandu hasn’t synced yet.

Why Multi-Leader?

  • Offline Support: Users in remote areas can still transact.
  • Locality: Pokhara users interact with Pokhara nodes.

Conflict Handling:

  • Last-Write-Wins: If two transfers for ₹5000 happen simultaneously, the one with the higher timestamp wins.
  • Application Logic: For critical transfers (e.g., loan repayments), eSewa uses two-phase commit to ensure atomicity.

In the Real World

  1. eSewa (Multi-Leader Replication)

    • Idea: Uses multi-leader replication with transaction IDs to resolve conflicts.
    • How: Kathmandu and Pokhara nodes accept writes independently. Conflicts are detected via version vectors and resolved by the application logic (e.g., "the transfer with the higher timestamp is valid").
  2. Daraz (Eventual Consistency + Quorums)

    • Idea: Eventual consistency for inventory with quorum-based reads/writes.
    • How: Daraz’s global inventory uses Cassandra (R=3, W=2 out of 6 replicas). A user in Pokhara may see a product as "in stock" even if Kathmandu’s primary hasn’t updated yet (but the system guarantees eventual consistency).
  3. Pathao (Causal Consistency for Ride Assignments)

    • Idea: Causal consistency to ensure ride assignments respect causality (e.g., if Driver A is assigned to User X, Driver B cannot be assigned to User X’s same request).
    • How: Pathao’s backend uses vector clocks to track dependencies between ride requests and driver assignments.

Exam Tip

  1. Compare Consistency Models:

    • Always link models to real systems (e.g., "Google Spanner uses linearizability for financial transactions").
    • Draw the spectrum figure in exams to show trade-offs.
  2. Replication Strategies:

    • For primary-backup, emphasize bottlenecks and synchronous vs. asynchronous.
    • For multi-leader, highlight conflict resolution (LWW, timestamps, CRDTs).
    • Use Daraz/eSewa examples to explain trade-offs.
  3. CAP Theorem:

    • Memorize the triangle and give one example per quadrant (e.g., "Cassandra is AP, Spanner is CP").
    • Relate to network partitions (e.g., "If NTC’s network splits, it must choose between consistency and availability").
  4. Quorums:

    • Calculate R + W > N for given N, R, W (e.g., "For N=5, R=3, W=3, explain consistency").
    • Mention tunable consistency (e.g., "Cassandra allows adjusting R/W for performance").
  5. Conflict Resolution:

    • Describe vector clocks or CRDTs with a small diagram (e.g., "Show how two nodes merge changes using causality").
    • Use eSewa’s transaction ID as a concrete example.
  6. Real-World Mapping:

    • Always tie theory to Nepali apps (eSewa, Daraz, Pathao) or global systems (Google, WhatsApp).
    • For replication protocols, mention synchronous vs. asynchronous and their impact on latency vs. durability.

Key Formula to Remember:

Consistency = (R + W > N)  // Quorum condition for strong consistency

Based on the TU BCA syllabus for Distributed System (CACS352), unit 10.

Discussion

Loading…