DEV Community

NARESH
NARESH

Posted on

Designing a Production RAG Retrieval System for 500 Million Vectors

Designing a Production RAG Retrieval System for 500 Million Vectors

Banner

Assume there is a fictional company called Atlas.

Atlas started as a fairly typical enterprise RAG product. Companies uploaded internal documents, Atlas converted them into embeddings, stored them in a vector index, and used retrieval to provide relevant context to an LLM. With around 10 million vectors, the architecture was simple enough to operate and fast enough for the traffic it was receiving.

Then the product started growing.

Larger customers brought millions of documents with them. Existing customers began indexing more of their internal knowledge. Uploads were happening throughout the day while search traffic was increasing at the same time. Atlas also had to keep tenant data isolated, respect document permissions, handle updates and deletions, and make newly uploaded information searchable without making users wait for a large indexing job to finish.

The team estimated that the retrieval layer could eventually reach around 500 million vectors.

They could have responded by choosing a larger vector database and hoping it scaled. That would avoid the more important question.

What would the system itself need to look like at that scale?

In this article, I want to design Atlas the way I would approach a real system design problem. We will start with the architecture that works today, push it until something fails, and introduce a new component only when that failure gives us a reason to.

By the end, those individual decisions will come together into a complete production retrieval architecture for Atlas.


Before Designing Anything, Lock the Requirements

Before changing the architecture, Atlas needs to be clear about what it is actually designing for.

For this system, I would start with a few working assumptions:

  • the corpus can grow to roughly 500 million vectors
  • documents are uploaded throughout the day
  • search traffic is much higher than write traffic
  • each request belongs to a tenant and may include metadata or ACL filters
  • newly uploaded documents should become searchable quickly
  • updates and deletes must not corrupt older versions
  • search should continue even if an individual node fails
  • retrieval quality matters, so lowering latency by aggressively sacrificing Recall@K is not acceptable

There are also two very different workloads hiding inside the same product.

The first is the write path:

Document
→ processed
→ embedded
→ stored
→ becomes searchable
Enter fullscreen mode Exit fullscreen mode

The second is the read path:

Query
→ find relevant vectors
→ fetch the corresponding chunks
→ provide them to the LLM
Enter fullscreen mode Exit fullscreen mode

Keeping those paths separate in our head is useful because they optimize for different things. Writes care about durability, freshness, and background indexing. Reads care about latency, recall, filtering, and avoiding unnecessary work.

With those requirements in place, we can now look at Atlas's current 10-million-vector setup and ask the first useful system-design question:

What breaks if we try to stretch the same architecture all the way to 500 million vectors?


What Atlas Looks Like Today

Before redesigning Atlas, it helps to see the version that already works.

At its current scale, Atlas does not need a very complicated retrieval architecture. A customer uploads documents, the system chunks them, generates embeddings, stores the vectors with metadata, and uses those vectors during retrieval to provide context to the LLM.

V1 Architecture

That is good enough when the corpus is still manageable, the working set fits comfortably within the available infrastructure, and the rate of change is not too high.

At this stage, Atlas can afford a relatively straightforward setup. Search runs against one logical retrieval layer, ingestion is continuous but still manageable, and the operational overhead stays low. There is no strong reason to introduce a router, shard management, tiered storage, or a more elaborate write path before the workload demands it.

That matters because this is how a lot of systems actually begin. They do not start with a full distributed retrieval architecture. They start with the simplest design that works, and they evolve only when a clear bottleneck appears.

For Atlas, the first bottleneck shows up when the current setup is stretched toward 500 million vectors.


Failure #1: One Search Unit Stops Making Sense

The current Atlas design is fine while the corpus is still around 10 million vectors. There is one logical vector store, the write path feeds into it, and the read path searches it directly.

The first serious problem appears when Atlas starts pushing toward 500 million vectors.

A 768-dimensional float32 vector takes about 3 KB. At 500 million vectors, the raw vector data alone is roughly 1.5 TB. Once we include the ANN structure, metadata, replicas, caches, and spare capacity, keeping the entire search system on one machine becomes difficult to justify.

So the first redesign is about capacity.

Atlas splits the corpus into shards and adds a router in front of them.

                 ┌→ Shard 1
Query → Router ──┼→ Shard 2
                 ├→ Shard 3
                 └→ ...
Enter fullscreen mode Exit fullscreen mode

Each shard owns only part of the corpus and builds its own local search index. As the dataset grows, Atlas can add more shards instead of continuously scaling one machine vertically.

This also gives us a natural place to think about storage differently. The hottest search state can remain close to memory, while larger vector data and index structures can move to cheaper storage such as SSD.

At this point, Atlas has solved the first scaling limit: the corpus no longer has to fit inside one search unit.

But this redesign creates the next problem.

Those shards are optimized for search, while customers are still uploading documents all day. If every new vector forces Atlas to modify the optimized index immediately, the write path starts interfering with the read path.


Failures 1 and 2: Capacity Stops Scaling, and Writes Start Fighting With Reads

Atlas works well enough with its current architecture while the corpus is still around 10 million vectors. There is one logical retrieval layer, the write path feeds into it, and queries search it directly.

The first problem appears when Atlas starts moving toward 500 million vectors.

A 768-dimensional float32 vector takes about 3 KB. At that scale, the raw vectors alone are already around 1.5 TB. Once metadata, ANN structures, replicas, and spare capacity are added, keeping the entire retrieval system as one search unit becomes difficult to justify.

So the first redesign is about capacity. Atlas stops treating the corpus as one large index and starts partitioning it across shards, with a router in front to direct traffic to the right places. That gives the system a way to grow horizontally instead of relying on one increasingly expensive machine.

That solves the first bottleneck, but it exposes the next one.

Atlas is not indexing a static corpus. Customers are uploading documents throughout the day, which means the system is constantly producing new chunks, new embeddings, and new vector records. If every incoming write tries to update the fully optimized search structures immediately, ingestion starts competing with query traffic for CPU, memory, and I/O.

Shard

At that point, Atlas needs to separate accepting a write from fully optimizing it for search.

The write path becomes more structured. A document is chunked, embedded, and written with metadata and version information into a durable log. From there, the new vectors land in a mutable segment, where they can become part of the searchable corpus quickly. The heavier index construction work happens later, in the background, when that data is turned into an optimized segment.

This gives Atlas a much healthier shape.

Sharding solves the capacity problem. The mutable path solves the write-contention problem. Atlas no longer has to keep the entire corpus on one machine, and it no longer has to rebuild or heavily mutate its read-optimized structures for every upload.

That leaves the next question.

If recent data is living in mutable segments while older data is already in optimized segments, how should the read path search both without sacrificing freshness or wasting work?


Failures 3 and 4: Fresh Data Gets Missed, and Global Search Becomes Wasteful

Atlas now has shards, mutable segments for recent writes, and optimized segments for older data.

The next problem is freshness.

If a query searches only the optimized ANN segments, a document uploaded a few seconds ago may not appear yet. The write succeeded, but the background indexing step has not finished.

So each shard has to search both places:

Shard Query
→ Mutable Segment
→ Optimized Segment
→ Merge Local Candidates
Enter fullscreen mode Exit fullscreen mode

That keeps recent writes visible without forcing Atlas to wait for full ANN optimization.

Once this works, another inefficiency becomes obvious.

Atlas is a multitenant product. A query from one customer should not search every vector in the platform. The request already carries useful context such as tenant ID, document permissions, and metadata filters, so Atlas can use that information before doing expensive retrieval.

The read path now becomes:

Query
→ Authentication + Tenant Context
→ Query Router
→ Relevant Shards
→ Mutable + Optimized Search
→ Local Candidates
Enter fullscreen mode Exit fullscreen mode

Query → Auth/Tenant Context → Router → Relevant Shards → Mutable + Optimized Search → Local Candidates

For a small tenant, that may mean touching only a limited part of the cluster. A larger tenant can span several shards, but the router still avoids sending the request to places that cannot contain valid results.

ACLs and metadata filters are applied as part of that search strategy, not as an afterthought after scanning the entire corpus.

With these two changes, Atlas solves two different problems at once.

Searching mutable and optimized data preserves freshness. Tenant-aware routing and filtering reduce unnecessary fanout.

But once a query starts touching several shards, Atlas gets a new problem: every shard has its own idea of the best results. Someone now has to combine them into one global ranking.


Failure #5: Distributed Search Creates a New Bottleneck

Once Atlas starts routing a query to multiple shards, each shard returns its own local candidates.

That creates a new problem. The best result inside Shard A is only the best result inside Shard A. Atlas still needs to compare candidates from every participating shard before it can produce the final top results.

So the read path gets one more component:

multiple shards → local candidates → candidate aggregator → global top-K → reranking.

The aggregator collects candidates from each shard, compares them using a common scoring strategy, and builds the global result set.

This works, but distributed search introduces a latency problem that did not exist before.

If a query touches ten shards, the response can be delayed by the slowest one. A single overloaded shard can push the p99 latency of the whole request higher even when the rest of the cluster is healthy.

Hot tenants make this worse. One customer may generate enough traffic to overload the shard group holding most of its data while the rest of the cluster remains underused.

Atlas now needs to think about routing and capacity together. Large tenants may need to be split across more shards, frequently accessed data may benefit from caching, and replicas can help absorb additional read traffic.

The important point is that sharding solves the storage problem, but it creates a coordination problem.

By this stage, Atlas is no longer running several independent vector indexes. It is running one distributed retrieval system that has to produce a single answer from many local searches.


What Happens When Something Actually Fails?

So far, Atlas has been designed for scale. Now it needs to survive failure.

Suppose a shard accepts new vectors and crashes before those vectors are fully indexed. The WAL becomes the recovery point. A replacement process can replay the accepted operations, rebuild the mutable state, and continue indexing without asking the customer to upload the document again.

Search availability needs a separate answer. If one shard disappears, Atlas should not depend on that single copy of the data. Replicas allow another node to continue serving reads while the failed shard recovers.

There is also a less dramatic failure mode that matters just as much: indexing falls behind.

If documents arrive faster than Atlas can process them, the mutable segments and indexing queue keep growing. Allowing that backlog to grow without control eventually hurts search latency too.

So Atlas needs backpressure.

Ingestion rate > indexing capacity
→ queue grows
→ slow or limit new ingestion
→ protect search traffic
Enter fullscreen mode Exit fullscreen mode

This is an important shift in the design. Reliability is not only about recovering from a crashed server. It is also about making sure one overloaded part of the system does not drag the rest of the retrieval path down with it.

At this point, Atlas has enough pieces to assemble the complete production architecture.


Putting It Together: The Atlas Architecture

Atlas V2 Architecture

Atlas now looks very different from the system we started with.

The original design had one simple retrieval layer. The production version has evolved into three cooperating parts: ingestion, retrieval, and distributed storage.

The ingestion side is responsible for turning continuously changing documents into searchable state:

Upload
→ Chunking
→ Embedding
→ WAL
→ Mutable Segment
→ Background Indexing
→ Optimized Segment

Enter fullscreen mode Exit fullscreen mode

The retrieval side begins with the user's identity and tenant context, routes the query only to the shards that matter, searches both recent and optimized data, and combines the results:

Query
→ Query Embedding
→ Auth + Tenant Context
→ Shard Router
→ Shard-local Search
→ Candidate Aggregator
→ Global Top-K
→ Reranking
→ ACL Validation
→ Fetch Chunks
→ LLM
Enter fullscreen mode Exit fullscreen mode

Underneath both paths sits the distributed storage layer.

Each shard owns part of the corpus, replicas provide redundancy, and different parts of the search state can live on different storage tiers depending on how frequently they are accessed.

The important part is that none of these components were added just to make the architecture look more sophisticated. Every one of them exists because the simpler version eventually hit a specific limit.

That is the architecture Atlas would now take into production.


Walk One Request Through the Final System

At this point, the architecture is easier to understand by following two real flows through it.

First, take a document upload.

A customer uploads policy.pdf. Atlas chunks the document, generates embeddings for those chunks, attaches metadata and version information, and writes the operation into the WAL. From there, the chunks enter a mutable segment so they can become part of the searchable corpus quickly. Later, the background index builder converts that data into an optimized segment and folds it into the longer-lived search structure.

Now take a search request.

A user asks, "What is our refund policy?" Atlas generates the query embedding, attaches the tenant and access context, and sends the request through the shard router. Only the relevant shards are searched. Inside each shard, Atlas looks at both the mutable segment and the optimized segment, then produces local candidates. Those candidates are sent to the candidate aggregator, which forms the global top results. The reranker refines that result set, ACL validation makes sure the returned chunks are actually visible to the user, and the final chunks are passed to the LLM.

What matters here is that the architecture is no longer one large vector index with a query attached to it.

Atlas is coordinating ingestion, recent searchable state, distributed candidate generation, global ranking, and final authorization checks in one retrieval flow.

That is the real shift from the original 10-million-vector design. Search is still the center of the system, but it is no longer the whole system.


Conclusion

Atlas did not reach its final architecture by starting with shards, WALs, replicas, mutable segments, candidate aggregators, and tiered storage.

It started with a simple retrieval system that worked at 10 million vectors.

Then the workload changed.

The corpus became too large for one search unit, so Atlas introduced sharding. Continuous ingestion started interfering with read-optimized structures, so it separated mutable state from optimized state. Fresh data had to remain visible, so queries searched both. Multitenancy made global search wasteful, so routing became tenant-aware. Distributed search produced local answers, so Atlas added candidate aggregation and global ranking. Reliability then pushed the design toward WAL recovery, replicas, and backpressure.

The useful lesson is not the final diagram itself.

It is the sequence of decisions that produced it.

Good system design usually works this way. Start with the simplest architecture that satisfies the current workload. When a real constraint appears, understand what is failing and why. Then add the smallest architectural change that fixes that specific problem.

For Atlas, the final system looks complex only when viewed all at once.

Taken step by step, every component has a reason to exist.


📖 Blog by Naresh B. A.

👨‍💻 Backend & AI Systems Engineer | Distributed Systems · Production ML

🌐 Portfolio: [Naresh B A]

📫 Let's connect on [LinkedIn] | GitHub: [Naresh B A]

Thanks for reading. This is my personal engineering perspective, and I'd genuinely be interested in hearing where you agree or disagree. ❤️

Top comments (0)