When partitioning a search index across machines, how does a document-partitioned layout differ from a term-partitioned one?
answer
- which axis of the map is cut
- self-contained shard vs whole postings list
- fan-out vs cross-machine intersection
- one write vs many writes
- skewed term popularity
basics
~20 sA document-partitioned index gives each shard a complete index over some of the documents, so every query hits every shard. A term-partitioned index gives each shard whole postings lists for some of the terms, so a query hits only its terms' shards.
solid answer
~40 sIn a **document-partitioned** (local) layout, each shard is a self-contained index over its own documents. A new document goes to one shard, load stays even, and multi-term intersections run locally. The cost is that a plain keyword query must fan out to every shard. In a **term-partitioned** (global) layout, each term's complete postings list lives on one shard, so a two-term query touches at most two shards. That layout brings other costs. Intersecting lists that sit on different machines means shipping postings across the network. Indexing one document touches the shard of every distinct term it contains. And very common terms create huge, hot shards. Fan-out is paid in small parallel messages, while these costs are harder to engineer around, so most general-purpose search systems partition by document.
go deeper
Recall the two cuts. Document partitioning gives each shard some of the documents, and term partitioning gives each shard some of the terms with their full postings lists.
Explain the query path and the write path under each layout: fan-out to all shards versus touching only term owners, and one shard per write versus one shard per distinct term.
Show the operational consequences: hot shards for common terms, postings shipped between machines for intersections, what a lost shard removes, and why teams make fan-out cheap instead.
Frame the choice as paying a predictable, parallel fan-out cost versus structural skew and multi-shard writes. Know where hybrid layouts earn their complexity.
## Two ways to cut an inverted index An **inverted index** maps each **term** to a **postings list**: the documents that contain the term, often with positions and frequencies. When the index no longer fits on one machine, it can be cut along either axis of that map: - **By document.** Each machine gets a subset of the documents and builds a complete index over just those. - **By term.** Each machine gets a subset of the terms and holds the complete postings lists for those terms, covering all documents. Textbooks call these the **local** (document-partitioned) and **global** (term-partitioned) index layouts. ## Document partitioning Each shard is a small, self-contained search engine. - **Query path:** any shard may hold matches, so a keyword query must visit every shard. The coordinator then merges the per-shard top-k lists (scatter-gather). - **Write path:** a new or updated document goes to exactly one shard, chosen by hashing its ID or by a routing key. - **Balance:** documents hash evenly, so shard sizes and loads stay similar. - **Multi-term queries:** intersections and scoring run locally on each shard. No postings cross the network, and only small top-k lists travel. - **Failures:** losing a shard removes one slice of documents from every result, not every result for some term. ## Term partitioning Each term's whole postings list lives on one shard, which is usually chosen by hashing the term. - **Query path:** a query touches only the shards that own its terms, so a two-term query touches at most two shards. - **Multi-term queries:** postings lists for different terms may sit on different machines. Intersecting them means shipping lists or partial results across the network, often in a pipeline that starts from the rarest term. - **Write path:** indexing one document touches the shard of every distinct term in it. A long document with thousands of distinct terms can touch many shards, which makes atomic updates hard. - **Balance:** term popularity is heavily skewed, so the shards that own common terms hold the longest lists and receive the most queries. - **Failures:** losing a shard removes every result for the terms it owns. ## Side by side | Property | Document-partitioned | Term-partitioned | |---|---|---| | Shards touched per keyword query | all of them | only those owning the query's terms | | Shards touched per document write | one | one per distinct term (possibly many) | | Data moved across machines per query | small top-k lists | possibly whole postings lists | | Load balance | even by construction | skewed by term popularity | | Effect of losing a shard | some documents missing | some terms missing entirely | | Adding capacity | add shards, move documents | move terms, split long lists | ## Why document partitioning is the usual default The fan-out cost of document partitioning is paid in small messages that run in parallel. Techniques for that cost are well understood: replicas, latency-aware replica selection, hedged requests, deadlines, and routing when a query names a partition key. The costs of term partitioning are structural. Hot terms produce hot shards, intersections produce bulky transfers, and writes touch many machines. For these reasons most general-purpose search systems partition by document and invest in making fan-out cheap. ## Where term partitioning still shows up - **Single-term or rare-term lookups**, where touching one shard instead of all of them is a real saving. - **Hybrid designs** that keep a global structure for a small, special set of terms or keys and a document-partitioned index for everything else. - **Very common terms**, whose lists are often split further by document-ID range across machines. That split brings fan-out back for exactly the terms that are queried most, which is one reason pure term partitioning rarely holds up at scale. Partitioned databases face the same local-versus-global choice for secondary indexes. The general theory is covered separately; this question is about serving keyword search queries.
- In a term-partitioned search index, why is indexing a single new document expensive?Each distinct term in the document has its postings list on the shard that owns that term, so one document write fans out to many shards. A long document can have thousands of distinct terms. Keeping all those updates atomic needs cross-shard coordination, or the system must tolerate a window where the document is findable by some terms but not others. In a document-partitioned index, the same write touches exactly one shard.
- How does a term-partitioned search index cope with a very common term?A very common term has a huge postings list and receives a large share of queries, so its shard becomes both oversized and hot. The usual responses are to replicate that shard more heavily and to split the list by document-ID range across several machines. Splitting brings fan-out back for the most-queried terms, which erodes the main advantage of term partitioning.
saying these in an interview costs you the question
- Document partitioning lets a keyword query skip shards that lack its terms.
- Term partitioning makes indexing cheaper because each document goes to one shard.
- Term partitioning removes all cross-machine work from multi-term queries.
- Term-partitioned shards stay balanced because terms hash evenly.
- The two layouts differ only in shard naming, not in query or write cost.