Study notes · 12% of the exam

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. 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. 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. 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. 4

    hash mod N moves 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. 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. 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.

Test yourself on Data Modeling, Indexing and Partitioning

Ten questions, with the answer and explanation after each one.