A global social network storing user feeds in Cassandra (Wide-Column NoSQL with Quorum writes), session caches in Redis (Key-Value), and user payment accounts in Google Cloud Spanner (TrueTime synchronized distributed SQL with ACID guarantees).
-- Sharding Strategy & Quorum Configuration in Distributed Databases -- 1. Consistent Hashing Token Ring Concept: -- Hash(user_id) mod 2^32 maps to specific Shard Node (e.g. Node A, B, C, D) -- 2. Quorum Consistency Equation (Strong Consistency / Linearizability): -- R (Read Quorum) + W (Write Quorum) > N (Replication Factor) -- Example: N = 3 Replicas -- Write Quorum W = 2 (waits for 2 ACK from replicas) -- Read Quorum R = 2 (reads from 2 replicas to pick latest timestamp) -- Since 2 + 2 = 4 > 3, at least 1 read node is guaranteed to have the latest write! -- 3. Distributed Two-Phase Commit (2PC) Protocol: -- Phase 1 (Prepare): Coordinator asks all Participant Shards: "Can you commit?" -> Participants acquire locks and vote YES/NO. -- Phase 2 (Commit): If ALL vote YES, Coordinator sends GLOBAL_COMMIT. If ANY vote NO or timeout, sends GLOBAL_ABORT.
Visual representation of control loops, memory layout, and execution flow for Distributed Databases, Replication, Sharding & NoSQL.
CAP Theorem: In any distributed asynchronous data store across a network subject to Network Partitions (P), you MUST choose between Consistency (C: every read receives the most recent write or an error) and Availability (A: every non-failing node returns a non-error response without guarantee of latest write). You CANNOT have CA during a network partition! PACELC Theorem Extension: If there is a Partition (P), how does the system trade Availability (A) vs Consistency (C); ELSE (E), when the network is running normally, how does the system trade Latency (L) vs Consistency (C)? (e.g. DynamoDB is PA/EL; Spanner is PC/EC).
1. Single-Leader (Primary-Replica): Writes go to Primary; Primary streams changes to Read Replicas (Synchronous = zero data loss but high latency; Asynchronous = fast writes but risk of Replication Lag and stale reads). 2. Multi-Leader: Writes accepted at multiple regional primary nodes; requires conflict resolution (Last-Write-Wins, CRDTs). 3. Leaderless (Dynamo-Style): Client writes directly to N replica nodes. Quorum Rule for Strong Consistency: R + W > N. If N=3, setting W=2 and R=2 guarantees the read set overlaps with the write set on at least one replica.
Sharding splits a large database table into smaller subsets distributed across multiple physical server nodes. - Range-Based Sharding: Partition by key range (e.g. A-M, N-Z). Risk: Hotspotting if traffic clusters on popular ranges. - Hash-Based Sharding: Partition by Hash(Key) % Number_of_Nodes. Problem: Adding/removing a node forces rehashing and moving 100% of keys! - Consistent Hashing: Maps both Nodes and Keys to a virtual 360-degree hash ring (0 to 2^32 - 1). A key is assigned to the first node clockwise from its hash position. Adding/removing a node only moves K/N keys (where K is total keys and N is total nodes)! Virtual nodes prevent uneven data skew.
Two-Phase Commit (2PC): Coordinates atomic transactions across distributed shards. Phase 1 (Prepare): Coordinator asks all nodes to prepare log changes; nodes lock resources and reply YES/NO. Phase 2 (Commit/Abort): Coordinator decides and broadcasts outcome. Weakness: Blocking protocol if coordinator crashes mid-commit. NoSQL Family Taxonomy: 1. Key-Value Stores (Redis, Memcached): Ultra-fast O(1) in-memory lookups for sessions and caching. 2. Document Stores (MongoDB, Couchbase): Flexible semi-structured JSON/BSON schemas for nested data. 3. Wide-Column Stores (Cassandra, ScyllaDB, HBase): High write throughput with sparse column families for time-series/sensor data. 4. Graph Databases (Neo4j): Direct node-and-edge pointer traversals for social graphs and fraud detection networks.
| Feature / Dimension | Relational Databases (SQL) | NoSQL Databases |
|---|---|---|
| Data Model & Schema | Rigid tabular schema with predefined columns, types, and foreign keys | Flexible schema (JSON documents, key-value, column families, graphs) |
| Transactions & Guarantees | Strong ACID guarantees with serializable isolation and immediate consistency | BASE model (Basically Available, Soft state, Eventual consistency) |
| Scaling Strategy | Primarily Vertical Scaling (Scale-Up); Horizontal sharding is complex | Native Horizontal Scaling (Scale-Out) across commodity server clusters |
Detailed answers, interviewer pro tips, key takeaway summaries, and code examples formulated for technical rounds.
✅ Correction: Network partitions are physical inevitabilities in distributed networks (cables cut, switch failures). When a partition occurs, a database MUST choose either Consistency (rejecting writes) or Availability (accepting writes on both sides). You cannot choose CA.
✅ Correction: 2PC is a blocking protocol: if the coordinator crashes during the prepare phase, participant shards remain locked indefinitely, destroying throughput. Modern microservices use asynchronous Saga patterns or Raft/Paxos consensus instead.
Architectures, replication protocols, sharding algorithms, and distributed consistency models.