Learn Labs
7. Sharding

7.5 Sharding and Secondary Indexes

Secondary indexes

index type
8 shards
read fan-out
1 shard
write fan-out
read one term
write one document
Problem

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.

A secondary index has to be partitioned too, and there are only two choices: partition it by document, or partition it by term. Each makes one operation cheap and the other expensive.

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–499Shard 1 — IDs 500–999
Records191 red Honda
306 black Honda
214 red Ford
515 red Nissan
768 silver Ford
893 red BMW
Local indexcolor: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:

SituationCost
You know the partition keySearch the appropriate shard only ✔
You want only SOME results, not allSend to any shard ✔
You want ALL results and don't know the partition keySend 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 shardTermPostings list
0 — colors a–r, makes a–fcolor:black[306, 806]
0color:red[191, 214, 515, 893] ◀ drawn from all data shards
0make:Ford[214, 768]
1 — colors s–z, makes h–zcolor:silver[768]
1make:Honda[191, 306]
1make: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:

QueryCost
Single condition (color = red)Read from ONE index shard to fetch the postings list ✔
Fetch the actual records, not just IDsStill 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 bythe record's partition keythe indexed value (term)
WriteONE 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 localstill multiple data shards
Read (multi-condition AND)each shard ANDs locallycross-shard postings-list intersection over the network
Scalabilitymore shards ≠ more query throughputscales with terms
Consistencynaturally consistenthard — needs distributed txn or accepts staleness
Used byMongoDB, Riak, Cassandra, Elasticsearch, SolrCloud, VoltDBCockroachDB, TiDB, YugabyteDB; DynamoDB (both)

On this page