Distributed Transactions and Sagas
Introductions, exercises and summaries stay visible.
38.0 What this chapter gives you#
- A local transaction can protect changes inside one database boundary. A workflow touching several independent resources needs an additional agreement or recovery design.
- You will distinguish two-phase commit from a sequence of local transactions with compensation, understand prepared and in-doubt work, and design explicit policies for retries and irreversible effects.
- The examples use fictional reservations and delivery arrangements. The companion models are state-transition exercises, not an implementation of a distributed transaction manager or permission to run prepared transactions on a live server.
38.1 One transaction across participants#
38.1.1 PLAIN — in simple words#
- Suppose an operation must change two databases as one logical decision. A successful commit in the first database does not force the second database to commit.
- Network failures create intermediate situations: one participant may be ready, another may refuse, and the coordinator may disappear before everyone learns the decision.
- The design must say whether it needs one coordinated commit decision or whether intermediate completed steps are allowed and later corrected through business actions.
38.1.2 PLAIN — a picture in your head#
- Mira needs both a stock reservation and a delivery slot. Two separate clerks control those resources.
- Asking one clerk and then the other can leave a stock reservation without a delivery slot. The workflow needs a rule for that partial outcome.
- Where the comparison breaks: a database can sometimes hold prepared transactional state under a formal protocol. A delivery company may offer only ordinary requests and cancellations, not participation in the same transaction mechanism.
38.1.3 PLAIN — a worked example#
- A fictional order needs three units from database A and one delivery slot from database B. Initially both are available.
- A local transaction in A commits the stock reservation. Before B commits, the connection fails. Rolling back the caller’s remaining code does not reverse A’s completed transaction.
- A single Python
tryblock around both calls therefore does not provide distributed atomicity. The program can catch an exception while the first resource has already changed. - Before choosing a mechanism, ask whether both resources support coordinated prepare/commit, whether temporary partial states are acceptable and what a valid compensation would mean.
38.1.4 PLAIN — what is really happening inside#
- Independent resource managers own their local locks, logs and recovery. A coordinator must connect their decisions through a protocol if one global atomic decision is required.
- Atomic commitment and isolation are distinct. Agreement that all participants commit does not automatically give arbitrary external readers a single synchronized multi-database snapshot.
- A simpler architecture may keep tightly coupled data within one transactional owner. Distribution should follow requirements, not create a cross-resource invariant without a plan to enforce it.
38.1.5 TECHNICAL — the engineer’s version#
- A distributed transaction coordinates multiple transactional resources. Two-phase commit is an atomic-commit protocol; it is not identical to replicated-log consensus or to a complete distributed isolation model. [S155] [S156]
- PostgreSQL prepared transactions are intended for external transaction managers. Ordinary application sessions should not create long-lived prepared work without a reliable mechanism to resolve it. [S156]
- Review resource capability, coordinator durability, identity, locks, failure detection and recovery ownership. A local connection library cannot promise guarantees unsupported by the participants.
38.1.6 WORDS — remember these#
Participant: a resource involved in a coordinated transaction — a resource manager responsible for its local prepared and final state. Coordinator: the component managing a shared decision — a protocol participant that gathers votes and records or distributes the transaction outcome. Atomic commitment: one compatible commit-or-abort decision — a property distinct from simultaneous visibility or complete transaction isolation.
38.2 Prepare and commit#
38.2.1 PLAIN — in simple words#
- Two-phase commit first asks whether each participant can promise to complete its part. A participant that says yes prepares enough state to honor the eventual valid decision.
- If the required participants prepare successfully, the coordinator can record the commit decision and tell them to commit. A refusal can lead to an abort decision instead.
- Preparing is not the same as finishing. The participant may hold resources while waiting for the final decision.
38.2.2 PLAIN — a picture in your head#
- The stock clerk and delivery clerk each set aside the necessary resources and say, “I am ready to follow the final decision.” The coordinator then announces the agreed outcome.
- Once a clerk has made that promise, independently releasing the resource while others commit can break the all-or-nothing agreement.
- Where the comparison breaks: a real prepared promise must survive the relevant crashes and be tied to a durable transaction identity. Verbal readiness alone is not enough.
38.2.3 PLAIN — a worked example#
- Give the synthetic transaction identity
TX-AB-1. A and B begin in an initial state. Each successfully records a prepared state for that same transaction. - The coordinator durably records
COMMITunder the protocol. A receives the decision and commits. B’s decision message is delayed. - B remains prepared until it obtains the valid decision; it must not
invent
ABORTmerely because the reply took longer than expected. - When B later receives the commit decision, it completes. Repeated delivery of that same decision should be recognized as the same transaction outcome, not applied as a new business operation.
38.2.4 PLAIN — what is really happening inside#
- The prepare phase establishes the participants’ ability to follow the eventual decision. The decision and recovery records must remain available after failures under the implementation’s assumptions.
- Participants can learn the same final outcome at different times. Protocol atomicity is about compatible decisions, not all disks changing at the same physical instant.
- The coordinator cannot safely forget a transaction merely because it sent the decision once. It needs a completion and recovery policy for participants that were disconnected.
38.2.5 TECHNICAL — the engineer’s version#
- PostgreSQL
PREPARE TRANSACTIONdetaches the prepared transaction from the session and stores state needed for laterCOMMIT PREPAREDorROLLBACK PREPARED. The transaction remains a resource-management obligation. [S156] - Required coordinator and participant records must be durable before the protocol makes dependent promises. The precise recovery algorithm belongs to the transaction manager, not to an ad hoc timeout handler.
- The book’s state model accepts only valid transitions and repeated matching decisions. It does not persist coordinator logs or implement network recovery, so it is not a two-phase-commit service.
38.2.6 WORDS — remember these#
Prepare: promise readiness for a later decision — a state in which a participant retains the resources and recovery information needed by the commit protocol. Two-phase commit: prepare followed by a final decision — an atomic-commit protocol across participating transactional resources. Decision record: retained evidence of commit or abort — information needed to resolve participants consistently after interruption.
38.3 In-doubt work#
38.3.1 PLAIN — in simple words#
- A participant is in doubt when it has made a prepared promise but does not know the final decision. It may be unable to release resources safely on its own.
- This is one cost of coordinated atomic commitment. A failed coordinator or missing decision path can leave otherwise healthy resources waiting.
- An operator needs a documented resolution procedure. Guessing a decision to clear a lock can create a contradiction with participants that already completed the opposite outcome.
38.3.2 PLAIN — a picture in your head#
- The delivery clerk has promised to hold a slot until the coordinator decides. The phone line then fails. The stock clerk may already have been told to commit.
- Releasing the slot because “surely the coordinator gave up” can leave a confirmed order without the promised delivery resource.
- Where the comparison breaks: real protocols may replicate coordination state or use specialized recovery procedures. The simple story explains the uncertainty, not every implementation’s availability properties.
38.3.3 PLAIN — a worked example#
- Continue
TX-AB-1: A is committed, B is prepared and the coordinator is temporarily unreachable. B’s local timeout fires. - The timeout establishes that B has not learned the outcome within its waiting policy. It does not prove the global decision was abort.
- A valid recovery path must obtain authoritative decision evidence. If such evidence is unavailable, the system records the unresolved state and follows its approved incident procedure instead of inventing a successful resolution.
- A test confirming that B remains prepared in this model demonstrates the waiting limitation. It does not mean the missing recovery service has been implemented.
38.3.4 PLAIN — what is really happening inside#
- Prepared work can hold locks and prevent cleanup. Its age and resource footprint need monitoring, especially if a transaction manager fails repeatedly.
- The recovery procedure must bind an action to the exact transaction and evidence. A similarly named request or an old log from another environment is not enough.
- Prevention includes bounded operational handling, reliable coordinator state and ownership for abandoned work. Disabling a warning does not remove the retained locks or the uncertainty.
38.3.5 TECHNICAL — the engineer’s version#
- Two-phase commit can block in failure cases when a prepared participant cannot determine the decision safely. Replicating the coordinator may improve availability but does not erase the underlying protocol obligations. [S155] [S156]
- PostgreSQL warns that long-lived prepared transactions retain locks and interfere with vacuum and transaction-ID management. Without a responsible external transaction manager, the feature should remain disabled. [S156]
- Operational resolution is not a substitute for protocol evidence. Record transaction identity, participant states, coordinator decision evidence and any explicitly authorized exceptional disposition.
38.3.6 WORDS — remember these#
In-doubt transaction: prepared work lacks a known final decision — a state requiring authoritative recovery information before safe resolution. Blocking: progress waits for unavailable coordination — a limitation distinct from a participant process being physically stopped. Heuristic resolution: an exceptional locally chosen disposition — a decision outside normal coordinated evidence that can create inconsistency and requires explicit treatment.
38.4 Compensation#
38.4.1 PLAIN — in simple words#
- Compensation is a new business action intended to address a previously completed action. It is not time travel and not the same as rolling back an uncommitted local transaction.
- Cancelling a reservation can release stock. It does not erase the fact that the reservation existed, and it must not undo unrelated work performed meanwhile.
- Some effects cannot be fully undone. A parcel already delivered or a message already read may require a different response rather than pretending the original event disappeared.
38.4.2 PLAIN — a picture in your head#
- Mira reserves three notebooks for a delivery that later becomes impossible. She records a cancellation and releases those three notebooks back to availability.
- She does not reset the entire shelf count to what it was yesterday, because other customers may have bought or reserved items since then.
- Where the comparison breaks: the compensation must follow the real business rules and authority. An arithmetic inverse is not automatically a valid or permitted reversal.
38.4.3 PLAIN — a worked example#
- Start a separate stock fixture at 5. Saga reservation S1 takes 3, leaving 2. Another legitimate reservation S2 takes 1, leaving 1.
- If S1 is cancelled, releasing S1’s three units leaves 4. Resetting stock to S1’s old value of 5 would incorrectly erase S2’s reservation.
- The cancellation must be tied to S1 and applied once under its own idempotency rules. A repeated cancellation message must not add another three units and raise the quantity to 7.
- This fixture demonstrates why compensation is a state-aware new operation. It does not model actual warehouse movement or establish that every reservation may legally or commercially be cancelled.
38.4.4 PLAIN — what is really happening inside#
- A compensating operation checks the current state and the prior operation’s identity. It records what it reversed or adjusted and the outcome of that attempt.
- Other transactions may observe intermediate saga states. Reports must distinguish pending, completed, compensating and compensated work instead of forcing everything into a final-success column.
- Compensation itself can fail or time out. Its retry and escalation policy must be as deliberate as the original action’s policy.
38.4.5 TECHNICAL — the engineer’s version#
- The original saga formulation decomposes long-lived work into transactions that can interleave with other work and pairs appropriate steps with compensating transactions. [S157]
- Compensation preserves application meaning rather than necessarily restoring identical prior bytes. Its correctness depends on concurrent work, allowed transitions and business policy.
- Use stable step identities, conditional transitions and retained outcomes. A reversal amount or released quantity must be derived from the specific original operation, not from untrusted repeated input alone.
38.4.6 WORDS — remember these#
Compensation: a new action addressing an earlier completed action — an application-defined correction distinct from local transaction rollback. Saga: a workflow of local transactions with recovery behaviour — a sequence whose partial completion is handled through continuation or compensation. Semantic reversal: restoring an intended business condition — an operation whose meaning is not necessarily identical to restoring old storage bytes.
38.5 Saga state and retries#
38.5.1 PLAIN — in simple words#
- A saga needs a durable record of its progress. After a process restart, the system should know which steps completed, which are pending and which outcomes remain unknown.
- Each retry must refer to the same intended step unless a new operation is explicitly being requested. Otherwise the recovery process can create duplicates.
- A final label such as “cancelled” should be assigned only after the required cancellation conditions are established, not merely because a cancellation message was queued.
38.5.2 PLAIN — a picture in your head#
- Mira uses a checklist for each delivery: reserve stock, arrange delivery, confirm dispatch. If the delivery step fails, the checklist records whether stock still needs to be released.
- A new worker can resume from the recorded state rather than starting every task again from memory.
- Where the comparison breaks: durable workflow state must be connected atomically to local changes and outgoing intents where required. A handwritten checkmark after an external call has the same lost-reply ambiguity as other distributed actions.
38.5.3 PLAIN — a worked example#
- Define fictional saga states
NEW,STOCK_RESERVED,DELIVERY_CONFIRMED,COMPENSATING,COMPENSATEDandREVIEW_REQUIRED. - After the stock step commits, the delivery request times out. The
correct next state may be
REVIEW_REQUIREDor a pending status query because the slot could already have been confirmed remotely. - Blindly releasing stock and reporting cancellation while the remote delivery remains confirmed can create a new inconsistency. First reconcile the delivery request using its stable identity and supported cancellation/status contract.
- If delivery is definitively refused, transition to compensation, release the specific stock reservation once, then record the compensated outcome. The state names are example policy, not a universal workflow standard.
38.5.4 PLAIN — what is really happening inside#
- Local state transitions can use version checks or row locks. Outbox records can connect a committed local transition with a later message without pretending that message delivery is part of the same local transaction.
- Incoming messages need duplicate and conflict handling. A late success may arrive after a timeout initiated review; the workflow must reconcile it with the current state rather than ignoring it blindly.
- Keep retries bounded operationally and retain enough evidence for manual review. A permanently failing compensation must not remain an invisible infinite loop.
38.5.5 TECHNICAL — the engineer’s version#
- A saga is a state machine with durable step identities, transition guards and explicit recovery policy. Orchestration centralizes coordination; choreography distributes reactions, but neither removes the need for observable workflow state.
- The transactional outbox and idempotent consumer patterns address local intent and replay boundaries. They do not establish a globally atomic saga or guaranteed exactly-once external effects. [S125] [S126]
- Persist unknown outcomes distinctly from confirmed failures. Recovery must use authoritative status or supported idempotency mechanisms before choosing a potentially conflicting compensation.
38.5.6 WORDS — remember these#
Workflow state: the recorded stage and evidence of a process — durable information used to resume or reconcile a multi-step operation. Orchestration: a coordinator directs workflow steps — an arrangement that centralizes sequencing and recovery decisions. Choreography: participants react to events — an arrangement distributing workflow reactions while still requiring explicit correctness and observability.
38.6 Choosing an explicit failure policy#
38.6.1 PLAIN — in simple words#
- The mechanism should follow the requirement. Some work belongs in one local transaction; some needs coordinated atomic commitment; some can tolerate visible intermediate states with carefully designed compensation.
- “Use microservices” does not answer whether a half-completed order is acceptable. “Use transactions” does not make an external partner support a prepare protocol.
- The chosen policy must say what happens when a step fails, a reply is lost, a participant is unavailable or compensation cannot finish.
38.6.2 PLAIN — a picture in your head#
- Before opening a second counter, Mira writes down what a customer can be promised at each stage. A tentative reservation is not called a dispatched order.
- Staff know which interruptions can be retried, which need review and which require cancelling a specific prior commitment.
- Where the comparison breaks: written policy is necessary but not sufficient. The actual implementation must enforce it under concurrent requests and failures.
38.6.3 PLAIN — a worked example#
- Compare three designs for the fictional stock-and-delivery workflow. Design A keeps both resources in one local database transaction. Design B uses participants capable of managed two-phase commit. Design C records local steps and defined compensations.
- A may be simplest when the resources truly share one owner and database. B preserves a coordinated decision but introduces prepared-state and recovery obligations. C permits intermediate committed states and needs an explicit semantic recovery policy.
- None is universally best. A delivery partner exposing only a request API may rule out B; a rule forbidding visible partial completion may rule out a naive C.
- Evaluate the allowed states and failure traces before selecting the architecture, then test the implemented design rather than the diagram alone.
38.6.4 PLAIN — what is really happening inside#
- Write down the invariant, the transaction boundary and every external effect. Mark which resources can atomically participate and which can only acknowledge separate actions.
- For each interruption point, identify known committed state, unknown outcomes, recovery evidence and permitted next actions. Include compensation failures and stale responses.
- Review the user-visible labels against those states. A screen that says “complete” while compensation remains pending conceals the very uncertainty the protocol needs to manage.
38.6.5 TECHNICAL — the engineer’s version#
- Compare designs by their supported invariants and failure semantics, not by treating two-phase commit and sagas as interchangeable syntax. Prepared atomic commitment and compensating workflows solve different versions of the problem. [S156] [S157]
- The companion state models test valid prepare/decision transitions and idempotent compensation arithmetic. No external transaction manager, delivery provider or production distributed workflow is executed.
- A release review requires implementation-level evidence for concurrency, durable recovery, authorization and external-effect handling. Local educational tests are useful but narrower evidence.
38.6.6 WORDS — remember these#
Failure policy: the permitted response to an interruption — explicit rules for retry, waiting, reconciliation, compensation or escalation. Intermediate state: a visible stage before final completion — a condition that a saga-based application must interpret correctly. Recovery ownership: who is responsible for unresolved work — an operational assignment ensuring pending and in-doubt operations are not abandoned.
38.97 Practice and worked answers#
- Do two successful local transaction APIs automatically form one distributed transaction? No. Their decisions remain independent without a coordinating protocol.
- A is committed and B is prepared. Can B abort solely because its timeout fired? Not safely under the illustrated coordinated decision. It needs authoritative outcome evidence.
- Why monitor old prepared transactions? They retain resources and can interfere with cleanup; they require a responsible resolution mechanism.
- Stock starts at 5; S1 reserves 3; S2 reserves 1; S1 is compensated. What remains? Four units. Restoring the old value of 5 would erase S2’s effect.
- What should a repeated S1 compensation do? Return or recognize the already recorded compensation, not release the quantity again.
- A delivery request times out. Is cancellation automatically the next valid step? No. The remote request may have succeeded; reconcile its outcome and supported cancellation contract first.
- Does a saga restore exactly the same historical bytes? Not necessarily. Compensation is a new semantic action and preserves the fact of prior work.
- What does the educational model leave unimplemented? Durable distributed coordination, real network recovery, provider integration, authentication and production concurrency handling.
38.98 Common wrong ideas#
- Wrong: one exception handler makes remote calls atomic. Right: an earlier resource may already have committed.
- Wrong: prepared means finished. Right: prepared work awaits a valid final decision and can retain locks.
- Wrong: a timeout proves global abort. Right: it may leave a participant or caller in doubt.
- Wrong: two-phase commit is the same as consensus. Right: atomic commitment and replicated agreement have different roles and assumptions.
- Wrong: compensation is ordinary rollback. Right: it is a new action after earlier work has committed.
- Wrong: restoring an old total is always a valid compensation. Right: it can erase legitimate intervening work.
- Wrong: sagas remove the need for idempotency. Right: original and compensating steps both need replay handling.
- Wrong: every failure should trigger automatic cancellation. Right: unknown outcomes and irreversible effects require explicit policy.
38.99 Chapter summary in 20 lines#
- Local transaction guarantees stop at their resource boundary.
- Independent resources need an explicit coordination or recovery design.
- Atomic commitment is distinct from global isolation and simultaneous visibility.
- Two-phase commit separates preparation from the final decision.
- Prepared participants retain a promise and relevant recovery state.
- A durable decision must be recoverable after interruption.
- Participants can learn the same outcome at different times.
- In-doubt work cannot be resolved safely by guessing.
- Prepared transactions require monitoring and recovery ownership.
- Compensation is a new business action, not time travel.
- Correct compensation preserves unrelated intervening work.
- Repeated compensation needs stable identity and duplicate handling.
- Sagas allow intermediate committed states under defined policy.
- Workflow state must distinguish pending, failed and unknown outcomes.
- Lost external replies require reconciliation rather than blind repetition.
- Outboxes connect local intent to later delivery without global atomicity.
- Compensation can fail and needs bounded operational handling.
- Architecture should follow invariants and participant capabilities.
- User-visible labels must match the actual workflow state.
- Educational transition tests do not certify a distributed transaction service.