Distribute it — shards, scatter-gather, the LLM stack
Past one machine you split vectors across nodes and ask every node for its local best, then a coordinator merges to the global best — exact only if each node returns enough candidates; this is the box that sits beside the KV store in every 2026 LLM app.
Hybrid search gave us a feature-complete engine on one machine — but scene 1's wall returns: a billion vectors don't fit on one box. So we shard across nodes and scatter-gather: every shard returns its local best, a coordinator merges the global best. And once it's distributed, it's worth stepping back to see where this entire engine sits in a real 2026 system.
Scene 14
Distribute it — shards, scatter-gather, and the LLM stack
- Watch
- Try it
- Predict
- Capture
A billion vectors can't live on one machine, so we split — shard — them across several nodes, and copy (replicate) each shard so a node failure doesn't lose data. Each shard is a complete mini vector index from earlier scenes (an HNSW graph or an IVFPQ index). When a query arrives, the coordinator fans it out to every shard at once; each shard searches its own slice and returns its local closest songs. Watch the three shards report back, one by one. Notice the asymmetry that makes this affordable: building those indexes is slow and expensive, but it's done once and amortized — queries are cheap and run hot, millions of times over the same built index.
Highlighted lines are the ones running in the diagram right now.
def search(query, k, candidates):targets = route(query)# scatter: ask all shards in parallelreplies = parallel_map(targets,lambda s: s.local_search(query, candidates),)# gather: pool every returned candidatepooled = flatten(replies)pooled.sort(key=lambda c: dist(query, c))return pooled[:k]
def local_search(query, candidates):# self.index is a full mini vector index:# an HNSW graph or an IVFPQ index over this slicehits = self.index.query(query)hits.sort(key=lambda c: dist(query, c))# only the closest `candidates` leave the shardreturn hits[:candidates]
def route(query):if sharding == 'clustered':# similar vectors share a shard, so a query# only needs the nearest bucketsreturn nearest_shards(query)# random spread: every query must hit all shardsreturn all_shards
Where this sits in Build a vector database (Pinecone / Weaviate / pgvector style)
Scene 14 of 15, in the Production act — Filters, the graph-disconnection trap, hybrid RRF, and sharded scatter-gather.. Shard vectors across nodes, scatter-gather every query, merge per-shard top-k — exact only if each shard over-fetches. This is the box that sits next to the KV store in every 2026 LLM app.
Up next. You've now built every piece: the vector and metric, the indexes (Flat, IVF, HNSW, PQ, IVFPQ), the trilemma, filters, hybrid, and distribution. The last skill is choosing among them for a real workload — RAG, recommendation, semantic cache, or billion-on-a-budget — with every knob traceable to the scene that justified it.
All 15 scenes in Build a vector database (Pinecone / Weaviate / pgvector style) · Every curriculum