Sharding Failure Modes — Orphans, Partial Commits, Dual-Writes & Idempotency
Sharding turns single-node ACID into a distributed problem. Orphan rows, partial multi-shard commits, dual-write divergence, lost updates at cutover, and missing idempotency keys are the failures seniors name. Prefer one shard per transaction. Use a saga when you cannot. Two-phase commit is the rare exception.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
A workflow touches two shards
Prefer
Same shard if you can, otherwise a saga
Put rows that must commit together on one key. When the business process really crosses shards, record intent, compensate, and make every retry idempotent.
- Availability survives a partition better than a blocked coordinator.
- Compensations are explicit and have their own idempotency keys.
- You do not pretend the engine enforced a foreign key it cannot see.
Alternative
Distributed transactions by default
XA or 2PC gives a stronger cross-shard isolation story when the engine actually supports it. Coordinators block, and cloud OLTP teams rarely run them.
- A coordinator outage stalls participants.
- Latency is the protocol, not the query.
- It feels safe across microservices and is not.
Overview
Interviewers ask what happens when the network fails halfway through a multi-shard write, or when dual-write succeeds on the old shard and fails on the new one. The catalog below is the answer. The move procedure is resharding. The key that should have kept the transaction local is shard keys.
Schema dual-write during a column or table migration is Zero-downtime database migrations. Do not mix the runbooks. A messaging partition key is also not a database shard key. A log can carry a saga. It does not place the row.
By the end you should be able to:
- Match each failure to a symptom, a root cause, and a mitigation
- Choose saga over 2PC for ordinary multi-shard business work
- Place the idempotency record on the shard that owns the business key
- Repair divergence without deleting the old copy early
- Version the routing map so a reshard does not 404 forever
Decisions
- 1
Write spans shard A and B
- nextA commits?
- ?
A commits?
- NoRetry the whole unit
- YesB commits?
- 3
Retry the whole unit
- ?
B commits?
- YesBoth sides committed
- NoA committed, B did not
- 5
Both sides committed
- 6
A committed, B did not
- nextCompensate A or repair
- 7
Compensate A or repair
Lesson map
Sharding Failure Modes — Orphans, Partial Commits, Dual-Writes & Idempotency
Sharding turns single-node ACID into a distributed problem. Orphan rows, partial multi-shard commits, dual-write divergence, lost updates at cutover, and missing idempotency keys are the failures seniors name. Prefer one shard per transaction. Use a saga when you cannot.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB w["Write spans shard A and B"] aok["A commits?"] abort["Retry the whole unit"] bok["B commits?"] w -->|Write spans shard A and B| aok aok -->|No| abort aok -->|Yes| bok
If A never commits, retry the whole unit. If A commits and B does not, you already have an orphan or a one-sided effect. Compensation or a repair queue is the recovery. There is no silent rollback of A.
Failure catalog
| Failure | Symptom | Root | Mitigation |
|---|---|---|---|
| Orphan child row | Line items without an order | Cross-shard write without a saga | Same-shard key, or an outbox plus an orchestrator |
| Partial commit | Money moved on one side only | Multi-shard transaction aborted halfway | Saga with compensations; avoid the split |
| Dual-write divergence | Old and new differ during a reshard | Second write failed, or a bad clock rule | Idempotent upsert, repair job, shadow alerts |
| Phantom read across shards | A global report is wrong | No global snapshot | Accept staleness, or take an OLAP snapshot |
| Hot shard meltdown | One shard at 100 percent CPU | Whale or monotonic key | Split, dedicate, or salt |
| Idempotency gap | Double charge on retry | Retries without keys | Idempotency keys on the shard, plus a ledger |
| Routing sticky miss | 404 after a reshard | Client cached the shard map | Versioned routing config and a short TTL |
| Unique violation surprise | Duplicate email on two shards | UNIQUE is per shard | Global allocator or a directory |
Sagas versus two-phase commit
| 2PC | Saga | |
|---|---|---|
| Isolation | Stronger across shards | Eventual; compensations |
| Availability | Coordinator can block | Higher under partitions |
| Ops | XA is rare in modern OLTP SaaS | Common with an outbox |
| Prefer | Rare, short, critical sections where XA exists | Most business workflows that cross shards |
Pat Helland's argument is the design pressure: once entities live in different places, you stop pretending one transaction covers them. Vitess can offer a 2PC mode. It is optional for a reason. AWS's saga writeup is the compensation shape. microservices.io is the pattern page. Use them as the vocabulary. The database decision is still "did these rows need to be on one shard?"
Idempotency on sharded systems
Store the idempotency record on the shard that owns the business key, or in a dedicated idempotency store that routes with that key. A retry must hit the same shard. That is another reason the shard key should match the idempotency scope. A key stored on a random shard will not see the duplicate.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Press Run. Snippets must be self-contained — no network, files, or native modules.
The TypeScript check is a toy. It compares JSON text, so key order has to be stable. Production comparison is a canonical checksum or a field-level diff on the shadow-read path, not JSON.stringify of two maps you built differently. The signal is the same: old and new disagree, so you do not cut over.
Repair
- Reconciliation. Checksum ranges on a schedule. Re-copy mismatches.
- Outbox or inbox. Durable events for cross-shard side effects, so a crash after the local commit still has the intent.
- Compensating transactions. Explicit undo, with their own idempotency keys. A compensation that double-fires is a second incident.
- Quarantine routing. Pin a bad key to a repair path so the rest of the shard keeps serving.
Dual-write failed halfway: pause the cutover, alert from shadow metrics, repair from the queue. Do not delete the old shard early. A routing cache after reshard: version the shard map, make clients refresh, and treat a stale route as a miss that retries against the new map. A 404 that is actually "wrong shard" should not become a deleted account.
Partial aggregations during a failure: prefer a degraded read with a partial flag over a silent wrong total. That rule is the same as on the scatter-gather page, and it matters more when a shard is down.
A log can help the saga or the outbox. It does not fix a bad shard key, and a messaging partition is not this cluster's shard.
Interview Q&A
What is an orphan in a sharded database?
Answer
A child row whose parent lives on another shard or never committed. It is usually a multi-shard write that finished on one side. Foreign keys will not catch it across databases.
How do you prevent a double charge on retry?
Answer
An idempotency key stored before the side effect, on the shard that owns the business key. The same key must route to the same shard. The second call returns the first result.
Dual-write failed halfway. What now?
Answer
Repair queue, shadow metrics, pause the cutover. Do not delete the old copy. Do not flip reads while mismatch alerts are firing.
Why avoid a cross-shard foreign key?
Answer
The engine cannot enforce it across databases. The invariant is an application rule: same shard key, or a saga that can compensate.
What about a routing cache after a reshard?
Answer
Version the shard map. Clients refresh. A stale map should miss and retry, not stick forever and 404. Keep the TTL short during the move.
Is a log the fix?
Answer
A log helps a saga or an outbox leave the database transaction. A messaging partition key is not a database shard key. Do not conflate the two designs.
What should a partial aggregation do when a shard fails?
Answer
Return a degraded read with an explicit partial flag, or fail. Do not publish a quiet undercount.
Why is 2PC not the default?
Answer
Coordinators block, latency stacks, and most cloud OLTP setups do not run XA across microservices. Use it only for a short critical section where the engine supports it. Design the rest as single-shard transactions.
Pitfalls
- Starting with 2PC because the diagram looks like one transaction.
- Idempotency keys on a different shard from the balance they protect.
- Compensations that are not idempotent.
- Dropping the old shard while dual-write mismatches are non-zero.
- A client that pins the shard id for the process lifetime.
- Treating per-shard UNIQUE as global UNIQUE.
- Explaining a database orphan as a messaging consumer bug. Related, not the same page.
Shard A committed an order. Shard B never wrote the line items. A retry is in flight with the same client token. Say what the customer must not be charged twice for, where the idempotency row lives, and whether you compensate A or keep retrying B.
Go deeper
The saga pattern on microservices.io and the AWS prescriptive-guidance page are the compensation vocabulary. Vitess documents its transaction model, including optional 2PC. Pat Helland's "Life beyond Distributed Transactions" is the essay behind "entities that do not share a shard do not share a transaction."
Back to the map: Database sharding and partitioning.