Lesson 0010 · Capstone

Multi-document transactions & synthesis

The escape hatch from single-document atomicity — when you genuinely need it, how snapshot isolation makes it safe, what it costs, and why good schema design usually removes the need. Followed by a synthesis quiz spanning all ten lessons.

~15 minWarm-up + transactions + cross-arc synthesis quizCore arc complete after this lesson

Warm-up — apply lesson 0009 first (cold, from memory)

A query on a sharded collection that does not include the shard key field causes mongos to:

Without the shard key, mongos can't determine which chunk holds the data — it has no choice but to broadcast. Every shard runs the query, and mongos merges all responses. Latency = slowest shard. This is the COLLSCAN of the sharding tier — it doesn't get better as you add shards, it gets worse.

You shard an orders collection on { createdAt: 1 } with range sharding. All new orders have increasing timestamps. What happens?

Range sharding on a monotonically increasing key is the canonical anti-pattern. All new timestamps are larger than existing ones, so they fall into the last chunk — which lives on one shard. That shard takes 100% of write traffic. The fix: use hashed sharding (random distribution) or a compound key like { userId, createdAt } that distributes writes across users.

mongos holds:

mongos is stateless with respect to your data. It reads the chunk map (which shard key ranges live on which shards) from the config servers and uses it to route queries. No documents ever live in mongos. You can run multiple mongos instances for redundancy at zero data-synchronization cost.

The default: single-document atomicity is enough

Lesson 0001 established the keystone: the document is the unit of storage, atomicity, and schema design. A write to a single document — any field, any nested array, any sub-document — is all-or-nothing, always, for free. No BEGIN. No COMMIT. No coordinator. WiredTiger handles it at the document boundary.

When schema design is done right — embedding related data that must remain consistent (lesson 0006) — almost every operation that looks like a transaction in the relational world becomes a single-document write in MongoDB. — MongoDB Manual: "In most cases, multi-document transaction incurs a greater performance cost over single document writes, and the availability of multi-document transactions should not be a replacement for effective schema design."

The corollary: if you find yourself reaching for a multi-document transaction, first ask whether a schema change would remove the need. A transaction is a cost you pay for crossing a boundary your schema has drawn.

When you actually need a transaction

There are genuine cases where referencing is the right schema decision (lesson 0006) and the update must still be atomic:

Multi-document transactions became available on replica sets in MongoDB 4.0 and on sharded clusters in 4.2. — MongoDB Manual: "Multi-document transactions are available for replica sets and sharded clusters… Starting in version 4.2, MongoDB provides the ability to perform multi-document transactions on sharded clusters."

How transactions work: snapshot isolation

A MongoDB transaction opens with a session and runs under snapshot isolation: the transaction reads from a consistent snapshot taken at its start, exactly as WiredTiger takes a snapshot for a single read (lesson 0002). Concurrent writes from other operations are invisible to the transaction for its entire duration. The transaction's own writes are buffered — visible within the session, invisible to other clients — until commitTransaction().

session.startTransaction({ readConcern: { level: "snapshot" }, writeConcern: { w: "majority" } }); try { db.accounts.updateOne({ _id: from }, { $inc: { balance: -amount } }); db.accounts.updateOne({ _id: to }, { $inc: { balance: +amount } }); session.commitTransaction(); // all-or-nothing } catch (e) { session.abortTransaction(); // rolls back both writes }

At commit, MongoDB applies the buffered writes atomically using a two-phase commit protocol internally. For sharded clusters, this coordinates across shard primaries before the commit is visible. If any step fails, abort is called and all writes vanish — exactly like a Postgres transaction rollback.

Write conflicts are handled optimistically: if two transactions try to write the same document, the second one that reaches that document encounters a write conflict and must abort and retry. MongoDB drivers support retryable transactions for this. — MongoDB Manual: "If a document was modified by a non-transactional write or different transaction since the start of the current transaction, the current transaction will abort."

The cost — why transactions are the exception

ConcernSingle-document writeMulti-document transaction
Coordination None — WiredTiger handles it locally Two-phase commit across one or more shards; coordinator tracks state
Snapshot lifetime Operation duration (milliseconds) Transaction lifetime — held open until commit/abort. Max 60 s by default.
Write conflicts Not possible (document-level lock) Possible — another write to the same document aborts your transaction
Lock escalation Intent locks only above document level Collection-level intent locks held for the transaction duration
Memory Minimal Buffered writes held in memory until commit; large transactions can hit oplog constraints

The 60-second timeout is hard: if a transaction doesn't commit within transactionLifetimeLimitSeconds, the server aborts it automatically. This prevents long-lived snapshots from blocking MVCC cleanup and keeps lock pressure manageable. Keep transactions short, small, and targeted.

The Postgres contrast

In Postgres, BEGIN … COMMIT is the primitive — used casually for any multi-row update, essentially zero overhead for small transactions. The planner is designed around transactions; its statistics assume they're common.

In MongoDB, the single-document write is the primitive, and the multi-document transaction is the escape hatch. The cost difference is real: a Postgres transaction touching 5 rows in one table carries negligible overhead. A MongoDB multi-document transaction crossing two collections has coordination overhead that Postgres hides inside its engine. Neither is "better" — they reflect different design priorities. Postgres is designed for normalized schemas and free joins. MongoDB is designed for embedded schemas and free single-document atomicity. The transaction cost in each engine is the price of crossing the engine's preferred pattern.

Core arc complete. You've covered the full first-principles stack: the document model and BSON (0001), WiredTiger storage (0002), indexing and ESR (0003), the query planner and explain() (0004), the aggregation pipeline (0005), schema design (0006), replica sets and the oplog (0007), write and read concern (0008), sharding (0009), and multi-document transactions (0010). Every lesson was grounded in contrast with your Postgres and InnoDB knowledge. What's left is practice and wisdom — building real schemas, profiling real queries, and reading real oplog entries.
Synthesis quiz — spanning all ten lessons

An order document embeds 10 line items. You update the order total and a line item price in one call. Is a transaction needed for atomicity?

Single-document atomicity (0001, 0010): writing any combination of fields — including nested arrays — in one document is all-or-nothing. No transaction needed. If the line items were in a separate referenced collection, that would be a different story. Embedding them (0006) gives you atomicity for free.

A transaction starts and reads a product's stock: 50. Another session concurrently decrements it to 49 before our transaction commits. What does our transaction see?

Snapshot isolation (0002, 0010): the transaction received a consistent snapshot at its start. Concurrent writes made after that point are invisible for the transaction's lifetime. Our transaction sees 50 throughout. The conflict only materializes if we also try to write the same document — at that point one of the two transactions encounters a write conflict and must abort.

For find({ region: "US", price: { $lt: 100 } }).sort({ rating: -1 }), the best compound index is:

ESR (0003): Equality first (region = "US" — most selective, keeps the remaining index sorted), Sort second (rating — the index provides the sort order so no in-memory SORT stage), Range last (price < 100 — ranges break sorted run, so they go last). {region, rating, price} eliminates both a COLLSCAN and an in-memory SORT.

You write with w: "majority". After a primary election, can any acknowledged write be rolled back?

w: "majority" (0008, 0007): an election requires a majority of voting members. The winner is, by definition, a member that already received the majority-acked write. So any write acknowledged under w: "majority" is guaranteed to be on the new primary. Rollback under normal failure modes is impossible. j: true adds single-node disk durability (journal flush) but is orthogonal to rollback prevention.

explain("executionStats") shows IXSCAN → FETCH → SORT. The examined:returned ratio is 1:1. What is still wrong?

The 1:1 examined-to-returned ratio is excellent (0004) — the index is selective. But the SORT stage means the index doesn't provide the sort order, forcing an in-memory sort. The fix: include the sort field in the correct S position (before any range fields) per the ESR rule (0003). A FETCH stage alone is not a problem if the index is selective — the issue is the SORT.

Your $lookup-heavy aggregation pipeline is slow. MongoDB's own documentation recommends:

0005 + 0006: $lookup is MongoDB's left outer join — and the Manual explicitly warns against overuse: "If you are frequently using $lookup, consider whether you can restructure your schema to use embedded documents." Heavy $lookup is a schema smell (0006) that points toward over-normalization. The fix is to embed what's always accessed together — $lookup then disappears (0001: data accessed together stored together).
Optional — when you have a mongosh instance: open two mongosh sessions. In session A, start a transaction and update an account balance. In session B (without a transaction), try to read the same document — you'll see the pre-transaction value. Commit the transaction in session A. Re-read in session B — now you see the updated value. This directly demonstrates snapshot isolation and the transaction commit boundary.