Skip to the lesson
Little Builders system design, explained small

Advanced Database Architectures

Picture a school library that got too popular. One librarian cannot check out every book, and one room cannot hold them all. So the school opens more rooms and splits the books between them. It also makes copies of the books, so many kids can read at once and nothing is lost if one room floods. But copies take a moment to make, so sometimes a kid picks up a copy that is a little out of date. This page is about splitting data, copying it, finding the right room fast, and taming copies that run behind.

Horizontal sharding

Splitting the books between many library rooms

When one database gets too full or too busy, we split its rows across many databases. Each piece is called a shard. It is like a library putting books by authors A to H in room 1, I to P in room 2, and Q to Z in room 3. Each room holds fewer books and has its own librarian, so together they store more and can handle many more changes at once.

The value used to choose the room is the shard key, such as a user's name or ID. Range sharding gives each shard a range of keys, like A to H. Neighbors stay together, so a list of everyone from A to C is quick, but some ranges get crowded: if the key is the sign-up date, every new user lands on the newest shard. Hash sharding first scrambles the key into a number with a hash function, then uses that number to pick a shard. Keys spread out evenly, but a range question like A to C now has to ask every shard.

Choose the key with care. A hot shard is one room with a line out the door while the others sit empty. A celebrity with millions of fans can make one shard hot even with hashing, because all of their activity shares one key. Common fixes are splitting very busy keys into several pieces or giving them special treatment.

Sharding has real costs. A question that needs data from many shards must visit each room and combine the answers. A change that touches two shards at once needs extra teamwork to stay all or nothing. And moving rows when you add shards, called resharding, is slow and risky, which is one reason consistent hashing exists.

Three bookcases each hold one range of books, like shards each holding some rows. Point at a case to pull its books out.

Remember

Sharding gives more room and more writing power, so pick a shard key that spreads the work and keeps most questions inside one shard.

Leader and follower replication topologies

Who writes in the main notebook, and who copies it

Replication means keeping copies of the same data on several computers, so one breaking does not lose anything and many readers can be served at once. The big question is who is allowed to accept writes.

Single leader: one copy is the leader, and every write goes to it. The leader keeps an ordered list of changes, and the followers copy that list in order, like students copying the teacher's notes from the board. Followers can answer reads. If the leader breaks, a follower is promoted to be the new leader, which is called failover. This is simple and has no conflicts, but every write squeezes through one computer.

Multi-leader: several leaders accept writes, often one per region, such as one in Europe and one in America, and they send each other their changes. Writes are fast for everyone nearby, but two leaders can change the same thing at the same moment. That is a conflict, and the system must settle it, for example with last writer wins, which quietly throws one change away, or by merging both changes.

Leaderless: there is no boss. The writer sends each change to several copies and waits until W of them say saved. A reader asks R copies and keeps the newest answer. With N copies in total, if W plus R is more than N, the copies written and the copies read must share at least one, so the reader always reaches a copy that has the latest saved write. Dynamo-style databases like Cassandra work this way.

  • Synchronous copying: the leader waits until a follower confirms the change. Safer, but slower, and one stuck follower can hold up writes.
  • Asynchronous copying: the leader says done right away and copies afterward. Fast, but if the leader dies, its newest changes can be lost.
  • Semi-synchronous copying: wait for one follower, copy to the rest later. A middle path many systems use.
A big leader board takes every new line, and small follower slates copy it. Point at a slate to send the newest change to it.

Remember

One leader is simple, many leaders need conflict rules, and no leader needs reads and writes that overlap.

Consistent hashing

Servers sitting around a clock face

With hash sharding, the simple way to pick a shard is to scramble the key into a number, divide by the number of servers, and use the remainder. Grown-ups write this as hash mod N. It works well until you add a server. Change N from 3 to 4 and most keys get a new remainder, so about 3 of every 4 keys must move. That is like reshelving almost the whole library because one new room opened.

Consistent hashing places everything on a circle, like a clock face. Each server is hashed to a spot on the circle, and each key is hashed to a spot too. To find a key's home, start at the key and walk clockwise until you meet a server.

Now add a new server. It lands on one spot and takes over only the keys between it and the server before it, which is one arc of the circle. Every other key stays exactly where it was. On average only about 1 out of N keys moves, instead of nearly all of them. Removing a server is just as gentle: its keys simply move on to the next server clockwise.

With only a few spots, some servers get lucky and own a huge arc while others own a tiny one. The fix is virtual nodes: each real server sits at many spots around the circle, often a hundred or more. The arcs average out, so the load is even, and when a server joins or leaves, its share is spread over many servers instead of landing on one neighbor.

Keys sit round a dial and each walks clockwise to the next server. Move the new server round: only the keys in one arc go to it.

Remember

Walk clockwise to the next server, so adding a server only moves the keys in one arc.

Replication lag

The copy that is a few seconds behind

With asynchronous copying, followers are always a little behind the leader. Usually that is a few milliseconds, too short to notice. But when the system is busy, or a follower is slow or far away, the gap can grow to seconds or even minutes. That gap is called replication lag.

Lag causes strange moments. You post a photo, and it is saved on the leader. You refresh the page, and your read happens to go to a follower that has not copied the photo yet. Your photo seems to vanish. Nothing is broken, the copy is just late, but it feels broken.

If nobody writes for a while, every follower catches up and all copies agree. That is why this style is called eventual consistency. The work is to hide the confusing moments in between, without giving up the speed that followers bring.

The leader stack gets each new block first, and the follower copies it a moment later. Raise the pointer to add blocks and watch it lag.

Remember

A follower can be behind, so a read from a follower may show the past.

Replication lag mitigation

Little rules that stop time from going backward

Read your own writes: after you change something, read your own things from the leader for a short while, say a minute. Or remember the spot in the change list where your write landed, and only read from a follower that has caught up to that spot. Everyone else can keep using followers.

Monotonic reads: imagine refreshing a page and seeing 3 comments, then refreshing again and seeing only 2, because the second read went to a follower that is further behind. Time went backward. The fix is to keep each user on the same follower, for example by choosing it from their user ID, so their view only moves forward.

Consistent prefix reads: if a question is written before its answer, everyone should see them in that order. When data is split across shards, the question and the answer may sit on different shards with different lag, so a reader could see the answer first. The fix is to keep related writes in one ordered stream, such as the same shard, or to track which writes depend on which.

Finally, watch the lag itself. Measure how far behind each follower is, take very slow followers out of rotation until they catch up, and for data that must never look old, like a bank balance, use synchronous or semi-synchronous copying or read from the leader. Each of these costs some speed, so use them only where it matters.

A reader keeps one seat at one copy desk, so the pages never jump back in time. Point at a desk to make it the reader's seat.

Remember

Read your own changes from an up-to-date copy, keep each user on one copy, and keep related writes in order.

Quick recap

  1. Sharding splits rows across databases by a shard key, adding room and write power, but questions across shards and resharding are hard.
  2. Range sharding keeps neighbors together, hash sharding spreads evenly, and a poor key makes a hot shard.
  3. Single leader is simple, multi-leader must settle conflicts, and leaderless needs W plus R to be more than N.
  4. Synchronous copying is safer but slower. Asynchronous copying is fast but can lose the newest changes.
  5. Consistent hashing puts servers on a ring, so adding one moves only about 1 out of N keys. Virtual nodes even out the load.
  6. Followers lag behind. Read-your-writes, monotonic reads, consistent prefix reads and lag monitoring hide the confusing moments.

Grown-up words

and what they mean in plain words

Shard
One piece of a big database, holding some of the rows.
Shard key
The value that decides which shard a row lives in.
Hot shard
A shard that gets far more work than the others.
Replica
A copy of the data kept on another computer.
Leader
The copy that accepts writes and tells the others what changed.
Follower
A copy that repeats the leader's changes and can answer reads.
Failover
Promoting a follower to leader when the leader breaks.
Quorum (W and R)
How many copies must confirm a write (W) or answer a read (R).
Virtual node
One of many spots a single server takes on the hash ring.
Replication lag
How far a follower is behind the leader.