Atlas Infinite: MongoDB's Disaggregated Architecture

Atlas Infinite, which dropped in public preview on September 29 alongside MongoDB 9.0, delivers a significant architectural update payload: compute-storage disaggregation. Separating compute from storage is not new, and there are well-known ways to do it. But what makes Atlas Infinite interesting is how MongoDB departs from the standard design to keep WiredTiger's efficiency and to ensure that the shared storage layer never sees customer data unencrypted.

I am giving a high level overview of the architecture here distilling from the MongoDB Engineering blog post by Abhishek Chauhan, who led the effort. For detailed coverage, I refer you to that post. It is a great read, and parts 2 and 3 will drop soon.


Why not earlier, and why now

Disaggregated databases have existed for some time, so why did MongoDB wait? This is mainly because MongoDB didn't need disaggregation to scale out: when a workload outgrew one machine, sharding spread it across many nodes in a seamless manner.  MongoDB was among the first databases to ship sharding with automatic chunk balancing and routing. Since horizontal scaling has been a first-class feature for over a decade,  MongoDB never faced the single-node limit that pushed other systems toward  disaggregation early. That bought us time to do disaggregation right, without regressing on WiredTiger's performance or its security model.

Disaggregation is a good idea for several reasons anyways, and MongoDB adopted it to address customer needs. For one, storage gravity has become the  bottleneck for very large clusters. In Atlas's shared-nothing architecture, each node in a three-node replica set has its own disk with a full copy of the  data. That means, operations like growing storage, adding a replica, or restoring a backup, would require copying the whole dataset. On a large cluster  that can take hours when customers expect minutes. And since compute and  storage come bundled, needing more of one means paying for more of both. This  is why customers want compute and storage to scale independently. Disaggregation helps in right-sizing the customers and giving them flexibility for scaling. Finally, recent agentic development patterns want cheap snapshots,  clones, and branches of production data.


How MongoDB departs from the traditional disaggregated design

The standard approach works something like this. Compute is stateless. It ships a physical redo log (a record of every change to every page) to a quorum-replicated log service. Page servers read that log and rebuild pages by replaying the changes, so they can produce any page as of any point in the log. Object storage sits underneath as the durable copy.

MongoDB adopted much of this approach, but it broke with the traditional design in several places in order to keep WiredTiger's performance characteristics and the security guarantees customers expect.

Let's start with WiredTiger, one of the most sophisticated storage engines in production. WiredTiger achieves high performance by deferring and batching work. It is built around checkpoint-based consistency: rather than logging each page modification one by one, it lets writes accumulate in memory. A page can absorb many mutations and get written to disk just once, at a checkpoint. This amortizes the cost of turning writes into readable state. (This is increasingly where the field is heading as modern engines keep converging on copy-on-write, multi-version, checkpointed designs.) Forcing WiredTiger to log every page change for disaggregation sake would cripple its efficiency.

Secondly, because WiredTiger is lock-free and parallel and evicts pages speculatively, there is no canonical page image: two nodes holding the same data hold different bytes. So "replicate the pages" has no single correct target to replicate toward, which breaks from the traditional disaggregated design.

Finally, MongoDB refused to let shared storage read customer data. A storage tier that rebuilds pages from a log must be able to read them, so security ends up resting on operational controls. The better approach is to treat unencrypted user data as toxic. Holding it is a big responsibility, so keep it encrypted for as much of its lifetime as you can. That matters most when storage is shared across tenants, where one mistake no longer stops at one customer.

Because MongoDB owns the engine, the server, and the cloud service, it could co-design across all three layers. While the traditional design makes storage smart (by reading the data and rebuilding any page on demand), MongoDB went the other way and made storage blind to preserve WiredTiger's efficiency and to guarantee that storage can never read customer data unencrypted.


The architecture

Compute (primary, hot standby, and optional read replicas) runs queries and transactions and keeps only caches and recent changes, all of which can be rebuilt from storage if the node is lost. The shared storage layer, written in Rust, has three parts. A Log Service does only one thing, consensus on the log, which keeps it simple and steady. A Page Service spreads each customer's pages across many servers in every zone, with page materializers feeding it from the log. Object storage is responsible for long-term durability, organized by an Object Index Service.


Difference 1: two logs instead of one. A write is durable once the oplog, MongoDB's logical operation log, is on a majority of Log Service replicas across three zones. Separately, at checkpoint time, WiredTiger emits a phylog of physical page changes, mostly deltas. Page servers organize these by page after receiving them through the log service. This allows WiredTiger to keep its batching. Pages are written once per checkpoint, not once per change. Durability and page materialization are fully separate, with only the oplog on the critical path.

Difference 2: storage never sees your unencrypted data. In Atlas Infinite, data is encrypted on the compute node before any byte goes down, and decrypted only there, inside the customer's network, with keys the customer can hold and revoke. Storage can't see documents, collection names, indexes, or even tenant boundaries. Even obtaining root access to the whole storage fleet would yield an attacker only opaque bytes. Contrast this to the traditional design where storage must read data to rebuild pages, so security is partly "trust us" and partly "best intentions".

Difference 3: fresh replicas without page-level logging. In contrast to the traditional disaggregated approach where storage can build any page at any log point, in MongoDB pages only advance at checkpoints, so a replica reading only shared pages could potentially lag by seconds. This is easy to address however on the standby and read replica. Note that page servers build checkpoint-consistent pages from the phylog. On top of that, the compute node applies the oplog (as it arrives from the Log Service) to a local ingest table. Reads combine both to construct the fresh data. This way MongoDB gets millisecond replica lag while keeping checkpoint batching.


Finally, a word on availability. Every cluster has a hot standby, so a primary failure is a quick handoff rather than starting a new node. Data lives in three forms at once (the replicated log, the page servers, and object storage), the fleet is split into isolated cells so trouble in one can't cascade into the next, and recovery parallelizes across many machines rather than falling to one.

Comments

Popular posts from this blog

The Safest Job from AI may be Writing

The Two Abstractions of System Design: Hide or Reduce

The Agentic Self: Parallels Between AI and Self-Improvement

In Search of a Compositional Theory of Self-Stabilization

Learning about distributed systems: where to start?

Building a Database on S3

Foundational distributed systems papers

Cloudspecs: Cloud Hardware Evolution Through the Looking Glass

Specula: Scaling formal specifications for autonomous model checking of system code

Hints for Distributed Systems Design