Designing Data-Intensive Applications
Ch. 6

Partitioning and Secondary Indexes

Secondary indexes on partitioned data require either local indexes per partition or a global index.

Primary key partitioning is straightforward — route by key hash. Secondary indexes (search by email, filter by date) complicate routing because the indexed value may live on a different partition than the record.

In practice

MongoDB sharded clusters maintain local indexes per shard — a query on an unsharded field must scatter-gather to every shard. Elasticsearch distributes inverted indexes across nodes; a search hits all shards and merges results. PostgreSQL on Citus uses co-location and reference tables to reduce scatter-gather overhead.

Airbnb at scale

Search by city or neighborhood hits secondary indexes scattered across listing shards. Each shard returns partial results; the router merges and ranks them — scatter-gather is the cost of partitioning by listing_id.

typescript — Scatter-gather secondary index query
// Airbnb search — secondary index scatter-gather across shards
async function searchByCity(city: string): Promise<Listing[]> {
  const shards = await router.allShards();
  const partial = await Promise.all(
    shards.map((s) => s.query("SELECT * FROM listings WHERE city = $1", [city]))
  );
  return partial.flat().sort((a, b) => b.rating - a.rating).slice(0, 50);
}
Key Takeaways
  • Partitioned by document: each partition maintains its own secondary indexes.
  • Partitioned by term: the index itself is partitioned by the indexed value.
  • Global indexes enable efficient lookups but add coordination overhead.
  • Secondary index queries may scatter-gather across all partitions.
  • Index maintenance on writes adds latency to every insert and update.
  • Elasticsearch scatter-gather queries and MongoDB compound indexes illustrate the trade-offs.
secondary indexElasticsearchMongoDBscatter-gatherglobal indexlocal index