Part 4 · 1 chapters · ~12 min
Partitioning
Splitting data across machines: key-range and hash partitioning and compound keys, why mod N fails, fixed partitions and consistent hashing for rebalancing, hot spots and their fixes, and local versus global secondary indexes.
9
Keys, rebalancing, hot spots and secondary indexes
split data so load spreads and queries stay local
- Key ranges: cheap range scans, but sequential keys create a hot spot at the end.
- Hashing spreads load evenly. A compound key gets both even spread and ordered scans.
- Never use mod N over the number of nodes.
- Use fixed partitions or consistent hashing so adding a node moves only part of the data.
- Hot spots: split the hot key, cache it, or isolate it. Watch the hottest partition, not the average.
- Secondary indexes: local or global, chosen by which is more frequent, reads or writes.
code
// DynamoDB single-table design: partition by account, sort by time
// PK = "ACCT#7f2" SK = "TXN#2026-10-06T09:12:44Z#tx_91a"
await ddb.query({
TableName: 'ledger',
KeyConditionExpression: 'PK = :pk AND begins_with(SK, :month)',
ExpressionAttributeValues: { ':pk': 'ACCT#7f2', ':month': 'TXN#2026-10' },
}); // one partition, ordered: a statement in one query
// a hot merchant: write-sharding the counter
const shard = Math.floor(Math.random() * 10);
await ddb.update({ Key: { PK: `MERCH#42#${shard}`, SK: 'DAILY#2026-10-25' },
UpdateExpression: 'ADD total_kobo :a', ExpressionAttributeValues: { ':a': 150000 } });
// read: sum the 10 shardswhen to shard at all
A single well-indexed Postgres primary on a large instance handles tens of thousands of writes per second and terabytes of data. Shard when you have measured that you need to, and shard by the key your most frequent queries already filter on, usually tenant or account.
PARTITIONING: SPLITTING DATA ACROSS MACHINES
key ranges, hashes, consistent hashing, rebalancing, hot spots and the secondary-index problem
swipe the figure sideways, or tap expand for full screen
1/6
key range
By key range: partition 1 holds keys A to F, partition 2 G to M, and so on (Bigtable, HBase, CockroachDB, Spanner splits). Range scans are efficient (all transactions for one account in date order are adjacent), but sequential keys (timestamps, auto-increment ids) send every new write to the last partition: a hot spot.