AnnouncementBuilt for unpredictable AI demand: Atlas Infinite and MongoDB 9.0 are here. Read more >  >>
AnnouncementMeet the intelligent data platform built for the AI era. Read more > >>
NewNow in public preview: Atlas Agent Engine, the secure way to run AI agents at scale. Read more > >>
Blog home
arrow-left

The Ground Beneath the Database

October 8, 2026 ・ 7 min read

From attached disks to a shared, elastic foundation—the story of disaggregating MongoDB and turning a database you manage into one you consume.

Summary

On September 29th, MongoDB announced the public preview of Atlas Infinite, the largest architectural innovation in MongoDB Atlas ever, built on disaggregated storage: compute and storage run on separate machines. Holding no data of their own, compute nodes become ephemeral and stateless. Durability, replication, and page serving move into a shared storage layer built from three services: a Log Service that reaches consensus on every write, a Page Service that keeps materialized pages warm in every zone, and an Object Index Service managing object storage underneath as the durability floor. From an application’s perspective, the wire protocol, the drivers, and the semantics do not change. 

With independently scaled storage decoupled from compute, storage scalability becomes virtually unlimited. Read replicas come up without the need to copy a disk. Vertical scaling can swap compute atop the same storage. Snapshots, clones, and restores become forks of a log: constant time, zero copy, at any database size.

The shared layer does all of this without the ability to read your data. Pages are encrypted on your compute node before a byte is handed down, with keys that can stay yours. This post covers the architecture. The next two posts will cover how a storage fleet operates on data it cannot understand, and what it takes inside the database engine to live without an attached disk.

Prologue

In October 2024, six weeks into my tenure at MongoDB, I found myself in a meeting with fellow server and support engineers, drawing straws. The data growth in a large MongoDB Atlas cluster had been accelerating. Now, that is great news: this is the problem you strive to have! However, our back-of-the-envelope estimate for storage expansion came to several days. The desired answer would have been several minutes. The question on the table was who would go and tell this to the customer.

I’ve thought about that meeting a lot. The database was behaving as designed. Like every traditional database, standard MongoDB Atlas is architected with coupled storage: each database node has its storage tightly bound to its compute.

It was this coupling that was in the way, because storage has gravity.

Data doesn’t want to move. The bigger it gets, the harder it pulls down. Anyone who has managed a large database knows the symptoms: scaling a cluster, resharding, rebalancing, restoring from backup, and adding a read replica. They are slow. They are slow because they all reduce to the same storage-shaped primitive: copy a very large amount of data from one place to another. Nothing about better software makes a terabyte weigh less.

And because compute and storage are coupled, regardless of whether you need more storage or more compute, you get more of both.

So we ask the obvious question: if the problem at hand is tight coupling, then would decoupling be the answer?

Turns out the answer is yes. And that answer is Atlas Infinite, a new deployment option within MongoDB Atlas, launched in public preview at on September 29th at MongoDB’s Investor Day. This post examines it in detail: what we set out to make possible; the architecture Atlas runs today and the disaggregated one we launched; the life of a write and of a read through the new machinery; what an append-only history buys you; what happens when things break; and the one line we would not cross to get any of that.

Decoupling is how a database stops being infrastructure you manage and starts being a capability you consume. But that's the marketing promise—here's the machinery that makes it real.

A peek under the hood

How do you move storage out from under a highly optimized, lock-free, parallel storage engine with non-deterministic, speculative eviction? Carefully.

Commercial databases built on decoupled compute and storage have existed for at least a decade. So it is reasonable to ask, "Why did it take MongoDB this long, and what’s new here?" The honest answer is that the properties that make our storage engine fast are exactly the ones that make it hard to pull apart—and that we refused to change it in ways that would give up the performance or the security our customers already had.

Besides, our customers have always had horizontal scale-out. When a workload outgrew a machine, sharding was the natural succession: one that several disaggregated architectures in the market lacked until recently. That freedom bought us the room to get this right rather than ship it first.

The challenge

WiredTiger™ is one of the most sophisticated storage engines in production. It earns its speed by deferring and batching work. It is built around checkpoint-based consistency: rather than logging each page modification, blow by blow, 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 a readable state, and it is increasingly where the field is heading, with modern engines converging on copy-on-write, multi-version, checkpointed designs precisely because of their efficiency. Something WiredTiger has been doing for years.

The decade-old mainstream disaggregation designs work by shipping a physical redo log that is obligated to record every change to a page, so the shared storage tier can reconstruct an arbitrary page version. Adopt that directly, and you’d have traded WiredTiger’s core efficiency for the newfound elasticity.

Furthermore, conventional disaggregation relies on the shared storage layer to read, compose, and rewrite your data in the clear. That is something we rejected. In Atlas Infinite, if your data is unencrypted cleartext, it is not shared. And if it is shared, it sure ain’t clear. We don’t hold cleartext in a shared service. Period.

The final wrinkle: because WiredTiger is lock-free, parallel, and evicts pages speculatively, two nodes holding identical logical data will not hold identical bytes on disk. Layout depends on timing, eviction order, and allocation state, so there is no canonical page image to ship and no single layout for copies to converge to. "Replicate the pages" has no well-defined target.

This is where being MongoDB helped. We are the one company with its hands on the whole stack: the storage engine, the database server, and the cloud that runs it. We didn't have to bolt disaggregation on from the outside and work around the engine—we could reach into the engine. This is what allowed us to solve something sharper than the textbook problem: keep WiredTiger's deferred, batched writes intact, keep the storage layer completely blind to customer data, and still deliver read replicas that lag the primary by milliseconds, not by checkpoint intervals. None of the existing recipes delivers all three at once.

The rest of this series is about what we had to invent to get there. We’ll take a glimpse at MongoDB Atlas Core’s architecture, followed by an in-depth look at the Atlas Infinite design, service decomposition, write and read paths, as well as how they inform our security, durability, availability, and performance: SDAP—the yardsticks by which we measure everything we do. Subsequent posts will expand further into the security and durability aspects in particular.

Let's start with a brief look at Atlas Core.

MongoDB Atlas Core: The attached-storage architecture

A MongoDB cluster is a replica set: three nodes, each a complete, self-sufficient machine. Each node holds a full copy of the data on its own attached disk. Each stays in step with the others by replicating the oplog: our logical operation log, from a primary to its secondaries. If the primary fails, a secondary is elected to replace it. Every node is doing every job: durability, availability, and serving queries, all bound together in the same box. It is shared-nothing. Its simplicity affords a near absolute security and isolation model: it simply relies on virtual private cloud (VPC) and virtual machine (VM) guaranteed isolation. One tenant’s actions cannot bleed into another.

But because compute and storage live on one machine, anything you want to do to the data, you have to do to the machine. Scaling reads means cloning an entire node before it can answer a query. High availability means paying for full physical copies that mostly stand and wait. Growing a node means rebuilding it and warming a cold cache. The design is simple, robust, and cleanly isolated tenant-from-tenant, and it makes every storage-shaped operation an act of moving the whole machine.

Figure 1. MongoDB Atlas Core: Three nodes across zones A/B/C, each a mongod bound to its own full-copy disk, joined by oplog replication.

A diagram illustrating MongoDB Atlas Core architecture, showing three nodes across Availability Zones A, B, and C. Zone A contains the Primary node, while Zones B and C each contain a Secondary node. Each node is connected vertically to its own dedicated storage disk containing a full copy of the data, and the nodes are joined horizontally across zones via oplog replication.

Atlas Infinite: The disaggregated architecture

We took a design so simple—three nodes in a replica set—and replaced it with at least fourteen! Database nodes, log nodes, page materializers, and page serving nodes. That sounds like the wrong direction, until you notice who runs them: not you. The cluster goes from three machines doing every job to fourteen sharing the work – extra hands, not extra burden. Managing data is hard. But done right, managing a million times more data is not a million times harder. Disaggregation splits the architecture into two halves.

The compute layer runs the database duties: parsing queries, running transactions, and handling aggregations. The change is that it becomes ephemeral and without storage: it becomes weightless. It offloads its other duties (oplog replication, consensus, page copying, and backups) to the storage layer. In place of Atlas Core’s three full replicas, Atlas Infinite runs a primary database node, a standby node for high availability, and optionally, additional read replicas for read scaling. Just as with Atlas Core, the compute layer runs isolated in its own VPC.

The storage layer is made up of multi-tenant storage layer services, responsible for managing storage for that “million times more data.” The machinery here is sophisticated, as we'll see. But none of it is yours to run. From the customer's side of the glass, storage is simply managed by MongoDB.

This layer consists of three main services: the Log Service, the Page Service, and the Object Index Service, each able to scale, fail, and recover independently. The entire storage layer was purpose-built from the ground up in Rust.

The Log Service is the shortest path to durability. Every write becomes an append to a replicated log. The service does exactly one thing: reach consensus on the log. And a service that does exactly one thing can do it with verifiable correctness, higher performance, and a tighter tail latency than one juggling durability alongside the full duties of a database.

The Page Service is a pool of page servers holding the materialized pages in each zone, serving page reads. A new compute node reads from a cache that is already warm instead of cloning a disk of its own. Every log’s pages are distributed across many page servers in each zone, so a read spike is absorbed across many instead of hammering one. This is the moment the fleet can finally bring more than one node's worth of muscle to bear on a single customer's workload.

A Page Materializer reads entries being written to the log and distributes them to the page server group responsible for that log. In parallel, it does one more thing: it continuously spools the log onto object storage.

Object Index Service and Object Storage—S3, Google Cloud Storage, Microsoft Azure Blob—sits underneath as the durability floor: designed to provide eleven nines, continuously maintained, the bottom of the world. Object storage is organized and managed by an Object Index Service, and access is mediated by an Object Read Proxy, which can answer the same shape of read queries as a page server—except from object storage. This is what allows us to deliver continuous data protection, time travel, and fast backups and restores.

Through all this, the front door does not move. Same wire protocol, same drivers, same read and write semantics. The application cannot tell the difference. Everything that changed is underneath it.

Figure 2. Atlas Infinite: Ephemeral compute on top, over three shared services—Log Service (consensus/durability), Page Service (shared warm cache), Object Storage (durability floor)—spanning zones A/B/C.

A diagram illustrating the MongoDB Atlas Infinite disaggregated architecture, showing a compute layer at the top with a Primary node in Zone A, a Standby node in Zone B, and an optional Read Replica in Zone C, separated by a security boundary from the multi-tenant storage services below. The storage layer comprises a Log Service for consensus and durability, a Page Service with warm page servers across all zones, and an underlying Object Index Service with Object Read Proxy resting on Object Storage (S3, GCS, or Azure Blob).

Control plane and heat management

I won't get into the details in this blog, but suffice it to say that the storage layer services are all provisioned, monitored, and heat-managed by control plane services that ensure automated fault tolerance and recovery as well as load spreading and tail-at-scale operations.

Scale anything independently

A quick reality check: if compute and storage are really separate, three operations that were expensive in the coupled world should become cheap without special-casing. Let's have a look.

Adding a read replica is zero-copy. In MongoDB Atlas Core, a secondary needs to clone a full copy of the disk before it can answer a single query, which is where the wait comes from. In Atlas Infinite, the data was never on the node to begin with. A new compute node starts up, tails the committed oplog from the Log Service to catch up to the present, and reads its pages from the shared Page Service, whose cache is already warm. Its cost is bounded by the length of the oplog tail rather than the size of the database, so it is serving reads in seconds regardless of how large that database is.

Vertical scaling swaps the compute without touching storage. In Atlas Core, resizing a node means standing up a bigger machine and then waiting out a cold cache, because both the cache and the data live on the machine you are replacing. In Atlas Infinite, they don’t. A larger compute node launches, takes over as primary, and the old one steps aside. Nothing moves, and nothing has to re-warm, and the swap closes in under a minute.

Storage performance scales on its own axis. With attached storage in Atlas Core, IOPS are a property of the attached disk. In Atlas Infinite, IOPS scale with the storage pool rather than with any one disk. If a workload is storage-bound, you request more throughput or more IOPS, and the storage layer grants it, with no need to touch compute at all.

The common thread, and the real test of the design, is that each of these touches only the layer that has to change, and its cost scales with the size of that change rather than the size of the database.

A deeper dive into the write path

Let us take a closer look at the life of a single write. With consensus and replication delegated to the Log Service, and the separation of durability from readability, there is a lot happening on the write path.

A write arrives at the primary. The first thing that happens, as always, is that the primary appends an entry to the oplog describing what changed. That record goes to the Log Service, which replicates it across the three zones, and the write is acknowledged the instant a majority of the log replicas hold it durably.

As the storage engine settles its pages in the ordinary course of checkpointing, it emits a second, physical log we call the phylog. Where the oplog says what changed in logical terms, the phylog carries the physical manifestation of the change: page images, and far more often, page deltas.

The standby and read replicas are clients of the Log Service and tail the log. At the same time the page materializer—one in each zone—tails the log as well, and distributes the phylog records across the page servers.

The page servers receive the phylog records, organize images and deltas by page, and store them on fast, directly attached drives.

In parallel, the page materializer uploads both the oplog and the phylog to object storage.

Notice the shape of everything downstream of the primary: logs—append-structured, versioned, immutable. Hold on to that; it is about to pay off.

Figure 3. Life of a write: critical path solid (client→primary→log→ack); async phylog + materialization dashed.

A diagram illustrating the write path in MongoDB Atlas Infinite, showing synchronous steps (solid lines) where a write goes from the client to the primary compute node, is appended and replicated in the Log Service to reach consensus and acknowledge the client, alongside asynchronous steps (dashed lines) where phylog records are sent to page materializers, routed to page servers across zones, and spooled to the Object Index Service and Object Storage.

The read path and solving standby lag: Layered tables

The read path on the primary is ... uninteresting, in that it is substantially the same as before. On the standby and read replicas, it gets interesting.

Checkpoint-based consistency in WiredTiger means that durable page state only moves forward at checkpoint boundaries. That's bad news for a standby or read replica that reads only materialized pages: it would trail the primary by as much as a full checkpoint interval: multiple seconds. That is a deal-breaker for a read replica serving reads. We wanted replica lag measured in milliseconds.

We accomplish this by layering two tables. Underneath sits the stable table: the checkpoint-consistent, shared, materialized data, the same page state the primary authors as it checkpoints. On top of it sits an ingest table (more a table cloth, really), into which the standby applies the oplog as it arrives. The ingest table holds the recent tail of changes—the writes that have landed since the last checkpoint the standby adopted. As newer checkpoints arrive and the stable table catches up, the standby garbage-collects the now-redundant entries from its ingest table, keeping the tail bounded.

A read composes both. The stable table supplies everything settled as of the last checkpoint. The ingest table supplies the fresh changes on top. Together, they give the standby a current, consistent view without waiting for the primary to publish its next checkpoint. That is what turns standby lag from a checkpoint-interval problem into an oplog-replication-lag one, measured in milliseconds.

Branching logs and ancestry

At the end of the write path I asked you to hold on to one property: everything downstream of the primary is a log, append-structured and versioned, immutable. Here is where that pays off.

When history is an append-only log, the storage-shaped operations a database dreads most no longer require moving data at all. A snapshot, a clone, a branch, a rewind, a restore: in the coupled world, every one of these is a copy. Over a versioned log, they are metadata. You are naming a point in history that already exists. Constant time, zero copy, independent of how large the database is.

Start with restoring a backup. A restore creates a new log derived from an existing one, branching off at the chosen point in time and proceeding independently from there. The original keeps going untouched; the restored copy is its own independent timeline from that instant forward. Do this a few times, and the logs form a lineage. We call it log ancestry: an ancestry tree rooted at the original log, each restore a new branch hanging off the point it was taken from.

Restore is not special. It is just one use of a more general primitive: take any log, mark a point, and fork. A snapshot is a marker on the timeline, nothing copied, just a named instant you can return to. A clone, or a writable snapshot, is a fork in that timeline: a branch that shares all history up to the fork point and then diverges.

Atlas Infinite sets the foundation for a production database you can branch like a Git repository: to test a migration against real data and throw the branch away. A rewind to any moment inside a retention window. A full-fidelity copy in every engineer’s harness, costing almost nothing to create because it copies almost nothing.

When things break

Drama, they say, needs three things: a villain, a victim, and a hero. A failing database has the first two in abundance. The villain is failure itself, in all its usual garbs: a node that dies mid-write, a network that partitions, a disk that goes bad, an availability zone that drops off the map. The victims are the machines snared. With Atlas Infinite, there is no hero in operations. We designed the system so the situation never rises to the level of needing one. Nothing rescues the database, because nothing put it in peril.

If a database primary fails, the standby is hot and ready to take over. If the network partitions, the Log Service mediates consensus, and the storage layer control plane may even be able to heal the raft group. If a log server fails, a new one gets commissioned, while writes never stop. If a page server fails, there are others with staged copies of the data. If a disk fails, there are other copies. If an entire availability zone goes dark, the cluster keeps serving: the log commits on a majority spread across three zones, each zone has its own page servers, and object storage that sits beneath it all, is zone failure-proof as well. The durability and availability guarantees stem from data that is held in three independent forms at once: the replicated log, the materialized pages, and the object store. And since the storage layer is multitenant, most recovery parallelizes many-to-many: many machines pitch in to take over the work of the departed machine, and to help a new node get up to speed.

That non-event is the visible tip of a great deal of deliberately unglamorous work that happens behind the scenes. The fleet is carved into cells, spread across availability zones, so trouble in one cell is contained and cannot cascade into the next. The converse also holds: a cell is designed to insulate itself from most trouble happening outside it. We call it Anti-Vegas: what happens outside a cell, stays outside a cell.

A resource scheduler watches the heatmap continuously, a tenant or a shard drawing more than its share, a node running over the utilization threshold, resource pressures of any kind, and moves load before anyone becomes a noisy neighbor, or feels one. Storage volumes are checked for integrity and repaired, with nothing visible from above.

None of this is drama. Boring is the highest compliment you can pay infrastructure. It means the system runs predictably and safely through spikes and failures and the ordinary daily chaos of a very large fleet, and asks our customers to think about none of it.

The line we wouldn't cross

Every disaggregated database has to answer one question about its shared storage layer: can that shared layer read your data? For most of them, the answer is yes: a natural consequence of what conventional disaggregation demands of storage: to materialize pages from a redo log.

The database ships a log of changes, and the storage fleet reconstructs pages by interpreting that log: find the page, apply the change, produce the bytes. This composition has a consequence that no amount of engineering can wish away. To interpret data, you must be able to read it. A storage tier that reconstructs your pages is, by definition, a storage tier that can see your pages.

The defenses built around that fact tend to be operational. Mandatory access controls, sacrificial nodes, scoped credentials, periodic audits. In other words, part of the security promise is “trust us.”

In contrast, we wanted to say: trust our algorithm, not us. In Atlas Infinite, the shared storage layer is not semantically aware—on purpose. Your data is encrypted on the compute node, at the database level, before a single byte is handed down. It is decrypted in exactly one place: the single-tenant compute node running inside your own isolated network, using keys that can be yours to hold, rotate, and revoke. The storage layer services cannot read your documents. They cannot see your collection names, your indexes, or where one tenant's data ends and the next begins. An operator with root on the entire storage fleet would hold nothing but opaque bytes.

That is a different kind of boundary. We augmented the full network isolation of Atlas Core, with cryptographic segmentation. This is the reason we could take a system that had always been isolated by machine and network boundaries and put it on shared storage without compromising its trustworthiness.

Making an entire storage fleet operate, heal, compact, and garbage-collect data it cannot understand is genuinely hard, and it is where a lot of the real invention in this project lives. How we did it, and the tenant-isolation guardrails around it, is the subject of the next post in this series.

Epilogue

I studied databases at the University of Wisconsin–Madison, back when the field was still arguing about B-trees and data layouts. Then I wandered off for twenty-five years: high-performance computing, networks, security, and a long stretch building load balancers, whose entire job is to stand in front of machines that fail and make them, together, into something that doesn't.

I look at Atlas Infinite through that lens. We stood a service in front of the database's storage and made it something that is designed not to fail, not to run out, and not to require trust. And the old question of how data should sit on your disk got a new answer: it shouldn't sit on your disk at all.

Storage had gravity. Less so now.

Atlas Infinite is in public preview. This post is the first of three. Next, my colleagues will show how a storage fleet operates, heals, and compacts data it cannot understand. After that, we open up mongod itself: what it took to make a database engine born in the era of attached disks feel at home in a world without them.

If you've read this far, you're who we built it for. Come kick the tires. Tell us how you experience it.

megaphone
Next Steps

Get started with these new capabilities today for free at mongodb.com/atlas. For Atlas Infinite, click here. 

To learn more about Atlas Infinite and MongoDB 9.0, check out the Atlas Infinite documentation, version release page, or the MongoDB 9.0 documentation.

MongoDB Resources
Documentation|MongoDB Community|MongoDB Skill Badges