Loading
0x60Lesson 7 of 13

Partition data with sharding

Split data across machines and add nodes without reshuffling everything.

20 min 6-question quiz 1 code exercise
By the end of this lesson you can
  • Compare range and hash partitioning.
  • Choose a shard key that avoids hotspots.
  • Explain how consistent hashing limits data movement.

When data or write traffic outgrows one machine, you shard (partition) it: each node holds a subset of the rows. The shard key decides where each row lives, and it is one of the most consequential choices in a design. A good key has many distinct values, spreads load evenly, and keeps data that is queried together on the same shard.

Partitioning strategies

  • Range partitioning: split by key ranges (A-F, G-M, …). Range scans are efficient, but sequential keys such as timestamps send all new writes to one shard.
  • Hash partitioning: place rows by hash(key). Load spreads evenly, but range queries must ask every shard.
  • Directory-based: a lookup service maps keys to shards. Flexible, but the directory must be fast and highly available.

The naive hash scheme hash(key) % N has a nasty property: when N changes, most keys map to a different node, forcing a massive data migration.

design.py
1keys = range(12)
2before = {key: key % 3 for key in keys}
3after = {key: key % 4 for key in keys}
4moved = sum(1 for key in keys if before[key] != after[key])
5print(f"{moved} of 12 keys moved")
Output
9 of 12 keys moved

Consistent hashing places both nodes and keys on a ring of hash values. Each key belongs to the first node clockwise from it. Adding a node only takes over the keys between it and its predecessor - about 1/N of the data - and removing one hands its keys to the next node. Each physical node usually owns many virtual nodes on the ring so load stays even. Systems such as Cassandra and DynamoDB are built on this idea.

Key takeaways

  • Pick a high-cardinality shard key that matches your main queries.

  • Range partitioning helps scans but risks hotspots; hashing spreads load but scatters ranges.

  • Consistent hashing with virtual nodes moves only ~1/N of keys when membership changes.

Lesson quiz

6 questions · pass with 5 correct · up to 50 XP

Passing this quiz completes the lesson and keeps your streak going. Questions you miss come back in review sessions later.

Practice: simulate system design building blocks

Use small Python programs to estimate capacity and simulate caches, load balancers, hash rings, and rate limiters. These exercises run locally in your browser.

Exercise 1

Look up keys on a hash ring

+25 XP

Read a line of nodes as position:name (in any order), then a line of key hashes. For each key, print hash -> name, where the owner is the node with the smallest position greater than or equal to the hash. If no node is at or after the hash, wrap around to the node with the smallest position.

  • Three nodes
  • Unsorted nodes
main.py
Loading editor…

Python runs in a sandboxed browser worker with a 60 second time limit. Its runtime loads from the Pyodide CDN; your code stays in this browser.

Questions about this lesson

Stuck? Ask. Figured something out? Share it. Explaining is one of the best ways to learn.

Loading posts…

Did you like the lesson? 😆👍
Consider a donation to support our work: