Data Modeling, Indexing and Partitioning
Choose storage engines, indexes and shard keys for the queries you actually run. Every choice makes some queries cheap and others expensive.
Key points
- 1
A composite index serves any leftmost prefix of its columns. Put equality-filtered columns first and the range or ORDER BY column last.
- 2
Every index is extra work on every write. B-trees update pages in place for predictable reads; LSM-trees buffer and flush sorted files for fast writes, paying with compaction and multi-file reads.
- 3
The shard key decides which queries stay on one shard. Hashing spreads writes evenly but destroys range locality; range partitioning keeps scans local but can create a hotspot on the newest range.
- 4
hash mod Nmoves almost every key when N changes (10 to 11 nodes moves about 91%). Consistent hashing moves about 1/N, and virtual nodes even out the load. - 5
A single hot key beats any partitioning scheme. Split it with key suffixes and aggregate on read, or batch writes before they reach storage.
- 6
Lookups by a non-key attribute need a global index partitioned by that attribute. Local per-shard indexes force a scatter-gather and can't enforce uniqueness.
Common traps
Scatter-gather latency follows the slowest shard: with 64 shards, about 47% of requests hit at least one shard's p99 (1 − 0.99^64).
Low-cardinality columns such as status are poor leading index columns on their own.
Partitions keyed only by an entity ID grow without bound for time-series data; add a time bucket to the partition key.
Read the source
Test yourself on Data Modeling, Indexing and Partitioning
Ten questions, with the answer and explanation after each one.