Concept lesson · Foundations
Distributed transactions and sagas
Start here
Definition
A distributed transaction is one transaction whose operations span multiple databases or transactional resource managers. An atomic-commit protocol such as two-phase commit coordinates their commit-or-abort outcome. A saga instead coordinates a business operation through committed local transactions and explicit compensating actions.
Why it matters: A local database rollback cannot undo a payment or reservation already committed by another service.
A payment timeout is an unknown outcome, not proof of failure. Reconcile before retrying or compensating.
Read the diagram step by step
- Persist order intent, reserve inventory, and call payment with one stable attempt identity.
- If the payment response is lost, record UNKNOWN and query or retry that same identity.
- When authorization succeeds, atomically allocate H81 only if still valid, then confirm O81. A late authorization after hold expiry must be voided. If terminal failure is known, release the hold.
- A compensating action can itself fail and needs durable retry; a saga is not simultaneous rollback across services.
Worked example
Order O81 needs 2 mugs and a $24 authorization. Stock is held for 120 seconds. If the hold expires before a delayed authorization succeeds, the workflow voids the authorization rather than confirming an order without stock.
Key takeaways
- Keep related changes in one local transaction when the same database can atomically commit them.
- 2PC coordinates commit or abort; isolation still needs its own concurrency protocol.
- A saga records partial progress, unknown outcomes, and recoverable compensation.
You will learn to
- Distinguish atomic commit from isolation and from business compensation.
- Model a durable workflow with stable operation identifiers and explicit uncertain states.
- Handle delayed success after a resource hold expires without overselling or duplicating an external effect.
Practice in this chapter
8 interview questions with model answers and follow-ups.
Go to interview practiceUseful foundations: Transaction isolation · Message queues, event logs, delivery guarantees, and backpressure · Idempotency, retries, and timeouts
Workload and timing examples are interview assumptions.
Dotted concept links open the relevant explanation in a new tab.
01Distributed transaction: definition and local boundaries
A distributed transaction is one transaction whose operations span multiple databases or transactional resource managers, such as an order database and an inventory database. An atomic-commit protocol, such as two-phase commit, coordinates their commit-or-abort outcome. A saga instead coordinates a business operation through committed local transactions and explicit compensating actions. A local database transaction can atomically change its own records, but cannot automatically undo an HTTP request that already succeeded at another service. Crossing independently failing systems therefore requires a protocol for partial completion.
For example, order O81 requests two MUG9 items at $12 each. Inventory starts at five, and hold H81 reserves two units for 120 seconds. Payment action A81 authorizes $24: authorization reserves funds and is distinct from capture. The order can become confirmed only after inventory allocation and the required authorization are established; partial progress remains pending.
If inventory and orders share one database and ownership boundary, a short transaction is the simplest option. Splitting tables into services prematurely creates a harder problem. We study the split because the provider is external and inventory may have a separate owner, not because every application needs distributed transactions.
02Two-phase commit: prepare and commit or abort
Two-phase commit (2PC) makes participating databases agree to commit or abort together. A coordinator records the final decision. First it asks each database to prepare. A database voting “yes” durably saves enough state to finish later and keeps the necessary locks or other protections. If all vote yes, the coordinator durably records “commit” and tells them to commit; otherwise the protocol chooses abort. Each participant must support preparing and honoring that decision.
An abort vote leads to abort. A prepared participant cannot safely invent the global decision when the coordinator is unreachable; classic 2PC can block.
Remember: Prepare votes; a durable decision; then deliver it.
Read the diagram
- Coordinator to Participant A: PREPARE
- Coordinator to Participant B: PREPARE
- Participant A to Coordinator: Durably prepared; YES
- Participant B to Coordinator: Durably prepared; YES
- Coordinator to Coordinator: Persist global COMMIT decision
- Coordinator to Participant A: COMMIT (retry delivery if needed)
- Coordinator to Participant B: COMMIT (same decision)
For O81, suppose the order and inventory databases both support 2PC. They prepare their changes, then follow the same commit-or-abort decision. This prevents one from committing while the other aborts. Their concurrency controls must still provide the required isolation; atomic commit alone does not make all cross-database transactions serializable.
03Saga: local transactions and compensation
Green arrows move the workflow forward. Rust arrows perform compensating business actions after a definite failure.
Remember: A compensation is another action, not erasure of a past commit.
Read the diagram
- Follow successful reservation and payment steps, then reverse the business effects after shipment fails.
- Reserve stock, authorize payment, then encounter a definitive shipment failure.
- Void the authorization and release stock when their state and business rules permit.
Try from memoryDoes voiding payment mean the authorization never occurred?
No. The authorization occurred and committed. Voiding it is a new action with its own outcome and recovery rules.
A durable workflow stores the business operation’s progress so another worker can continue after a crash. Model that progress as a state machine: named states and allowed transitions, such as awaiting authorization, ready to allocate, or cancellation pending. Each transition records what happened and which action is now permitted.
Our provider does not participate in the database’s prepare/commit protocol, so I choose a durable workflow. The order coordinator records O81’s state, the inventory hold identifier H81, and authorization operation A81. Each transition checks the expected previous state and records the next outgoing intent in the same local transaction.
| Approach | What it offers | Cost or limitation |
|---|---|---|
| One database transaction | One atomic local change | All protected data must fit that ownership boundary |
| 2PC | One commit/abort decision across capable participants | Prepared resources and recovery dependency |
| Saga/workflow | Recoverable progress across independent APIs | Intermediate states and explicit compensation |
The workflow needs durably stored state and a service responsible for advancing it. It does not need one process to remain alive throughout: a replacement worker can resume from the stored state.
- 1 → 2reserveO81 requested → Hold H81: 2 mugs
- 2 → 3send once logicallyHold H81: 2 mugs → Authorize A81: unknown
- 3 → 4deadline passesAuthorize A81: unknown → H81 expires at 120 s
- 3 → 5reconcile by A81Authorize A81: unknown → Late A81 success at 125 s
- 4 → 6cannot allocateH81 expires at 120 s → Cancel order; void A81 pending
- 5 → 6compensating actionLate A81 success at 125 s → Cancel order; void A81 pending
- 6 → 7retry or reconcile same voidCancel order; void A81 pending → Void confirmed; cleanup complete
04Successful saga: reserve, authorize, allocate, confirm
At time 0, inventory conditionally creates H81 for two mugs: available stock becomes three, and H81 expires at time 120. At time 1, the workflow asks the provider to authorize $24 using stable operation A81. At time 2, it records authorization success. It next asks inventory to convert H81 into an allocation for O81, only if the hold still exists and is valid. Inventory performs that check and transition atomically.
If allocation succeeds, a later coordinator transaction records the order as confirmed and publishes its event through an outbox. If the coordinator crashes after allocation but before recording confirmation, retrying the allocation request with O81 returns the existing allocation. It must not remove another two mugs. The same rule applies to authorization A81.
A distributed workflow is therefore a state machine: a set of allowed states and transitions. “Already allocated to O81” is a meaningful result. A vague boolean success loses the identity needed for recovery. The confirmation contract should also specify authorization validity and any later capture/shipping steps; those are separate transitions with their own failure handling.
Success and cancellation can race, so each state change must atomically check that the order or hold is still in a state that allows it. An allocation request names the order and hold; the inventory owner atomically returns the existing allocation, converts a still-valid hold, or rejects expiry/cancellation. The coordinator accepts confirmation only from its expected pending state. If cancellation won locally but allocation had already committed remotely, recovery records that allocation and releases it through an idempotent compensating transition; simply ignoring the late reply would strand stock. Start shipping only after checking that the order is confirmed and remains eligible for fulfillment.
05Unknown outcomes: lost replies and expired reservations
Now let the authorization response disappear. The provider may have processed A81 even though the coordinator received nothing. The workflow records authorization_unknown, queries by A81 or retries under the provider’s idempotency contract, and avoids inventing a new authorization identifier.
| Time | Durable or external fact | Correct reaction |
|---|---|---|
| 0 s | H81 reserves two mugs until 120 s | O81 remains pending |
| 1 s | A81 sent; response lost | Record uncertainty and reconcile |
| 120 s | H81 expires before allocation | Two mugs become available again |
| 125 s | Reconciliation finds A81 succeeded | Do not confirm from this fact alone |
| After 125 s | Allocation is no longer possible through H81 | Void A81 and finish cancellation |
06Saga recovery: compensation, retries, and outbox
Suppose voiding A81 times out too. Marking O81 simply “cancelled” and forgetting it would hide unfinished work. Record cancellation requested, authorization cleanup pending, and a stable void operation identifier. Retry or query that operation, and retain enough evidence for an operator to resolve a permanently unclear provider outcome. The client can see that the order will not ship while the authorization release is still processing.
For every step, specify how recovery works if a process crashes: before the local commit there is no recorded intent; after commit a worker can rediscover it; after remote success but before recording the result, a stable key or status query resolves ambiguity. The outbox closes the local database/publication gap, but it does not make the provider part of the local transaction.
Compensation is not always a valid business remedy. Shipping an irreplaceable item twice cannot be made correct merely by scheduling a refund. Protect scarce inventory with conditional allocation and gate irreversible steps carefully. If the business rule forbids exposing partial completion and no compensating action can repair it, reconsider which service owns the data or use participants that can commit the required changes together.
For implementation, a small workflow can use a transactional state table, an outbox and leased workers. A durable workflow engine such as Temporal provides persisted event history and replay, but workflow code must follow its deterministic execution constraints. External calls belong in retryable activities with stable effect identities; the engine does not give a third-party API transactional rollback or unlimited deduplication.
07Orchestration versus choreography and interview explanation
In an interview I would say: “O81 has a durable coordinator record, and each remote step uses a stable operation identifier. Inventory owns the hold’s expiry and conversion. The order remains pending while authorization is uncertain. After the hold expires, a late authorization triggers voiding, not confirmation. Every outgoing step is recoverable from an outbox, and every incoming result is checked against the current workflow state.”
An orchestrated workflow puts these transitions in one explicit coordinator. An event choreography distributes reactions among services; it may reduce central coupling but makes the overall progress and compensation path harder to inspect. Either approach needs ownership of timeouts, retries and terminal outcomes.
Measure the age and count of stuck pending orders, unknown external outcomes, failed compensations and expired holds. Alert on old unfinished work, not only HTTP errors. A successful request log does not prove that the multi-step business operation finished. The design is complete when another worker can recover O81 from persisted facts without guessing what the previous worker intended.
Practise the interview questions
Say your answer aloud before opening the model answer. Then answer the follow-up and compare the reasoning.
What problem does a distributed transaction or saga solve?
Reveal a model answer
A distributed transaction spans multiple transactional participants and needs a coordinated commit-or-abort outcome. A saga addresses a related business need through separately committed local transactions and compensation. For O81, creating an order, reserving 2 mugs, and authorizing $24 can succeed or fail separately. A capable 2PC system coordinates one commit decision; a saga records local progress and compensates failures. I first ask whether the work could remain in one simpler database transaction.
Interviewer follow-up
When would you avoid a saga?
Reveal the follow-up answer
When the invariant and data already fit one database ownership boundary. A saga adds visible intermediate states and recovery work. Putting separate databases in one application process does not make them one transaction domain.
What the answer must demonstrate: Identify the actual independent commit boundaries.
What does a yes vote in 2PC mean?
Reveal a model answer
“The participant has prepared enough durable state and retained the necessary protections to honor a later commit decision. It is stronger than saying the request looks valid right now.”
Interviewer follow-up
Why can coordinator failure block progress?
Reveal the follow-up answer
“A yes voter may not know whether commit was already decided, so it cannot safely invent an abort solely from a timeout.”
What the answer must demonstrate: Prepared is a durable protocol state, not a best-effort check.
Does 2PC guarantee serializable transactions?
Reveal a model answer
“2PC coordinates the final commit or abort outcome. Isolation depends on the concurrency-control protocol over the affected reads and writes. I would not claim serializability just because every participant votes on one decision.”
Interviewer follow-up
What would you inspect?
Reveal the follow-up answer
“I would inspect locking or validation across participants and show whether concurrent transactions admit a valid serial order.”
What the answer must demonstrate: Atomic commit and isolation solve different parts of correctness.
A payment authorization A81 times out with no known result. What should a durable workflow do next?
Reveal a model answer
“Save the outcome as unknown and use A81 to check with the provider. I do not create A82 just to retry: A81 may already have succeeded.”
Interviewer follow-up
What if the provider has neither idempotency nor a status query?
Reveal the follow-up answer
“The ambiguity cannot be eliminated by our local database alone. I need another provider-supported reconciliation mechanism or a product process that explicitly handles unresolved outcomes.”
What the answer must demonstrate: Do not promise exactly-once effects across an unsupported boundary.
An inventory hold expires at 120 seconds and payment authorization succeeds at 125. Can the order be confirmed?
Reveal a model answer
“Not from the authorization alone. Inventory must atomically verify or convert a valid hold, and H81 is expired. I keep confirmation conditional and void the authorization while cancelling the order.”
Interviewer follow-up
What if a callback races with cancellation?
Reveal the follow-up answer
“Both transitions check the durable workflow state, and inventory independently checks the allocation condition. A late callback cannot overwrite a terminal cancellation.”
What the answer must demonstrate: Two authorities must enforce their own conditions.
What happens if the compensating void also fails?
Reveal a model answer
“The cancellation has an outstanding cleanup state with a stable void identifier. A worker retries or queries it, and an age-based alert exposes work that cannot finish automatically.”
Interviewer follow-up
Can you report that no authorization exists yet?
Reveal the follow-up answer
No. The order can be blocked from fulfillment while authorization cleanup remains pending. A timeout does not prove the void failed, and retrying after the provider’s deduplication window expires may create a different external effect.
What the answer must demonstrate: Do not hide unfinished compensation behind a terminal label.
What guarantee does a transactional outbox add to a distributed workflow?
Reveal a model answer
“It atomically records the local state transition and the intent to send the next message. After a crash, the relay can find that intent. The relay may publish twice, so consumers still need idempotent handling.”
Interviewer follow-up
Does it atomically commit the provider’s authorization?
Reveal the follow-up answer
“No. That remains a remote effect whose uncertain outcome must be reconciled separately.”
What the answer must demonstrate: Keep the outbox guarantee within its actual transaction boundary.
Would you use orchestration or choreography for an order workflow with inventory holds, payment authorization, and compensation?
Reveal a model answer
“I would start with an explicit coordinator because the order’s deadlines, compensation and user-visible status form one workflow that operators must inspect. Services still own inventory and authorization details.”
Interviewer follow-up
When could choreography fit?
Reveal the follow-up answer
“A few independent reactions to a completed fact, such as analytics and notification, may work well as subscribers. I would still assign ownership for failures and avoid an implicit cycle of events nobody can reconstruct.”
What the answer must demonstrate: Explain operational ownership instead of declaring one style universally better.
Blank-page exercise · 18 minutes
Build the answer yourself
Draw O81’s workflow through inventory hold, authorization, confirmation, cancellation and recovery. Inject a crash after every remote success and a late authorization after hold expiry.
- Distinguish a timeout from a confirmed rejection.
- Persist each transition and outgoing intent atomically.
- Give every retried effect a stable identifier and conditional state rule.
- Show who retries compensation and how unresolved work becomes visible.
Check that each component and design decision follows from your requirements and workload.
Recall the key ideas
Answer from memory before opening each card. Explain why the choice works and what it costs. Revisit missed cards tomorrow.
Distributed transactions and sagasWhat does a remote timeout prove?Recall first, then reveal
Only that the caller did not receive a timely answer; the remote effect may already have succeeded.
Timeout means unknown
Return to lessonDistributed transactions and sagasCan a prepared 2PC participant simply time out and abort?Recall first, then reveal
After voting yes it must learn a safe final decision; unilateral timeout abort can contradict an existing commit decision.
Prepared means promised
Return to lessonDistributed transactions and sagasIs compensation a rollback?Recall first, then reveal
It is a new business action after earlier steps committed, so intermediate observations and irreversible effects remain.
Repair forward, not rewind
Return to lessonDistributed transactions and sagasWhat must survive a worker crash halfway through checkout?Recall first, then reveal
The current workflow step, stable IDs for remote actions, saved pending requests and rules for advancing state safely. A replacement worker can then check uncertain results and continue.
Save progress → retry the same action → check the outcome.
Return to lessonFinal revision
Summary and interview notes
Keep related changes in one local transaction when possible. Across databases, use atomic commit if participants support it, or a saga that saves progress after each local step. A saga must recover uncertain results and perform compensating actions when later steps fail.
Remember these points
- 2PC coordinates one commit-or-abort outcome; it does not by itself prove cross-participant isolation.
- A prepared yes voter cannot unilaterally abort merely because the coordinator timed out.
- Saga steps commit locally, so compensation is new business work and can fail too.
- Save an uncertain action as UNKNOWN and reuse its stable ID while checking its result. Do not invent a second action because the first reply was lost.
- Late payment success must not revive an expired hold. Check the current reservation state atomically when allocating or cancelling.
Interview tips
- Draw one crash after remote success but before saving its reply, then show recovery from persisted facts.
- Show the normal path, timeout path and failed-compensation path on the same state machine.
- Before selecting a saga engine, ask whether keeping the related records in one database would let a local transaction satisfy the requirement.
Important qualifications
- An outbox atomically records local state and sending intent; it does not atomically perform the remote effect.
- Provider idempotency retention limits automatic retry safety; a durable workflow engine cannot extend that external contract.
Technical references
- Oracle Database: Distributed Transactions ConceptsOfficial definition of transactions spanning distinct database nodes and coordinated commit or rollback; used here to distinguish transaction terminology from a compensating workflow.
- Garcia-Molina and Salem: SagasPrimary research on composing local transactions with compensating work.
- PostgreSQL: Two-Phase TransactionsVerified reference for externally coordinated prepared transactions.
- PostgreSQL: PREPARE TRANSACTIONPrepared state retains locks and requires an external transaction manager to resolve it.
- AWS: Transactional Outbox PatternVerified reference for atomically recording a state change and its publication intent. All order timings are hypothetical.
- Temporal: Workflow ExecutionPersisted event history and replay support durable progress; application rules and external-effect contracts remain explicit.
- Temporal: Activity DefinitionActivities contain external work and must be designed for retry and idempotency.
Practice marks stay in this browser.