MongoDB arguments usually generate more heat than light. The useful question is narrower: does your workload read and write data in self-contained units that match one document, and do your consistency and reporting needs fit what a document database does cheaply?
This article builds that test from the architecture up. It explains what MongoDB guarantees and why, walks through a modelling exercise for an order system, shows the settings that decide durability, and lists the cases where you should choose something else. For the relational counterpart, read when to pick Postgres alongside it.
What MongoDB is, architecturally
MongoDB stores documents: nested key-value structures encoded as BSON, a binary JSON with extra types such as dates, 64-bit integers, decimals and ObjectIds. Documents live in collections, and a document may be at most 16 MB. The storage engine, WiredTiger, provides document-level concurrency control, compression and snapshots.
The single most important fact follows from that: a write to one document is atomic, however many fields and nested arrays it touches. That is why modelling matters so much. If everything an operation must change sits inside one document, you get atomicity, a single index lookup and no join, without paying for a transaction.
For availability, a replica set holds copies of the data on several nodes. One primary accepts writes and records them in the oplog, a capped collection that secondaries tail and replay. If the primary fails, the remaining members elect a new one, usually within seconds. For scale, a sharded cluster splits collections across several replica sets by a shard key. Stateless mongos routers consult the config servers, which hold the map of key ranges, and send each query to the shards that own the relevant keys.
The fit test: are your operations aggregate-shaped?
Domain-driven design has a useful word for this: an aggregate is a cluster of data that changes together and is loaded together, with one root that outsiders refer to. An order with its line items, shipping address and status history is an aggregate. A user profile with preferences and devices is another. A product with attributes that vary by category is a third.
Write down your ten most frequent operations and ask, for each one, how many aggregates it touches. If nearly all of them read or write a single aggregate by its key or a few indexed fields, the document model will feel natural and fast: one document per aggregate, one round trip per operation. If many of them combine several entity types with conditions on each, such as which customers in which region bought which products from which supplier, you are describing a relational query workload, and MongoDB will make you either denormalise heavily or run aggregation pipelines with $lookup that a relational planner would handle better.
Two further questions sharpen the test. First, are invariants local? A rule such as an order's total equals the sum of its lines is local to one document and free to enforce. A rule such as stock never goes negative across all orders spans documents and needs a transaction or a different design. Second, does the shape genuinely vary? Heterogeneous attributes are a real reason to use documents; a team simply not wanting to write migrations is not, because the application still assumes a shape and old documents still have to be handled.
Where it fits well
- Catalogues and content with varying attributes. A shoe and a laptop share a name and a price and little else. In MongoDB each product carries its own attribute set, and the attribute pattern (an array of name-value pairs with one compound index) keeps them all searchable.
- Per-user or per-tenant state read by key. Profiles, carts, settings, session-like state and game saves are loaded whole by one ID and written back. One document is one read.
- High write volume that must scale out. When the dataset or write rate outgrows one machine, sharding is built in, and a key that matches your access pattern spreads load evenly. Online resharding, available since 5.0, makes a bad first key recoverable rather than fatal.
- Events and measurements. Time series collections, added in 5.0, bucket measurements that share metadata internally, which cuts storage and speeds range queries over time.
- Evolving APIs that store what they serve. When the stored document is close to the JSON the API returns, the mapping layer almost disappears, which genuinely speeds up iteration.
Where it strains
- Ad hoc, cross-entity reporting. Analysts who join five entity types in new ways each week will be happier with SQL. Common practice is to stream changes into a warehouse and keep MongoDB for the operational path.
- Invariants spanning many documents. Multi-document ACID transactions exist, across replica sets since 4.0 and sharded clusters since 4.2, but they hold locks and snapshot history while they run, they are aborted by default after 60 seconds, and cross-shard ones need a two-phase commit. Use them for the occasional cross-document rule, not as the main write path.
- Unbounded growth inside one document. Comments on a viral post, or events on a long-lived device, will eventually approach 16 MB, and every update rewrites a growing document. Bound arrays, or move them to their own collection.
- Heavily many-to-many data. Students and courses, or users and groups with permissions, either duplicate data on both sides, which you must keep consistent, or turn into joins.
- Teams that need strict schemas everywhere. You can enforce JSON Schema validation per collection, as shown below, but the guarantees are weaker and less familiar than typed columns, foreign keys and check constraints.
Worked example: modelling an order system
Consider an online shop. Its main operations are: place an order; show a customer their orders, newest first; show one order's details; move an order through paid, shipped and cancelled; feed a fulfilment queue of orders waiting to ship; and report revenue by product monthly. Customers, products and inventory also exist.
Apply the aggregate test. Showing, placing and updating an order all touch one order with its lines, so lines are embedded as an array. The customer is a separate aggregate referenced by customerId, but the order page needs the name and email, so the order also holds an extended reference: a copy of just those fields, taken at order time. That copy is a feature, not a bug, because an invoice should show the address it was shipped to, not today's address. Status history is a bounded array inside the order. Inventory is its own collection, because stock is shared by all orders.
db.createCollection("orders", {
validator: { $jsonSchema: {
bsonType: "object",
required: ["customerId", "status", "placedAt", "lines", "total"],
properties: {
customerId: { bsonType: "objectId" },
status: { enum: ["placed", "paid", "shipped", "cancelled"] },
placedAt: { bsonType: "date" },
customer: { bsonType: "object", // extended reference: a copy of
required: ["name", "email"] }, // the fields the order page shows
lines: { bsonType: "array", minItems: 1, maxItems: 500,
items: { bsonType: "object",
required: ["sku", "qty", "unitPrice"],
properties: { qty: { bsonType: "number", minimum: 1 },
unitPrice: { bsonType: "decimal" } } } },
total: { bsonType: "decimal" }
} } },
validationAction: "error"
});
// The access paths, written down before the indexes
db.orders.createIndex({ customerId: 1, placedAt: -1 }); // "my orders", newest first
db.orders.createIndex({ status: 1, placedAt: 1 }, // fulfilment queue
{ partialFilterExpression: { status: { $in: ["placed", "paid"] } } });Each index corresponds to one listed access path. The compound index on customer and date serves the customer's order history in sorted order without a separate sort step. The partial index covers only orders still waiting to ship, so it stays small however many historical orders accumulate. Status changes and monthly reports then look like this:
// One atomic write: status change + audit entry in the same document. No transaction needed.
db.orders.updateOne(
{ _id: orderId, status: "paid" }, // guard = optimistic check
{ $set: { status: "shipped", shippedAt: new Date() },
$push: { events: { $each: [{ at: new Date(), type: "shipped" }], $slice: -50 } } },
{ writeConcern: { w: "majority" } }
);
// Reporting: revenue per SKU last 30 days. Fine on a secondary or an analytics node.
db.orders.aggregate([
{ $match: { status: { $in: ["paid", "shipped"] },
placedAt: { $gte: ISODate("2026-09-02") } } },
{ $unwind: "$lines" },
{ $group: { _id: "$lines.sku",
units: { $sum: "$lines.qty" },
revenue: { $sum: { $multiply: ["$lines.qty", "$lines.unitPrice"] } } } },
{ $sort: { revenue: -1 } }, { $limit: 20 }
], { readPreference: "secondaryPreferred" });The status update is a single atomic operation that includes its own guard: it only matches if the order is still paid, so two workers cannot both ship it. The $slice keeps the events array bounded. The report scans one month of orders on a secondary, which is acceptable for a monthly job; if analysts want it daily and sliced many ways, that is the signal to ship changes to a warehouse.
Placing an order is the one operation that crosses aggregates, because it must reserve stock. Here a transaction is the honest tool:
# Python (PyMongo): an invariant that spans two documents needs a transaction.
from pymongo import MongoClient, WriteConcern, ReadPreference
from pymongo.read_concern import ReadConcern
client = MongoClient(URI) # retryable writes are on by default
inv, orders = client.shop.inventory, client.shop.orders
def place_order(order):
def txn(session):
for line in order["lines"]:
r = inv.update_one({"_id": line["sku"], "onHand": {"$gte": line["qty"]}},
{"$inc": {"onHand": -line["qty"]}}, session=session)
if r.modified_count != 1:
raise ValueError(f"out of stock: {line['sku']}") # aborts everything
orders.insert_one(order, session=session)
with client.start_session() as s:
# with_transaction retries on transient errors; keep the body short and idempotent
s.with_transaction(txn, read_concern=ReadConcern("snapshot"),
write_concern=WriteConcern("majority"),
read_preference=ReadPreference.PRIMARY)If transactions start appearing in most of your write paths, revisit the model or the database choice; that is the clearest sign the data is relational.
Consistency settings that decide correctness
MongoDB's durability and freshness are chosen per operation, and the defaults have changed over time, which is the root of many old horror stories. Write concern says how many members must acknowledge a write. Since 5.0 the implicit default is w: "majority" for most topologies, but replica sets with arbiters can fall back to w: 1, and a w: 1 write acknowledged by a primary that then fails before replicating is rolled back. Set majority explicitly for anything that matters, and avoid arbiters.
Read concern says what a read may see: local returns the node's latest data, which might later be rolled back; majority returns only majority-committed data; linearizable adds a guarantee of recency for single-document reads at a latency cost. Read preference chooses the node. Reading from secondaries scales reads but returns stale data, so a user who writes and immediately reads from a secondary may not see their own write unless you use a causally consistent session. Retryable writes, on by default in current drivers, make a single write safe to retry after a network blip.
Choosing a shard key
You do not need sharding until a replica set cannot keep up, and many applications never reach that point. When you do, the shard key decides everything. A good key has high cardinality, spreads writes evenly, and appears in the filter of your most frequent queries so they can be routed to one shard. For the shop, customerId is a reasonable choice: lookups by customer go to one shard, and orders spread across many customers. A monotonically increasing key such as placedAt or an ObjectId sends every new insert to the same chunk and the same shard; hashing the key spreads writes but turns range queries into scatter-gather. Sharding strategies compared covers range, hash and directory partitioning in general, and database sharding covers rebalancing and the operational side.
Failure modes
- Rolled-back writes after failover because a write used
w: 1or the topology included arbiters. - Document bloat: unbounded arrays approach 16 MB and make every update slower long before that.
- Missing indexes show up as collection scans that are fine at ten thousand documents and catastrophic at fifty million. Read
explain()output for every hot query. - Inconsistent denormalised copies when an extended reference was meant to track its source and nobody wrote the update job.
- Schema drift: documents written by five application versions with different shapes, and code that assumes only the newest one.
- Scatter-gather queries on a sharded cluster because the shard key is absent from the filter, which makes latency grow with shard count.
- Oplog window too short: a secondary that falls behind further than the oplog covers needs a full resync.
A decision table
| Your workload | Lean towards | Why |
|---|---|---|
| Operations load and save one aggregate by key | MongoDB | One document per operation, atomic without transactions |
| Heterogeneous attributes per item | MongoDB | No sparse tables or entity-attribute-value designs |
| Write volume beyond one primary, key-routed | MongoDB | Sharding and resharding are built in |
| Many-entity joins and ad hoc SQL | Postgres | Planner and joins are the core strength |
| Cross-entity invariants on most writes | Postgres | Constraints and transactions are the default path |
| Both: operational aggregates plus heavy analytics | MongoDB plus a warehouse | Keep each system on the workload it handles cheaply |
What to do next
- List your ten most frequent operations with their rates and the entities each one touches.
- Mark each as single-aggregate or cross-aggregate; if most are cross-aggregate, choose a relational database.
- Sketch documents for the aggregates, embedding what is read together and referencing what has an independent life.
- Bound every array and estimate the largest document you expect in five years.
- Write a JSON Schema validator for each collection and one index per listed access path.
- Set write concern majority explicitly, avoid arbiters, and decide where reads may be stale.
- Load-test with production-scale data, read
explain()for hot queries, and choose a shard key only when one replica set is no longer enough.