Partitioning & Consistent Hashing

8.Splitting the Load Without Losing the Plot

M

In this chapter

We'll trace a naive partitioning scheme's real cost — adding one machine moving most of the dataset — and meet consistent hashing as the real, general technique behind Act 5's own Cassandra ring, keeping data movement small and localized when a cluster changes size, plus the data-locality payoff that buys.

9–11 min

The Problem in Real Life

GreenMart's catalog keeps growing past what any single machine can comfortably hold, in every region — exactly the horizontal-scaling path the earlier chapter already named. Sarah now has to actually decide how the data gets split across many machines, in a way that doesn't fall apart every time one machine gets added or removed.

This isn't a new problem. It's the exact same one Act 4's DynamoDB and Act 5's Cassandra already solved concretely, one database at a time — this chapter is where the general technique underneath both of those finally gets its own name.

S

I've built this before, twice, without ever naming what I was actually doing.

Sarah

A Naive Split vs. A Split That Survives Change

Partitioning splits data across machines

Sharding — the same idea Act 4's DynamoDB and Act 5's Cassandra already built, now named generally.

Naive hashing reshuffles almost everything

Verified: adding one machine to a 3-machine naive split moved 3 of 5 real items — a near-total reshuffle.

Consistent hashing keeps movement small

The exact mechanism behind Cassandra's own ring — adding/removing one machine only affects a small, adjacent portion.

Partition keys feed the whole mechanism

DynamoDB's own term — the specific field whose hash decides where on the ring an item actually lands.

Partitioning & Consistent Hashing

Partitioning (also called sharding) is the real, deliberate split: dividing a dataset across multiple machines, each one responsible for a real, distinct portion. This is the exact same idea Act 4's DynamoDB partition keys and Act 5's Cassandra ring already built concretely — this chapter is the general name for the technique both of those were specific instances of.

The real, hard problem partitioning has to solve isn't the initial split — it's what happens afterward, when a machine gets added or removed, which at real scale happens constantly, not rarely. As little data as possible should have to move when that happens.

Table — A Naive Split, Traced Through an Added Machine
Item Keyhash(key) mod 3 (3 machines)hash(key) mod 4 (4th machine added)Did It Move?
gift-box-0111No
gift-box-0222No
gift-box-0303Yes — moved
gift-box-0410Yes — moved
gift-box-0521Yes — moved

Adding just one machine (3 → 4) changed the assigned machine for 3 of these 5 real items — not a small, proportional share, a near-total reshuffle. At real catalog scale, this naive approach means adding one machine can mean moving almost the entire dataset at once.

Consistent hashing is the real fix, and it's the exact mechanism behind Act 5's Cassandra ring, generalized: instead of a hash directly deciding a machine number (which changes completely whenever the machine count changes), it maps both machines and data onto positions on a conceptual ring — Cassandra's own token ranges, precisely. Each machine owns the stretch of the ring between itself and the next machine going around. Adding a new machine only means it claims a portion of the ring near its own new position — only the data that used to belong to that specific stretch has to move, not data anywhere else on the ring. Removing a machine works the same way in reverse: its stretch of the ring is absorbed by its neighbor, and only that stretch's data has to move.

Data locality is the real payoff this buys, beyond just avoiding reshuffles: with a consistent, predictable placement scheme, related data (and the machines serving it) can be kept physically close to where it's actually used — the same real motivation behind GreenMart running each of its own regions with local copies in the first place. Partition keys — DynamoDB's own term, already met concretely in Act 4 — are exactly what gets fed into this whole mechanism: the specific field whose hash decides where on the ring a given item actually lands, which is exactly why choosing a good partition key was the single most consequential decision on a DynamoDB table, and exactly the same real decision Cassandra's own partition key made too.

Key Takeaway

Consistent hashing solves a real, specific problem naive partitioning doesn't: keeping the amount of data that has to move, when the cluster changes size, small and localized instead of nearly total — verified directly above, where adding one machine to a naive 3-machine split moved 3 of 5 real items, a near-total reshuffle for such a small change. Act 5's Cassandra ring was never a one-off design — it's this exact, general technique, and horizontal scaling only stays practical at real scale because of it.

Why This Matters

This is the actual mechanism behind "add more machines, not a bigger one" staying practical as GreenMart's real catalog and traffic keep growing — every earlier Act's own partition key decision (DynamoDB, Cassandra) was really this exact technique, applied once per specific database.

GreenMart now has the general name and real mechanism behind splitting data across many machines without a full reshuffle every time one changes. None of this has addressed who's actually in charge of coordinating all these machines once the network comes back together, though — exactly where the next chapter goes.

Next