Everything so far in this track assumed one relational database on one machine. Real systems often outgrow that: too much data for one disk, too many writes for one CPU, users on several continents, or data that does not fit neatly into tables. NoSQL databases and distributed databases are the answers. This lesson explains the four NoSQL data models with real examples, the CAP theorem and its more useful cousin PACELC, BASE versus ACID, the three replication styles, sharding strategies including consistent hashing, distributed transactions with two-phase commit, NewSQL, and a decision table for choosing SQL or NoSQL. Interviewers expect you to explain CAP correctly (most candidates do not), compare data models with concrete use cases, and justify a database choice for a given scenario. For large-scale patterns, follow up with databases at scale.
Why NoSQL?
NoSQL ("not only SQL") is an umbrella term for databases that do not use the relational table model as their primary interface. They became popular in the late 2000s when web companies hit limits of single-server relational databases.
The main motivations:
- Horizontal scaling. Vertical scaling (scaling up) means buying a bigger machine; it hits a ceiling and gets expensive. Horizontal scaling (scaling out) means adding more machines. Many NoSQL stores were designed from day one to spread data across many nodes automatically.
- Flexible schema. Relational tables need a fixed schema and a migration to change it. Document stores let each record carry its own fields, which suits rapidly changing products and data with many optional attributes.
- Data models that match the access pattern. A shopping cart is naturally a key-value lookup; a social graph is naturally a graph; an IoT sensor feed is naturally a time-ordered wide row. Forcing these into normalised tables can mean many joins.
- High write throughput and availability. Some workloads (logs, metrics, messaging) need huge write rates and must stay writable during network failures, even at the cost of temporarily inconsistent reads.
The trade-offs are real: most NoSQL stores give up some combination of joins, multi-row transactions, ad hoc queries, and strong consistency. Modern systems have narrowed the gap (MongoDB has multi-document transactions since 4.0, PostgreSQL has JSONB), so the choice is now about fit, not fashion.
The four NoSQL data models
Key-value stores
Model. A giant distributed hash map: each key maps to an opaque value (a string, a blob, sometimes a structure). The only fast operations are get, put and delete by key.
Examples. Redis (in-memory, with rich value types such as lists, sets, sorted sets and hashes), Amazon DynamoDB (key-value and document), Memcached (cache only), etcd (small, strongly consistent configuration data).
key value
------------------------- ---------------------------------
session:8f3a2c {"user_id": 42, "expires": ...}
cart:42 [{"sku": "KB-7", "qty": 1}]
rate:ip:203.0.113.9 17
# Redis commands
SET session:8f3a2c "{\"user_id\":42}" EX 3600 # expire in an hour
GET session:8f3a2c
INCR rate:ip:203.0.113.9 # atomic counter
ZADD leaderboard 980 "asha" 1210 "ravi" # sorted set
ZREVRANGE leaderboard 0 9 WITHSCORES # top 10
Fits: caching, sessions, rate-limit counters, leaderboards, feature flags, shopping carts. Poor fit: queries by anything other than the key, relationships, reporting.
Document stores
Model. Each record is a document: a self-describing, nested structure such as JSON (MongoDB stores BSON, a binary JSON). Documents live in collections. Unlike key-value stores, the database understands the document's fields, so you can query and index them.
Examples. MongoDB, Couchbase, Amazon DocumentDB, Firestore.
{
"_id": "order_1001",
"customer": { "id": 42, "name": "Asha" },
"status": "shipped",
"items": [
{ "sku": "KB-7", "name": "Keyboard", "qty": 1, "price": 2499 },
{ "sku": "MS-2", "name": "Mouse", "qty": 2, "price": 799 }
],
"created_at": "2026-10-10T09:15:00Z"
}
// MongoDB shell
db.orders.createIndex({ "customer.id": 1, created_at: -1 })
db.orders.find({ "customer.id": 42 }).sort({ created_at: -1 }).limit(20)
The key design idea is embedding: data that is read together is stored together. The order above includes its items, so one read returns everything, with no join. The cost is duplication (the customer's name is copied into every order) and awkward many-to-many relationships.
Fits: product catalogues with varied attributes, content management, user profiles, event data, aggregates read as a whole. Poor fit: highly relational data with many-to-many links, heavy cross-document reporting.
Wide-column stores
Model. Data is organised in tables of rows, but each row is identified by a partition key that decides which node stores it, and rows within a partition are sorted by clustering columns. Rows can have different sets of columns. Think of it as a distributed, sorted map of maps: partition key → (clustering key → columns).
Examples. Apache Cassandra, ScyllaDB, Apache HBase, Google Bigtable.
-- Cassandra Query Language (CQL)
CREATE TABLE sensor_readings (
sensor_id text,
day date,
reading_ts timestamp,
temperature double,
PRIMARY KEY ((sensor_id, day), reading_ts)
) WITH CLUSTERING ORDER BY (reading_ts DESC);
SELECT reading_ts, temperature
FROM sensor_readings
WHERE sensor_id = 's-17' AND day = '2026-10-10'
LIMIT 100;
Here (sensor_id, day) is the partition key (one partition per sensor per day, which keeps partitions bounded) and reading_ts is the clustering column (rows sorted newest first inside the partition).
In Cassandra you design tables per query: the query must supply the partition key. There are no joins, and filtering on other columns is discouraged. If you need the same data by another key, you write it to a second table (denormalisation).
Fits: very high write volumes, time-series data, event logs, messaging inboxes, IoT, anything with a natural partition key. Poor fit: ad hoc queries, joins, strong multi-row transactions.
Graph databases
Model. Nodes (entities) and relationships (edges) that both carry properties. Relationships are stored as direct pointers, so following them (traversal) costs about the same no matter how big the graph is. In a relational database, each hop is a join.
Examples. Neo4j (Cypher query language), Amazon Neptune, JanusGraph.
(Asha:Person) -[:FOLLOWS]-> (Ravi:Person) -[:FOLLOWS]-> (Meera:Person)
| ^
+-------------------[:FOLLOWS]----------------------+
// Cypher: people followed by people Asha follows,
// whom Asha does not already follow (friend-of-friend suggestions)
MATCH (a:Person {name: 'Asha'})-[:FOLLOWS]->(:Person)-[:FOLLOWS]->(s:Person)
WHERE NOT (a)-[:FOLLOWS]->(s) AND s <> a
RETURN s.name, count(*) AS mutual
ORDER BY mutual DESC
LIMIT 10;
Fits: social networks, recommendations, fraud detection (rings of accounts sharing devices or cards), knowledge graphs, network and dependency analysis. Poor fit: bulk aggregation over all records, simple CRUD with no relationships.
Comparison
| Model | Unit of data | Query by | Example | Typical use |
|---|---|---|---|---|
| Key-value | Opaque value | Key only | Redis, DynamoDB | Cache, sessions, counters |
| Document | JSON-like document | Key and indexed fields | MongoDB | Catalogues, profiles, content |
| Wide-column | Rows in partitions, sorted | Partition key (+ clustering range) | Cassandra | Time series, write-heavy logs |
| Graph | Nodes and edges | Traversals | Neo4j | Social graphs, fraud, recommendations |
| Relational | Rows in tables | Any column, joins | PostgreSQL, MySQL | General purpose, transactions |
The CAP theorem
The CAP theorem (conjectured by Eric Brewer in 2000, proved by Gilbert and Lynch in 2002) is about a data store replicated across several nodes. It considers three properties:
- Consistency (C): every read receives the most recent write or an error. Formally, linearizability: the system behaves as if there were one copy of the data. (Not the C of ACID.)
- Availability (A): every request to a non-failed node receives a non-error response, without a guarantee that it reflects the latest write.
- Partition tolerance (P): the system keeps working even when the network drops or delays messages between nodes (a network partition).
The theorem: when a network partition happens, a system must choose between consistency and availability.
Client X Client Y
| |
v v
+---------+ network cut +---------+
| Node 1 | ----- X X X ------- | Node 2 |
| bal=500 | | bal=500 |
+---------+ +---------+
X writes bal=300 on Node 1. Node 1 cannot tell Node 2.
Y reads bal on Node 2. Node 2 must either:
- answer 500 (stale) -> available, not consistent (AP)
- refuse / wait for Node 1 -> consistent, not available (CP)
Partitions are not optional in a distributed system: networks do fail. So the real choice is CP or AP during a partition. "CA" only describes a system that is not distributed (a single node), or one that simply stops working correctly when partitioned.
CP examples. Systems that refuse or block some requests during a partition to avoid returning stale data: ZooKeeper and etcd (a minority side cannot serve writes), HBase, MongoDB with majority write and read concerns (a primary cut off from the majority steps down), a relational database with synchronous replication that blocks commits when the replica is unreachable.
AP examples. Systems that keep accepting reads and writes on both sides and reconcile later: Cassandra and Riak with low consistency levels, DynamoDB with eventually consistent reads (the default), DNS, CouchDB.
Many systems are tunable: Cassandra with QUORUM reads and writes behaves much more like CP for those operations, and with ONE like AP.
Common mistake
"Pick two of C, A and P" is the classic oversimplification. You cannot give up P in a distributed system; the choice between C and A only bites during a partition. When the network is healthy, a well-built system can offer both. Saying this clearly is one of the easiest ways to stand out.
PACELC
CAP says nothing about the normal case, yet that is where systems spend almost all their time. PACELC (proposed by Daniel Abadi, 2010) extends it:
If there is a Partition, choose between Availability and Consistency; Else (normal operation), choose between Latency and Consistency.
The "else" half is the everyday trade-off: to make a write strongly consistent, it must reach other replicas (possibly in other regions) before acknowledging, which costs latency. To be fast, acknowledge locally and replicate later, accepting that a read elsewhere may be stale.
| System (typical configuration) | During partition | Normally | PACELC label |
|---|---|---|---|
| Cassandra, DynamoDB (default), Riak | Availability | Latency | PA/EL |
| ZooKeeper, etcd, HBase, Google Spanner | Consistency | Consistency | PC/EC |
These labels describe typical configurations. Many systems let you tune consistency per request, so the same database can sit in different cells depending on settings (MongoDB, for example, is classified differently depending on its read and write concerns).
Concrete example: an Indian payments app with servers in Mumbai and Hyderabad. If a payment write must be confirmed by both cities before the user sees "success" (EC), each payment pays an inter-city round trip of several milliseconds, but no replica ever shows a stale balance. If the write is confirmed by Mumbai alone (EL), the user gets a faster response, but a read served from Hyderabad a moment later might not show it, and a Mumbai failure in that window could lose it.
BASE versus ACID
BASE describes the guarantees of many AP systems:
- Basically Available: the system answers requests, even during failures, possibly with stale or partial data.
- Soft state: the state of the system can change over time without new input, as replicas converge.
- Eventual consistency: if no new writes arrive, all replicas eventually return the same value.
| ACID | BASE | |
|---|---|---|
| Goal | Correctness of every transaction | Availability and scale |
| Consistency | Strong, immediate | Eventual |
| Transactions | Multi-row, multi-table | Often single-key or single-document |
| Typical systems | PostgreSQL, MySQL, Oracle, Spanner | Cassandra, DynamoDB (default), Riak |
| Developer burden | Lower: the database guards invariants | Higher: handle stale reads, conflicts, retries |
Eventual consistency is weak on its own, so systems add session guarantees that are cheap and very useful:
- Read-your-writes: after you write, your own reads see that write (others may not yet).
- Monotonic reads: once you have seen a value, you never see an older one later.
- Consistent prefix: you see writes in the order they happened (no answer before its question).
Replication
Replication keeps copies of the same data on several nodes for availability (survive node failure), read scaling (serve reads from replicas) and latency (place copies near users).
Leader-follower (single-leader)
One node, the leader (also called primary or master), accepts all writes. It sends its change log to followers (replicas, secondaries), which apply the changes in the same order. Reads can go to the leader or followers.
writes
clients --------> +--------+
| Leader |
+--------+
/ | \ replication log
v v v
+------+ +------+ +------+
| F1 | | F2 | | F3 | <-- reads (possibly stale)
+------+ +------+ +------+
Used by PostgreSQL streaming replication, MySQL replication, MongoDB replica sets, and Redis replication. Simple and no write conflicts, but all writes go through one node, and failover (promoting a follower when the leader dies) must be handled carefully to avoid two leaders (split brain) or losing recently acknowledged writes.
Multi-leader
Several nodes accept writes (for example one leader per data centre) and replicate to each other.
- Pros: local, low-latency writes in each region; a region keeps working if another is cut off.
- Cons: write conflicts: two regions update the same row concurrently. They must be resolved, by last-write-wins (LWW, keep the write with the later timestamp, silently dropping the other), by merge logic, or by CRDTs (conflict-free replicated data types, structures such as counters and sets designed to merge automatically).
Used in multi-region setups, offline-first apps (each device is effectively a leader), and collaborative editing.
Leaderless
Any replica accepts reads and writes. A client (or coordinator) sends each write to all N replicas and waits for W acknowledgements; each read queries R replicas and takes the newest value.
Quorum rule: if R + W > N, every read set overlaps every write set in at least one replica, so a read sees the latest acknowledged write (under normal conditions).
Worked example with N = 3 replicas:
| W | R | R + W > N? | Behaviour |
|---|---|---|---|
| 2 | 2 | 4 > 3, yes | Balanced; tolerates 1 replica down for both reads and writes |
| 3 | 1 | 4 > 3, yes | Fast reads; a write fails if any replica is down |
| 1 | 3 | 4 > 3, yes | Fast writes; a read fails if any replica is down |
| 1 | 1 | 2, no | Fastest, but reads may be stale |
Stale replicas are repaired by read repair (a read that sees an old value writes the new one back) and anti-entropy (background comparison, often with Merkle trees). Used by Cassandra, Riak and the original Amazon Dynamo design.
Synchronous versus asynchronous replication
- Synchronous: the leader waits until a follower confirms the write before acknowledging the client. No acknowledged write is lost if the leader dies, but every write pays the replica's latency, and a slow or dead follower stalls writes.
- Asynchronous: the leader acknowledges immediately and ships changes in the background. Fast, but followers lag (replication lag), and a leader crash can lose the last acknowledged writes.
- Semi-synchronous: wait for one follower (or a quorum), ship to the rest asynchronously. A common compromise (PostgreSQL
synchronous_standby_nameswithANY 1, MySQL semi-sync replication).
Replication lag and read-your-writes
With asynchronous followers, this happens:
t=0 user updates profile name -> leader (ack)
t=5ms user reloads page -> follower F2 (lag 200 ms)
F2 still has the old name: "my change disappeared!"
Ways to provide read-your-writes:
- Read data the user can edit (their own profile) from the leader; read everything else from followers.
- Remember the time (or log position) of the user's last write and route their reads to the leader for a short window, or to a follower that has caught up to that position.
- Pin a user's session to one replica.
Partitioning (sharding)
Replication copies the same data to many nodes. Partitioning, called sharding when spread over machines, splits different data across nodes so that each holds a subset. Each piece is a partition or shard. You shard when one machine cannot hold the data or handle the write load. In practice systems do both: each shard is replicated.
The shard key (partition key) decides where each row goes. Choosing it well is the main design decision.
Range partitioning
Each shard holds a contiguous range of keys.
Shard 1: user names A-F Shard 2: G-M Shard 3: N-S Shard 4: T-Z
- Pros: range queries are efficient (all of "names starting with K" is on one shard); easy to understand.
- Cons: hotspots when access is not uniform. With a timestamp key, all of today's writes go to the newest shard while the rest sit idle.
- Used by HBase, Bigtable, and range-based sharding in MongoDB.
Hash partitioning
Apply a hash function to the key and assign by the hash: shard = hash(key) mod N.
- Pros: spreads keys evenly, avoiding sequential-key hotspots.
- Cons: range queries must ask every shard. And changing N reshuffles almost everything: going from 4 to 5 shards with
mod, a key stays only if hash mod 4 equals hash mod 5, which (for uniformly distributed hashes) happens for 4 out of every 20 values, so 80 % of keys move.
Consistent hashing (introduction)
Consistent hashing places both nodes and keys on a circle of hash values (a ring). Each key belongs to the first node found moving clockwise from the key's position.
0 / 2^32
|
node D .-*-. node A
*' '*
/ k1 -> \ k1 is stored on A
| |
\ <- k2 / k2 is stored on C
*. .*
node C '-*-' node B
When a node is added, it takes over only the keys between it and its counter-clockwise neighbour; every other key stays put. On average only about 1 / (new node count) of keys move: going from 4 to 5 nodes moves about 20 % of keys, compared with 80 % for mod. (A quick simulation with 20,000 keys and 100 virtual nodes per server moved 19 % of keys.)
To avoid uneven arcs, each physical node is placed at many points on the ring (virtual nodes), which also lets a bigger machine take more points. Consistent hashing is used by Cassandra, Riak, DynamoDB's internals, and many distributed caches. The databases at scale lesson goes deeper.
Hotspots and how to handle them
Even with good hashing, a single very popular key (a celebrity's profile, a viral product) sends all its traffic to one shard. Mitigations:
- Key salting: split a hot key into
key#0...key#9across shards; writes pick one at random, reads gather all ten. - Caching hot reads in front of the database.
- Better shard keys: avoid monotonically increasing keys (timestamps, auto-increment IDs) as range keys; use compound keys such as
(user_id, day). - Splitting hot ranges into smaller shards (automatic in many systems).
Costs of sharding
- Cross-shard queries must scatter to all shards and gather results.
- Cross-shard joins are expensive or unsupported; design so related data shares a shard key (an order and its items both keyed by
customer_id). - Cross-shard transactions need a distributed commit protocol (below).
- Rebalancing moves data while serving traffic.
- Unique constraints and auto-increment IDs no longer work globally; use UUIDs or ID services (for example Snowflake-style IDs).
Distributed transactions and two-phase commit
When one transaction touches data on several nodes (shards or different databases), all of them must commit or all must abort. Two-phase commit (2PC) coordinates this with a coordinator and several participants.
Coordinator Participant A Participant B
| --- PREPARE ------------> | |
| --- PREPARE -------------------------------> |
| write prepare record, |
| lock rows |
| <-- YES (vote commit) -- | |
| <-- YES ------------------------------------ |
write COMMIT decision to its log
| --- COMMIT -------------> | |
| --- COMMIT --------------------------------> |
| <-- ACK ---------------- | |
| <-- ACK ------------------------------------ |
- Phase 1 (prepare / voting). The coordinator asks each participant to prepare. A participant makes its changes durable (but not visible), keeps its locks, and votes yes, or votes no if it cannot commit. A "yes" is a promise: it will commit if told to, even after a crash.
- Phase 2 (commit / abort). If all voted yes, the coordinator logs a commit decision and tells everyone to commit; otherwise it tells everyone to abort.
The weakness is blocking: if the coordinator crashes after participants voted yes but before they hear the decision, they cannot decide on their own and must hold their locks until the coordinator recovers. 2PC also adds round trips to every transaction. Consensus-based approaches (replicating the coordinator's decision with Paxos or Raft, as Spanner and CockroachDB do) remove the single point of failure. Many microservice systems avoid distributed transactions altogether with the saga pattern (a sequence of local transactions with compensating actions), covered in transactions and CDC.
Relational databases expose 2PC as XA transactions (MySQL XA START, XA PREPARE, XA COMMIT) and PostgreSQL's PREPARE TRANSACTION / COMMIT PREPARED.
NewSQL
NewSQL (also called distributed SQL) databases aim to combine the SQL interface and ACID transactions of relational databases with the horizontal scaling of NoSQL. They shard data automatically, replicate each shard with a consensus protocol (Raft or Paxos), and run distributed transactions across shards.
| System | Notes |
|---|---|
| Google Spanner | Globally distributed; uses TrueTime (GPS and atomic clocks) to order transactions with external consistency |
| CockroachDB | PostgreSQL-compatible wire protocol; Raft per range; serializable by default |
| TiDB | MySQL-compatible; separate SQL layer and TiKV storage with Raft |
| YugabyteDB | PostgreSQL-compatible query layer on a distributed document store |
The cost is latency: a transaction touching several regions waits for consensus round trips, and these systems are operationally more complex than a single PostgreSQL server. They make sense when you truly need SQL semantics at a scale or availability level a single primary cannot provide.
SQL or NoSQL: a decision table
| Situation | Lean towards | Reason |
|---|---|---|
| Money, orders, inventory, bookings: strong invariants | Relational (PostgreSQL, MySQL) | ACID transactions, constraints, joins |
| Data is relational, queries are varied or unknown upfront | Relational | Ad hoc SQL, joins, mature tooling |
| Data fits on one large server (most apps, up to terabytes) | Relational, with read replicas | Simplest system that works |
| Cache, sessions, counters, leaderboards | Key-value (Redis) | Sub-millisecond lookups by key |
| Records with varied, nested, evolving attributes read as a whole | Document (MongoDB) or JSONB in PostgreSQL | Flexible schema, no joins for aggregates |
| Massive write volume, time series, known access by key | Wide-column (Cassandra) | Linear write scaling, partitioned by key |
| Multi-hop relationship queries (friends of friends, fraud rings) | Graph (Neo4j) | Traversals without join explosions |
| Must stay writable during regional partitions | AP stores (Cassandra, DynamoDB) | Availability over immediate consistency |
| Global scale and SQL transactions | NewSQL (Spanner, CockroachDB) | Distributed ACID, at a latency and cost premium |
| Full-text search, relevance ranking | Search engine (Elasticsearch, OpenSearch) beside the main DB | Inverted indexes |
Most real systems are polyglot: PostgreSQL as the source of truth, Redis as a cache, a search engine for search, and perhaps a wide-column store for event logs.
Interview tip
In a design interview, default to a relational database and justify any departure with a specific access pattern or scale requirement: "Orders go in PostgreSQL because we need transactions; the activity feed goes in Cassandra partitioned by user because it is write-heavy and always read by user id." Choosing NoSQL "because it scales" without saying what scales and why is a red flag.
Common mistake
"NoSQL is schemaless" is misleading. The schema still exists; it moves from the database into application code, which must handle every historical shape of a document. Schema-on-read is flexibility you pay for later, so document the expected structure and validate it (MongoDB supports JSON Schema validation on collections).
Interview questions
Q1. What is NoSQL and why was it created?
NoSQL is a family of non-relational databases (key-value, document, wide-column, graph) built for horizontal scaling, flexible schemas and data models that match specific access patterns. It emerged when web-scale applications outgrew single-server relational databases. The usual trade-off is giving up some joins, multi-record transactions or strong consistency.
Q2. Compare document and key-value stores.
A key-value store treats the value as opaque and only supports lookups by key, which makes it extremely fast and simple. A document store understands the structure of each document, so you can query and index fields inside it. Use key-value for caches and sessions; use documents for entities with rich, varying structure that you query by several fields.
Q3. When would you use a graph database?
When the core queries traverse relationships many hops deep, such as friend-of-friend suggestions, fraud rings sharing devices, or dependency analysis. Each hop in a relational database is a join whose cost grows with table size, while a graph database follows stored pointers. For simple CRUD or large aggregations, a relational database is usually better.
Q4. State the CAP theorem correctly.
In a replicated data store, when a network partition occurs, the system must choose between consistency (every read sees the latest write, linearizability) and availability (every non-failed node answers). Partition tolerance is not optional in a distributed system, so the real choice is CP or AP during a partition. When there is no partition, a system can provide both.
Q5. What does PACELC add to CAP?
It adds the normal case: else, when there is no partition, the system trades latency against consistency. Strong consistency requires waiting for other replicas, which adds latency. Cassandra is PA/EL (available and fast, eventually consistent), while Spanner is PC/EC (consistent always, paying latency).
Q6. What is eventual consistency?
A guarantee that if no new updates are made, all replicas will eventually converge to the same value. It says nothing about how long that takes or what intermediate reads return. Systems often add session guarantees such as read-your-writes and monotonic reads to make it usable.
Q7. Compare leader-follower, multi-leader and leaderless replication.
Leader-follower sends all writes through one leader, so there are no write conflicts but write capacity is limited and failover needs care. Multi-leader lets several nodes accept writes, useful across regions, but requires conflict resolution. Leaderless writes to and reads from quorums of replicas, staying available when some replicas fail, but needs read repair and gives weaker ordering guarantees.
Q8. What is a quorum, and why R + W > N?
With N replicas, a write waits for W acknowledgements and a read queries R replicas. If R + W > N, any read set and any write set share at least one replica, so the read will see the latest acknowledged write. For example N = 3, W = 2, R = 2 tolerates one failed replica for both reads and writes.
Q9. Synchronous vs asynchronous replication?
Synchronous replication waits for replicas to confirm before acknowledging, so no acknowledged write is lost on failover, at the cost of latency and availability if a replica is slow. Asynchronous replication acknowledges immediately, which is fast but causes replication lag and can lose recent writes if the leader fails. Semi-synchronous waits for one replica as a compromise.
Q10. How do you provide read-your-writes consistency with read replicas?
Route a user's reads of data they recently changed to the leader, either always for their own data or for a short window after each write. Alternatively, track the log position of the user's last write and only read from replicas that have caught up to it, or pin the session to one replica.
Q11. Range vs hash partitioning?
Range partitioning keeps contiguous key ranges together, so range scans are efficient, but sequential keys such as timestamps create hotspots. Hash partitioning spreads keys evenly but makes range queries hit every shard. The choice depends on whether your main queries are ranges or point lookups.
Q12. What problem does consistent hashing solve?
With hash mod N, changing N moves most keys (80 % when going from 4 to 5 nodes). Consistent hashing places nodes and keys on a ring so that adding or removing a node only moves the keys in one arc, about 1/N of them. Virtual nodes smooth out the distribution and allow weighting by capacity.
Q13. How does two-phase commit work and what is its weakness?
The coordinator asks all participants to prepare; each makes its changes durable, holds locks and votes. If all vote yes, the coordinator logs and sends commit, otherwise abort. Its weakness is blocking: if the coordinator fails after the votes, prepared participants hold their locks until it recovers. It also adds latency to every transaction.
Q14. What is NewSQL?
Distributed SQL databases such as Spanner, CockroachDB, TiDB and YugabyteDB that provide SQL and ACID transactions while scaling horizontally. They shard data automatically and replicate shards with Raft or Paxos. They cost more latency and operational complexity than a single-node database.
Q15. How would you choose between SQL and NoSQL for a new product?
Start from access patterns and invariants. If the data is relational, needs transactions or ad hoc queries, and fits on one powerful server with replicas, choose a relational database. Choose a NoSQL store for a specific need: key-value for caching, wide-column for massive partitioned writes, documents for flexible aggregates, graphs for deep traversals. Many systems combine them.
Key takeaways
- NoSQL trades some relational features (joins, multi-row transactions, strong consistency) for horizontal scale, flexible schemas and access-pattern-specific models.
- Key-value for lookups by key (Redis), document for nested aggregates (MongoDB), wide-column for partitioned write-heavy data (Cassandra), graph for traversals (Neo4j).
- CAP: during a partition choose consistency or availability; P is not optional.
- PACELC adds the everyday trade-off: latency versus consistency.
- BASE (basically available, soft state, eventual consistency) is the AP counterpart of ACID; session guarantees such as read-your-writes make it practical.
- Replication: leader-follower, multi-leader (conflicts), leaderless (quorums, R + W > N); synchronous is safe but slow, asynchronous is fast but lags.
- Sharding: range (good for ranges, risk of hotspots) or hash (even spread); consistent hashing moves about 1/N of keys on resize.
- 2PC gives atomic commit across nodes but blocks if the coordinator fails; NewSQL adds consensus to make distributed SQL practical.
- Default to a relational database and choose NoSQL for a stated reason.
Next lesson
Continue with Database design practice.

