Amazon DocumentDB is a managed document database that speaks the MongoDB wire protocol. Your application uses an ordinary MongoDB driver, sends BSON documents, and runs find, aggregate and updateOne as usual. Underneath, it is not MongoDB. AWS built its own engine on a storage design that separates compute from a shared, replicated volume, much like Aurora. That one decision explains most of DocumentDB's strengths and most of its surprises.

This article explains the architecture from first principles, the three deployment options, what compatibility means in practice, how to connect so that failover and quotas do not hurt you, how to index and model data within its limits, how to consume change streams, and how to operate a cluster. A worked example follows an order service end to end, and the article ends with a checklist. Numbers were checked against the DocumentDB developer guide in October 2026; quotas change, so confirm them for your engine version before you rely on them.

Advertisement

The architecture: compute instances on a shared volume

An instance-based DocumentDB cluster has one primary instance that accepts writes and up to 15 replica instances that serve reads. All of them attach to a single cluster storage volume that keeps six copies of the data across three Availability Zones. The primary sends its log records to storage; replicas read pages from the same volume and apply log records to their caches. No replica holds a separate full copy, which has three consequences.

First, adding a replica is fast and costs compute only, because no data is copied. Second, replica lag is usually small, but it is not zero, so a read from a replica can miss a write that the primary has just acknowledged. Third, durability and availability are separate: losing an instance loses no data, while losing the primary interrupts writes until a replica is promoted.

The cluster exposes endpoints. The cluster endpoint always points at the current primary. The reader endpoint distributes connections (not individual queries) across replicas. Instance endpoints address a single node. Drivers connect in replica set mode with the set name rs0 and discover members themselves. Storage grows automatically up to the cluster quota: 128 TiB on earlier engines and 256 TiB from engine 8.0.

Instance-based DocumentDB: compute instances share one storage volume; the driver sees a replica set named rs0ApplicationMongoDB driver + poolCluster endpointalways the primaryReader endpointspreads connectionswritesreadsPrimaryAZ a: reads + writesReplica 1AZ b: readsReplica 2AZ c: reads, failover targetShared distributed storage volumesix copies across three AZs; grows automatically; 128 TiB, or 256 TiB on engine 8.0+redo logpagespagesChange stream log3 h default, up to 7 daysReplicas do not replay writes into their own copy of the data; they read the same volume, so adding one does not copy dataFailover promotes a replica and repoints the cluster endpoint; the driver must reconnect, and writes in flight can fail
Instance-based cluster: endpoints, instances in three AZs and the shared storage volume.

Compare this with Aurora's storage layer, which uses the same idea for relational engines. If you know Aurora, you already know DocumentDB's failure behaviour.

Three ways to deploy

OptionScaling modelUse it when
Instance-basedYou pick instance classes; one primary, up to 15 replicasSteady workloads; you want predictable cost and every feature
Serverless (GA July 2025)Capacity in DocumentDB Capacity Units (about 2 GiB of memory each, with CPU and network), 0.5 to 256 DCUs, engine 5.0 and laterSpiky or unpredictable load; many small databases; dev and test
Elastic clustersHash-sharded across up to 32 shards; up to 4 PiB when data is evenly distributedWrite throughput or size beyond one primary; you can choose a good shard key

Storage also has two pricing configurations. Standard charges per I/O request; I/O-Optimized charges more for compute and storage but nothing per I/O. If I/O is a large share of your bill, compare the two using a month of real usage.

Elastic clusters are a different product with separate quotas: for example, the maximum document nesting depth is 100 levels there versus 200 on instance-based clusters, and limits on users and sharded collections are lower. Choose the shard key with the same care as in any sharded database; a monotonically increasing key concentrates writes on one shard.

Advertisement

What MongoDB compatibility means

Compatibility is defined per engine version. DocumentDB 3.6, 4.0, 5.0 and 8.0 accept drivers for the matching MongoDB API versions, and each release adds operators and stages. Version 8.0, announced in November 2025, accepts drivers for MongoDB API 6.0, 7.0 and 8.0 and added, among others, the $merge, $set, $unset, $replaceWith, $bucket and $vectorSearch stages, collation, views, a new query planner (version 3), Zstandard dictionary compression and a second-generation text index. AWS states up to 7 times lower query latency and up to 5 times better compression for 8.0; treat those as vendor figures and measure your own workload.

Two consequences matter more than any feature list. First, a feature list from an older engine is wrong for a newer one and vice versa, so check the functional differences page for the engine you actually run. Second, the semantics of features both systems share can still differ, for example in query planning, index use and error codes. The only reliable test is to run your application's own integration suite against a DocumentDB cluster before you migrate, including the aggregation pipelines and the error paths.

One difference affects every application: retryable writes. Engine 8.0.2 and later support them. On earlier engines they are not supported, and because modern drivers enable them by default, the connection string must say retryWrites=false. Your code must then handle a failed write during failover itself, idempotently.

Connecting correctly

TLS is enabled by default. Download the AWS certificate bundle global-bundle.pem and pass it to the driver. Connect in replica set mode so the driver follows the primary through failovers. Here is a Python client configured for an engine below 8.0.2:

import os
from pymongo import MongoClient
from pymongo.errors import AutoReconnect, NotPrimaryError

uri = (
    f"mongodb://{os.environ['DOCDB_USER']}:{os.environ['DOCDB_PASS']}"
    "@orders.cluster-xxxx.us-east-1.docdb.amazonaws.com:27017/"
    "?tls=true&tlsCAFile=global-bundle.pem&replicaSet=rs0"
    "&readPreference=secondaryPreferred&retryWrites=false"
)
client = MongoClient(
    uri,
    maxPoolSize=50,              # per process; multiply by replicas of your service
    serverSelectionTimeoutMS=15000,
    socketTimeoutMS=20000,
)
orders = client.shop.orders

def place(order):
    # Idempotent: _id is the client-generated order id, so a retry cannot duplicate it.
    for attempt in range(3):
        try:
            fields = {k: v for k, v in order.items() if k != "_id"}
            return orders.update_one({"_id": order["_id"]},
                                     {"$setOnInsert": fields}, upsert=True)
        except (AutoReconnect, NotPrimaryError):
            if attempt == 2:
                raise

Keep credentials out of the URI in real deployments, for example in Secrets Manager with rotation as described in AWS secrets rotation. Run the cluster in private subnets and allow port 27017 only from the application's security group.

Connections are a hard quota per instance, and the quota depends on the instance size. A db.r6g.large allows 3,400 connections in total, 1,100 of them active, and 450 open cursors; a db.t3.medium allows 1,000 connections and only 30 cursors. The arithmetic that breaks teams is pool size times processes: 40 pods with maxPoolSize=100 can open 4,000 connections to the primary. Size pools from the quota downward, as in connection pooling, and alarm on DatabaseConnectionsMax and DatabaseCursorsMax against their limit metrics.

Reads with secondaryPreferred are cheaper but may be stale. Send any read that must see the caller's own write to the primary, and use replicas for reporting, search pages and anything that tolerates a short delay. Engine 8.0.1 and later also support snappy wire compression with compressors=snappy, which helps when large documents saturate the network.

Modelling and indexing within the limits

The document model rules are the usual ones: embed data that is read together and bounded in size, reference data that grows without limit or is shared. DocumentDB's limits make the boundaries concrete. A document can be at most 16 MiB. A collection can have 64 indexes, a compound index 32 keys, and an index key 2,048 bytes. A single batch write can contain 100,000 operations.

Every query that runs often should be served by an index, and you should prove it with explain() rather than assume. A plan that shows COLLSCAN on a large collection reads every document, and on Standard storage you pay per I/O for it. Compound indexes follow the usual rule: equality fields first, then the sort field, then range fields. Each additional index adds write cost and storage, so drop the ones the profiler and index statistics show are unused.

TTL indexes work, but deletion is best effort. DocumentDB does not guarantee that expired documents disappear within a fixed time, and heavy load delays them further. If expired data must never be returned, filter on the expiry field in the query as well.

Change streams

Change streams let a consumer follow inserts, updates and deletes in order, which is how you feed search indexes, caches and event pipelines without dual writes. In DocumentDB they must be enabled per collection or database, and the change log is kept for 3 hours by default, extendable to 7 days with the change_stream_log_retention_duration cluster parameter. The retention window is your recovery budget: a consumer that is down longer than the window cannot resume and must rebuild from a full scan.

# Enable once, as an admin:
#   db.adminCommand({modifyChangeStreams: 1, database: "shop",
#                    collection: "orders", enable: true})

token = load_checkpoint()                      # persisted resume token, or None
with orders.watch(full_document="updateLookup", resume_after=token) as stream:
    for change in stream:
        publish_to_search(change)              # must be idempotent
        save_checkpoint(stream.resume_token)   # after the side effect, not before

Save the resume token after the side effect succeeds, never before, so a crash causes a replay rather than a gap; that makes the consumer at-least-once, so its writes must be idempotent. Monitor the age of the last processed event and alarm well inside the retention window.

Operating a cluster

  • Backups. Continuous backup supports point-in-time restore anywhere in a retention period of 1 to 35 days (the default is one day). A restore creates a new cluster; it does not roll back the existing one, so rehearse the cut-over, including DNS and secrets.
  • Failover. Put instances in at least two AZs, and set promotion tiers so the replica you want is promoted first. During failover the cluster endpoint moves and in-flight writes can fail; test it with a forced failover in staging and watch your error rate.
  • Monitoring. Alarm on CPU, FreeableMemory, BufferCacheHitRatio, DBInstanceReplicaLag, connection and cursor usage against limits, and the transaction limit metrics. A falling buffer cache hit ratio means the working set no longer fits in memory.
  • Slow queries. Enable the profiler with a threshold such as 100 ms and send logs to CloudWatch; most performance problems are missing indexes or unbounded queries.
  • Upgrades. Major engine upgrades change compatibility and quotas. Restore a snapshot into a test cluster on the new engine and run your integration suite before upgrading production.

Worked example: an order service

An order service writes 300 orders per second at peak and serves order history to customers. It runs on an instance-based cluster with an r6g.xlarge primary and two replicas in three AZs. Orders embed their line items, which are bounded, and reference the customer by id. Two indexes serve all hot queries: {customerId: 1, createdAt: -1} for history pages and {status: 1, updatedAt: 1} for the fulfilment worker.

The service has 30 pods with a pool of 40 connections each: 1,200 connections, well under the 7,000 allowed on the primary. History pages read from replicas, but the confirmation page reads the new order from the primary. A change stream consumer feeds the search index and the email service from a resume token stored in DynamoDB, with retention raised to 24 hours. A forced failover in staging should show a short burst of write errors while the endpoint moves; the idempotent upsert retries absorb them without creating duplicates, and the test proves it.

Trade-offs and failure modes

SituationWhat goes wrongMitigation
Driver defaults on an engine below 8.0.2Writes rejected because retryable writes are onretryWrites=false and idempotent app-level retries
Large pools on many podsConnection quota exhausted; new connections refusedPool budget per instance; alarm on usage against limit
Read-after-write on a replicaUser does not see their own updatePrimary reads for those paths
Consumer offline longer than retentionChange stream cannot resumeLonger retention, lag alarms, rebuild runbook
Feature works on MongoDB, not hereUnsupported operator or different semanticsRun the full test suite against the target engine
Single write-heavy primaryWrite throughput ceilingLarger instance, batching, or elastic clusters

Choose DocumentDB when you want the MongoDB API, AWS-native operations (IAM, VPC, KMS, CloudWatch) and Aurora-style storage, and your application fits the features of a specific engine version. Choose MongoDB itself when you depend on features DocumentDB lacks or need multi-cloud portability. Choose DynamoDB when access patterns are known and simple and you want serverless scale without managing connections at all.

What to do next

  1. Write down your engine version and check the functional differences and quotas pages for that version, not for MongoDB.
  2. Fix the connection string: TLS with global-bundle.pem, replicaSet=rs0, an explicit read preference, and retryWrites=false below 8.0.2.
  3. Compute your connection budget (pods times pool size) against the instance quota and add alarms on connection and cursor usage.
  4. Run explain() on your ten most frequent queries and add or drop indexes until none of them shows a collection scan.
  5. Set change stream retention to cover your longest plausible consumer outage, and alarm on consumer lag.
  6. Rehearse a forced failover and a point-in-time restore to a new cluster, and time both.
Key takeaway: Amazon DocumentDB gives you the MongoDB API on an AWS-built engine whose compute instances share one replicated storage volume, which makes replicas cheap and storage durable but leaves a single write primary, small replica lag and real per-instance connection quotas. Treat compatibility as a property of one engine version and prove it with your own tests, configure TLS, replica set mode and retryable writes correctly, budget connections, index every hot query, give change streams enough retention, and rehearse failover and restore before you need them.