System designby Learnastra

Concept lesson · Foundations

Distributed transactions and sagas

By Anup Rai

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.

The visual modelSaga state transitions and payment reconciliation

A payment timeout is an unknown outcome, not proof of failure. Reconcile before retrying or compensating.

Saga state transitions and payment reconciliationA payment timeout is an unknown outcome, not proof of failure. Reconcile before retrying or compensating. 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.O81: timeout does not prove payment failedPersist intentO81Reserve stockhold H81Authorize $24reply is lostRecord UNKNOWNreconcile the same A81Provider lookupA81 outcomeAuthorization known successful:allocate H81 only while validthen confirm O81If H81 expired: do not confirmvoid late authorization A81Known auth failure: release H81Record and retry compensation. Resolve an unknown payment outcome before compensating.
Read the diagram step by step
  1. Persist order intent, reserve inventory, and call payment with one stable attempt identity.
  2. If the payment response is lost, record UNKNOWN and query or retry that same identity.
  3. 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.
  4. 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 practice

Useful foundations: Transaction isolation · Message queues, event logs, delivery guarantees, and backpressure · Idempotency, retries, and timeouts

Workload and timing examples are interview assumptions.

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.

Concept in focusTwo-phase commit has a prepared middle state

An abort vote leads to abort. A prepared participant cannot safely invent the global decision when the coordinator is unreachable; classic 2PC can block.

Two-phase commit has a prepared middle stateAn abort vote leads to abort. A prepared participant cannot safely invent the global decision when the coordinator is unreachable; classic 2PC can block. 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)CoordinatorParticipant AParticipant BPREPAREPREPAREDurably prepared; YESDurably prepared; YESPersist global COMMIT decisionCOMMIT (retry delivery if needed)COMMIT (same decision)

Remember: Prepare votes; a durable decision; then deliver it.

Read the diagram
  1. Coordinator to Participant A: PREPARE
  2. Coordinator to Participant B: PREPARE
  3. Participant A to Coordinator: Durably prepared; YES
  4. Participant B to Coordinator: Durably prepared; YES
  5. Coordinator to Coordinator: Persist global COMMIT decision
  6. Coordinator to Participant A: COMMIT (retry delivery if needed)
  7. 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

Concept in focusCompensation travels back through completed work

Green arrows move the workflow forward. Rust arrows perform compensating business actions after a definite failure.

Compensation travels back through completed workGreen arrows move the workflow forward. Rust arrows perform compensating business actions after a definite failure. 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.Shipment fails after two local commitsReserve stockAuthorize payShipment failsVoid paymentRelease stockForward workCompensationEarlier commits happened. Compensation performs new business actions.Each action needs safe retries and an explicit unknown-outcome policy.

Remember: A compensation is another action, not erasure of a past commit.

Read the diagram
  1. Follow successful reservation and payment steps, then reverse the business effects after shipment fails.
  2. Reserve stock, authorize payment, then encounter a definitive shipment failure.
  3. 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.

Worked example diagramAfter both hold expiry and known late authorization success, cancellation needs a durable compensating void. A lost void response leaves cleanup pending until that operation is reconciled.
Distributed transactions and sagas: architecture diagram1. O81 requested to 2. Hold H81: 2 mugs: reserve; 2. Hold H81: 2 mugs to 3. Authorize A81: unknown: send once logically; 3. Authorize A81: unknown to 4. H81 expires at 120 s: deadline passes; 3. Authorize A81: unknown to 5. Late A81 success at 125 s: reconcile by A81; 4. H81 expires at 120 s to 6. Cancel order; void A81 pending: cannot allocate; 5. Late A81 success at 125 s to 6. Cancel order; void A81 pending: compensating action; 6. Cancel order; void A81 pending to 7. Void confirmed; cleanup complete: retry or reconcile same void1 → 2: reserve2 → 3: send once logically3 → 4: deadline passes3 → 5: reconcile by A814 → 6: cannot allocate5 → 6: compensating action6 → 7: retry or reconcile same void01O81 requested02Hold H81: 2 mugs03Authorize A81:unknown04H81 expires at 120 s05Late A81 success at125 s06Cancel order; voidA81 pending07Void confirmed;cleanup complete
  1. 1 → 2reserveO81 requested → Hold H81: 2 mugs
  2. 2 → 3send once logicallyHold H81: 2 mugs → Authorize A81: unknown
  3. 3 → 4deadline passesAuthorize A81: unknown → H81 expires at 120 s
  4. 3 → 5reconcile by A81Authorize A81: unknown → Late A81 success at 125 s
  5. 4 → 6cannot allocateH81 expires at 120 s → Cancel order; void A81 pending
  6. 5 → 6compensating actionLate A81 success at 125 s → Cancel order; void A81 pending
  7. 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.

Foundation · Question 1

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.

What the answer must demonstrate: Identify the actual independent commit boundaries.

Foundation · Question 2

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.”

What the answer must demonstrate: Prepared is a durable protocol state, not a best-effort check.

Follow-up · Question 3

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.”

What the answer must demonstrate: Atomic commit and isolation solve different parts of correctness.

Applied · Question 4

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.”

What the answer must demonstrate: Do not promise exactly-once effects across an unsupported boundary.

Applied · Question 5

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.”

What the answer must demonstrate: Two authorities must enforce their own conditions.

Applied · Question 6

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.”

What the answer must demonstrate: Do not hide unfinished compensation behind a terminal label.

Follow-up · Question 7

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.”

What the answer must demonstrate: Keep the outbox guarantee within its actual transaction boundary.

Applied · Question 8

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.”

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 lesson
Distributed 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 lesson
Distributed 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 lesson
Distributed 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 lesson

Final 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

Practice marks stay in this browser.