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).
Replication is a scaling technique because it:
- Offloads reads from the primary (e.g., Daraz’s inventory replicas in multiple data centers).
- Enables partition tolerance (CAP theorem: Availability over Consistency during network splits).
- 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:
- Primary node decrements stock to 9.
- Replica in Pokhara updates eventually (e.g., 5 seconds later).
- 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:
- Consistency (all reads return the most recent write).
- Availability (every request gets a response, even during partitions).
- Partition Tolerance (system works despite network failures).
| 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(whereN= 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
Synchronous Replication:
- Primary waits for acknowledgments from all backups before committing.
- Example: Oracle’s Data Guard.
- Downside: High latency (blocks writes).
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
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.
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
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").
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).
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
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.
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.
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").
Quorums:
- Calculate
R + W > Nfor givenN,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").
- Calculate
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.
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…