Designing Data-Intensive Applications
Ch. 6

Partitioning of Key-Value Data

Partitioning splits a dataset across nodes so each machine handles a manageable subset.

A single machine cannot hold all data forever. Partitioning (sharding) divides the key space so each node owns a subset. The partition scheme determines load balance and query efficiency.

Diagram
In practice

Kafka topics split into partitions for parallelism — each partition is an ordered log. MongoDB sharding routes documents by shard key through mongos routers. DynamoDB partitions by primary key hash; a hot partition (e.g., celebrity user ID) can throttle the whole table until you add a random suffix.

Uber at scale

Trip events partition by ride_id so all state transitions for one trip land on the same Kafka partition — preserving order. Driver location updates hash by geo-cell to spread load across shards.

typescript — Hash partition routing
// Kafka / DynamoDB-style partition routing
function partitionForKey(userId: string, partitionCount: number): number {
  let hash = 0;
  for (const ch of userId) hash = (hash * 31 + ch.charCodeAt(0)) >>> 0;
  return hash % partitionCount;
}

// WhatsApp routes messages by chat_id to the same shard for ordering
const partition = partitionForKey(chatId, 128);
Key Takeaways
  • Partitioning by key range assigns contiguous key ranges to nodes — risk of hot spots.
  • Hash partitioning distributes keys evenly but loses range-scan efficiency.
  • Skewed workloads need composite keys or random suffixes to spread hot keys.
  • Partitioning is almost always combined with replication for fault tolerance.
  • The partition function is hard to change after deployment.
  • Kafka partitions, MongoDB sharded clusters, and DynamoDB partition keys all shard by key.
partitioningshardingKafkaMongoDBDynamoDBhash partitioninghot spot