All resources
Daily Learning · Building Microservices: Designing Fine-Grained Systems

Why one giant database hits a wall — and how a "ring" of many small ones doesn't

Monday, 13 July 2026
Why one giant database hits a wall — and how a "ring" of many small ones doesn't

🎯 You'll understand why a single big database eventually chokes, how spreading data across many machines fixes it, and the real price you pay for that fix: consistency becomes your application's problem.

Let me tell you a story that happened at Amazon in 2004, right in the middle of holiday shopping. The databases behind their catalog started maxing out — at the exact moment they needed to work most. That pain led to one of the most important ideas in modern computing. Let's build it up slowly, one piece at a time.


Step 1: The comfortable starting point — one big relational database

Amazon began with a very reasonable choice. They put their data in big relational databases (they used Oracle).

What is a relational database?

A relational database stores data in tables — like neat spreadsheets — and lets you connect (or relate) those tables together with a language called SQL.

Example, in plain SQL:

SELECT * FROM carts WHERE user_id = 123;

That means: "Give me the cart belonging to user 123."

Why was this a smart bet? Because relational databases are:

  • Battle-tested — they'd been used for decades.
  • Transactional — they promise all-or-nothing changes (more on that soon).
  • Familiar — every engineer already knew SQL.

For years, this worked beautifully.


Step 2: The traffic curve goes vertical

Then Amazon grew. Fast. More and more requests hit the same database.

Here's the key thing to notice: the problem was NOT that queries were slow. The problem was the shape of the setup.

        millions of requests
                 │
                 ▼
        ┌─────────────────┐
        │  ONE PRIMARY DB  │   ← everything funnels here
        └─────────────────┘

What is a "primary node"?

The primary node is the single machine that's allowed to accept writes (changes) and is the source of truth. In a classic setup, there's just one of them.

And that one machine becomes the ceiling for the whole system. Everything has to pass through it.


Step 3: Why you can't just buy a bigger machine

The obvious fix is: make that one machine stronger. Add more CPU, more memory. This is called vertical scaling.

What is vertical scaling?

Vertical scaling means making one machine bigger and more powerful. Think of it as buying a taller ladder.

The trouble is that ladders have a maximum height.

One big box has a ceiling. Eventually you're buying the biggest computer on Earth — and it still chokes.

Compare it to a restaurant with one super-chef. You can give that chef the best knives and the fastest stove, but there's still only one pair of hands. At some point, the only real fix is more chefs.

That "more chefs" idea is called horizontal scaling — adding many machines instead of one giant one. Hold that thought.


Step 4: The clue hidden in the data

Amazon studied their own traffic (their postmortems — written reviews after things break). They found something surprising:

About 70% of all database access was simple key-value lookups.

What is a key-value lookup?

A key-value lookup is the simplest kind of data request: "Here's a key, give me the value." Like a coat check: hand over ticket , get back coat . No searching, no combining.

Examples of what Amazon was doing 70% of the time:

  • "Get this cart." → key = user, value = their cart
  • "Get this session." → key = session id, value = login info

None of these needed the fancy powers of a relational database. No joins, no complex transactions.

What is a join?

A join stitches two tables together — e.g. combining a users table with an orders table to see each user's orders. It's powerful, but expensive.

So here's the punchline: Amazon was paying the full cost of a heavy relational engine just to do what amounts to a hashmap lookup.

What is a hashmap?

A hashmap is a basic programming tool that stores key→value pairs and finds them instantly. It's the simplest, cheapest way to "get this by that."

They had a Ferrari engine pulling a wagon.


Step 5: Where the real breakdown happened — strong consistency

There was a deeper problem than wasted effort. The thing that actually fell over under peak load was the coordination needed to keep everything perfectly correct. That coordination is called strong consistency.

What is strong consistency?

Strong consistency means: the moment you write something, every reader everywhere sees the new value immediately. There's no lag.

That sounds ideal — until you realize the cost. To guarantee it, all the machines must constantly check in with each other and agree before anyone can move on. Under massive traffic, all that agreeing becomes the bottleneck.

Write "cart updated"
      │
      ▼
 Node A ──"do you agree?"──▶ Node B
      ◀──"agreed"───────────
      │  (wait for everyone...)
      ▼
   OK, now the read is allowed

At peak holiday load, this waiting-and-agreeing is exactly what buckled.


Step 6: The 2007 fix — the Dynamo paper

In 2007, Amazon published the Dynamo paper and rebuilt their system around a new philosophy. The change was almost a change of belief:

Stop demanding that every read sees the very latest write.

This is a huge deal, so let's unpack the three moves they made.

Move 1: Embrace eventual consistency

What is eventual consistency?

Eventual consistency means: after a write, readers might briefly see the old value — but given a little time, everyone catches up and agrees.

Analogy: you update your home address. Your bank sees it instantly, but a magazine you subscribe to might mail one more issue to the old address before catching up. It's not wrong forever — just briefly behind.

Move 2: Replicate across many nodes with consistent hashing

Now, remember "more chefs"? Amazon spread the data across many machines arranged in a ring, and decided which machine holds which data using consistent hashing.

What is consistent hashing?

Consistent hashing is a rule for placing data on machines. You take the data's key, run it through a math function (a hash) that turns it into a number, and that number points to a spot on a ring of machines. Same key → same machine, every time.

What is a hash?

A hash is a function that turns any input into a fixed scrambled number. hash("USER#123") might become 842910. The same input always gives the same output.

Picture the ring (this is the diagram at the heart of Dynamo):

              A
        H           B
      G               C
        F           D
              E

 hash("USER#123") → lands between C and D
 → data lives on node D

Each node owns a slice of the ring. To find your data, you hash the key and walk to the node responsible for that spot. No single machine is the ceiling anymore — the load is spread around the whole ring. This is horizontal scaling.

The neat part of consistent hashing: if one node dies or you add a new one, only a small slice of data moves, not everything.

Move 3: Let writes always succeed — even during partitions

What is a partition?

A partition is a network split — some machines temporarily can't talk to others, even though all are alive. Like a phone line going dead between two offices.

In a strongly consistent system, a partition can freeze writes (because machines can't agree). Dynamo made a different call:

Writes always succeed. Reads may return stale data — briefly.

For a shopping site, that's the right trade. Better to always accept "add to cart" than to reject a sale because two machines couldn't sync for a second.


Step 7: How this looks today — DynamoDB in practice

The modern child of that paper is DynamoDB. Here's the mechanics in real terms.

You choose a partition key — the key that decides where your data lives.

What is a partition key?

The partition key (often written PK) is the value that gets hashed to pick which shard stores your row.

What is a shard?

A shard is one slice of your data living on one group of machines. Sharding = splitting your data into pieces spread across the ring.

PK = "USER#123"
        │
   hash("USER#123")
        │
        ▼
   lands on shard 4  ──▶  stored there, read from there

So the flow is: pick a partition key → data is sharded by its hash → you accept that a read might briefly return stale data (an out-of-date value).


Step 8: The part most people skip — someone still pays

Here's the honest catch. Dropping strong consistency didn't make the hard problem disappear. It moved the problem — up out of the database and into your application layer (the code your own team writes).

What is the application layer?

The application layer is your own program — the code your engineers build on top of the database.

The specific problem your team now inherits is conflict resolution.

What is conflict resolution?

Conflict resolution is deciding what to do when two changes to the same thing disagree.

Real example — two writes to the same shopping cart:

Phone:   add "headphones"   ─┐
                             ├─▶  Which cart is correct?
Laptop:  add "charger"      ─┘

With strong consistency, the database would have sorted this out. Now your code must decide. Amazon's rule of thumb: keep the added items rather than lose a sale. So the merged cart holds both headphones and charger.

As Sam Newman (a well-known voice on microservices) points out: the moment you drop strong consistency, your developers inherit this merging work. It doesn't vanish — it changes hands.

Amazon made that trade with their eyes open:

A slightly wrong cart beats a cart that won't load.

Step 9: The one idea to carry everywhere

Here's the deep truth underneath all of it:

Consistency isn't a feature you switch on. It's a cost — you either pay it in the database engine, or you hand it to your application. Someone always pays.

If the engine pays: writes can stall, and the whole thing has a ceiling. If the app pays: it scales beautifully, but your developers must handle stale reads and merges.

Neither is "free." Good engineering is choosing who pays on purpose.


The core lessons

  • One big database has a ceiling. Vertical scaling (a bigger machine) eventually runs out, because everything funnels through one primary node.
  • Match the tool to the work. Amazon found 70% of their access was simple key-value lookups — no need for a heavy relational engine there.
  • Strong consistency is the hidden bottleneck. Making every machine agree instantly is exactly what collapses under peak load.
  • Spread the load with a ring. Consistent hashing places each key on a node by its hash, so many machines share the work (horizontal scaling) and losing one node only moves a little data.
  • Eventual consistency = always-on, briefly-behind. Writes always succeed, even during a network partition; reads may return stale data for a moment.
  • The hard problem moves, it doesn't vanish. Drop strong consistency and your application inherits conflict resolution — like merging two versions of a cart.
  • Someone always pays for consistency. Either the database engine pays (and hits a ceiling) or your app pays (and stays scalable). Choose deliberately.

Want this kind of thinking applied to your business?

Book a free 30-minute discovery call — we’ll show you your highest-value first automation, no jargon, no obligation.

Book a discovery call