7.5 Sharding and Secondary Indexes
Secondary indexes
- 8 shards
- read fan-out
- 1 shard
- write fan-out
Document-partitioned (local). Writes touch one shard, so ingest is cheap and consistent. Reads must scatter/gather across all 8 shards, and the request is as slow as the slowest shard — the tail-amplification problem again.
Everything so far relies on the client knowing the PARTITION KEY. That works in key-value models where the partition key is the first part of (or all of) the primary key.
A secondary index usually doesn't identify a record uniquely — it's a way of searching for occurrences of a value: find all actions by user 123, find all articles containing "hogwash", find all cars whose color is red.
The problem with secondary indexes is that THEY DON'T MAP NEATLY TO SHARDS. Two approaches: local and global.
5.1 Local secondary indexes (document-partitioned)
Each shard independently maintains its own secondary indexes, covering only the records in that shard.
| Shard 0 — IDs 0–499 | Shard 1 — IDs 500–999 | |
|---|---|---|
| Records | 191 red Honda306 black Honda214 red Ford | 515 red Nissan768 silver Ford893 red BMW |
| Local index | color:red → [191, 214]color:black → [306]make:Honda → [191, 306] | color:red → [515, 893]color:silver → [768]make:Ford → [768] |
Each shard indexes only the records it holds, so red cars appear in both local indexes — and finding all of them means visiting both shards.
Writes are easy: you deal only with the shard containing the record you're writing. (When a red car is added, that shard automatically adds its ID to the postings list for color:red.)
Reads:
| Situation | Cost |
|---|---|
| You know the partition key | Search the appropriate shard only ✔ |
| You want only SOME results, not all | Send to any shard ✔ |
| You want ALL results and don't know the partition key | Send the query to ALL shards and combine results — a scatter/gather ✗ |
This makes read queries on secondary indexes QUITE EXPENSIVE. Even querying shards in parallel, it is prone to TAIL LATENCY AMPLIFICATION (Ch 2 §2.6). It also LIMITS SCALABILITY: adding more shards lets you store more data, but it DOESN'T INCREASE QUERY THROUGHPUT if every shard has to process every query anyway.
Nevertheless widely used: MongoDB, Riak, Cassandra, Elasticsearch, SolrCloud, VoltDB.
⚠️ Warning on rolling your own: if your database supports only key-value, you may be tempted to implement a secondary index in application code. Take great care to ensure the indexes remain consistent with the underlying data. RACE CONDITIONS AND INTERMITTENT WRITE FAILURES (where some changes were saved but others weren't) CAN VERY EASILY CAUSE THE DATA TO GO OUT OF SYNC (Ch 8).
5.2 Global secondary indexes (term-partitioned)
Construct a global index covering data in all shards. But you can't store it on one node — it would become a bottleneck and defeat the purpose. So THE GLOBAL INDEX MUST ALSO BE SHARDED — but it can be sharded DIFFERENTLY from the primary-key index.
Primary data, sharded by record ID: shard 0 holds IDs 0–499, shard 1 holds 500–999, shard 2 holds 1000–1499.
Global secondary index, sharded by the indexed value — the term:
| Index shard | Term | Postings list |
|---|---|---|
| 0 — colors a–r, makes a–f | color:black | [306, 806] |
| 0 | color:red | [191, 214, 515, 893] ◀ drawn from all data shards |
| 0 | make:Ford | [214, 768] |
| 1 — colors s–z, makes h–z | color:silver | [768] |
| 1 | make:Honda | [191, 306] |
| 1 | make:Nissan | [515] |
(A "term" generalizes the full-text-search notion of a keyword to mean any value you can search for in the secondary index.)
The index shard can hold a contiguous range of terms, or terms can be assigned by HASH of the term.
Reads:
| Query | Cost |
|---|---|
Single condition (color = red) | Read from ONE index shard to fetch the postings list ✔ |
| Fetch the actual records, not just IDs | Still have to read from all the data shards responsible for those IDs ✗ |
| Multiple conditions (color AND make; or multiple words in a text) | Those terms will likely be on DIFFERENT index shards. To compute the logical AND, find all IDs in BOTH postings lists. Fine if the lists are short — but if they're long, it can be slow to send them over the network to compute the intersection ✗ |
Writes are the real problem:
Writing a single record might affect MULTIPLE SHARDS OF THE INDEX (every term in the document might be on a different shard). This makes it harder to keep the secondary index in sync with the underlying data. One option is a DISTRIBUTED TRANSACTION to atomically update the shards storing the primary record and its secondary indexes (Ch 8).
Used by CockroachDB, TiDB, YugabyteDB. DynamoDB supports BOTH local and global.
In DynamoDB, writes are ASYNCHRONOUSLY reflected in global indexes, so READS FROM A GLOBAL INDEX MAY BE STALE — similar to replication lag (Ch 6).
Nevertheless, global indexes are useful IF READ THROUGHPUT IS HIGHER THAN WRITE THROUGHPUT, and if the postings lists are not too long.
5.3 The comparison
| Local (document-partitioned) | Global (term-partitioned) | |
|---|---|---|
| Index sharded by | the record's partition key | the indexed value (term) |
| Write | ONE shard ✔ | Several index shards ✗ (may need a distributed transaction, or async ⇒ stale reads) |
| Read (single condition) | ALL shards — scatter/gather ✗ | ONE shard for the postings list ✔ |
| Read (fetch records) | already local | still multiple data shards |
| Read (multi-condition AND) | each shard ANDs locally | cross-shard postings-list intersection over the network |
| Scalability | more shards ≠ more query throughput | scales with terms |
| Consistency | naturally consistent | hard — needs distributed txn or accepts staleness |
| Used by | MongoDB, Riak, Cassandra, Elasticsearch, SolrCloud, VoltDB | CockroachDB, TiDB, YugabyteDB; DynamoDB (both) |