skip to content

How would you decide whether a sharded workload should use cross-shard MongoDB transactions or be remodelled to avoid them?

level: principalimportance: should knowfreq 33%

answer

  1. Ask how many shards the transaction touches
  2. One shard behaves like a replica set
  3. Multiple shards means a coordinator gets involved
  4. Prepared documents block other traffic
  5. The shard key is the real lever

basics

~20 s

Measure how many transactions actually span shards and what they cost. A transaction touching one shard is cheap; one touching several commits through two-phase commit with prepared state and held locks. Prefer a shard key that colocates the transactional unit.

solid answer

~50 s

Start from the data, not the doctrine. A transaction whose documents all live on one shard is essentially a replica-set transaction and is affordable. Once it writes to more than one shard, MongoDB runs a distributed transaction: one participating shard acts as coordinator, participants prepare, and only then does the commit land — extra round trips, durable coordinator state, and prepared documents that block conflicting readers and writers until the decision arrives. Cost grows with the number of participants and with commit latency. So I would measure the share of transactions that are multi-shard, their p99 commit latency against the single-shard baseline, and the abort rate. The strongest lever is the **shard key**: choose one that colocates a transactional unit — tenant, account, order — so the common write is single-shard by construction. Next, model the invariant into one document. Reserve genuine cross-shard transactions for low-rate, correctness-critical writes, and use an outbox or compensating workflow where the boundary is really a service boundary.

code

javascript · 9 lines
javascript
// Scattered: shard key spreads a customer's data across shards
sh.shardCollection("app.orders",  { _id: "hashed" });
sh.shardCollection("app.ledger", { _id: "hashed" });
// every order+ledger transaction is likely multi-shard

// Colocated: the transactional unit shares a shard-key prefix
sh.shardCollection("app.orders",  { customerId: 1, _id: 1 });
sh.shardCollection("app.ledger", { customerId: 1, ts: 1 });
// a transaction scoped to one customerId stays on one shard

go deeper

for a junior

Know that MongoDB supports transactions on sharded clusters, and that a transaction staying within one shard is much cheaper than one spanning several.

for a middle

Explain what changes across shards: a coordinator, a prepare step on each participant, and a commit decision, which adds round trips and leaves documents blocked until the outcome is known.

for a senior

Diagnose and reduce the cost — measure the multi-shard share and its p99 commit latency, find which access patterns scatter co-changing documents, and get the hot path onto a single shard.

for a principal

Own the call between reshaping the shard key, remodelling the invariant into one document, and accepting the distributed commit — including migration risk, the evidence that would justify resharding, and what you would do instead across service boundaries.

## What changes when a transaction crosses shards On a replica set, a transaction is local: one primary stages the writes, one commit makes them visible. In a sharded cluster, `mongos` routes each operation to the shards that own the documents. If every operation in the transaction lands on the same shard, the transaction is effectively a single-shard transaction and behaves like the replica-set case. When the transaction writes to more than one shard, it becomes a **distributed transaction**. MongoDB elects one of the participating shards as the transaction coordinator, and the commit runs as a two-phase commit: the coordinator asks each participant to prepare, each participant durably records that it is prepared and holds its locks, and only when all have voted does the coordinator record the decision and tell everyone to commit. The coordinator's state is persisted so the decision survives a failure — a prepared transaction is not something a crash can quietly forget. ## Where the cost actually lands - **Latency.** The commit is no longer one acknowledgement; it is a prepare round plus a decision round across every participant, each majority-acknowledged. Commit latency is set by the slowest participant, so p99 gets worse as the number of participants grows. - **Blocking on prepared documents.** Between prepare and decision, the documents a participant has prepared are in limbo. Other operations that need to see them may have to wait for the decision. A brief coordinator hiccup therefore radiates into unrelated traffic touching those documents. - **Failure surface.** More participants means more ways to hit a transient error and re-run the whole transaction, and more chance of hitting the transaction lifetime limit. Cluster activity such as a chunk migration involving the data can also abort the transaction. - **Operational opacity.** Diagnosing a slow distributed commit means correlating across several shards plus the coordinator, which is a materially harder investigation than one slow primary. None of this makes cross-shard transactions unusable. They are a supported, correct feature, and for a low-rate, high-value operation — settling a payment across two accounts on different shards — paying that cost is the right call. The judgment question is whether your *hot path* is paying it thousands of times a second. ## The decision framework **1. Measure the split.** How many transactions per second are single-shard versus multi-shard, and how many shards do the multi-shard ones touch? A workload where 2% of transactions are cross-shard has a very different answer from one where 60% are. Compare p99 commit latency for each class, and track the abort/retry rate. **2. Fix the shard key first.** This is the highest-leverage change and the hardest to reverse. If the transactional unit is "everything belonging to one tenant" or "one account and its ledger entries", a shard key prefixed by that identifier colocates the whole unit on one shard, and the common transaction becomes single-shard by construction. The classic failure is a shard key chosen purely for write distribution that scatters co-changing documents across every shard, turning every business operation into a distributed commit. **3. Then look at the model.** Many invariants that appear to need a transaction fit in one document — order plus line items, entity plus its counters — where atomicity is free and shard placement is irrelevant. Others can be restructured so that the cross-boundary part is a derived value that tolerates a short lag rather than a hard invariant. **4. Then consider a workflow instead of a transaction.** Where the boundary is genuinely a service or domain boundary, a transaction is often the wrong shape: write the local change and an outbox record in one atomic transaction on one shard, and drive the remote effect from the outbox with compensation on failure. You trade linearizable simplicity for availability and independence — a trade worth making when the two halves have different owners or availability requirements. **5. Keep the residual transactions small and boring.** Whatever cross-shard transactions survive should touch few documents, hold no external calls, run well under the lifetime limit, commit with an explicit write concern, and be instrumented separately so their cost stays visible. ## What would change my mind If the cross-shard rate is low and stable and the p99 is within budget, remodelling is not worth the migration risk — a resharding effort is expensive and disruptive, and correctness you already have is worth a lot. If the rate is growing with traffic, or commit latency dominates the request budget, or the abort rate is climbing under contention, then the shard key is the problem and no amount of tuning fixes it. ## The failure mode to name explicitly The worst version is a team that reaches for transactions to paper over a model that scattered co-changing data, then raises the transaction lifetime limit when the transactions start timing out. That converts a design problem into an operational one: longer-held prepared state, more blocking, more lag — and a much harder diagnosis when it finally breaks.

  • What does MongoDB do differently when a transaction writes to three shards instead of one?
    It runs a distributed commit. One participating shard acts as transaction coordinator, each participant is asked to prepare and durably records that it is prepared while holding its locks, and only after all have voted does the coordinator record and distribute the decision. That means additional majority-acknowledged round trips, coordinator state that must survive failure, and prepared documents that can block other operations until the decision lands.
  • Why is the shard key the highest-leverage lever here rather than tuning transaction settings?
    Because the shard key determines whether a transaction is distributed at all. If co-changing documents share a shard-key prefix — tenant, account, order — the routine transaction touches one shard and pays replica-set costs. No timeout, concern or retry setting can turn a multi-shard commit into a single-shard one. The catch is that changing a shard key is a heavyweight migration, which is why it deserves the analysis up front.
  • When would you accept cross-shard transactions rather than remodel?
    When the rate is low and stable, the p99 commit latency fits the request budget, and the invariant is genuinely valuable — money movement between accounts that cannot be colocated, for example. Resharding is disruptive and risky; correctness you already have is worth preserving. I would instrument those transactions separately, keep them small, and revisit the decision if their share of traffic starts growing.

saying these in an interview costs you the question

  • Treats cross-shard transactions as free once they are supported
  • Ignores shard-key choice when transactions are slow
  • Raises the transaction time limit instead of fixing placement
  • Cannot say what the commit does differently across shards
  • Proposes distributed transactions across service boundaries by default

context