Skip to content
KEDBYTE
Site navigation
How Data Works
Chapter
38

Distributed Transactions and Sagas

Part F · Data Across Machines|3,642 words|about 16 min read|Volume F

38.0 What this chapter gives you#

  1. A local transaction can protect changes inside one database boundary. A workflow touching several independent resources needs an additional agreement or recovery design.
  2. 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.
  3. 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#

  1. 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.
  2. Network failures create intermediate situations: one participant may be ready, another may refuse, and the coordinator may disappear before everyone learns the decision.
  3. 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#

  1. Mira needs both a stock reservation and a delivery slot. Two separate clerks control those resources.
  2. 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.
  3. 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#

  1. A fictional order needs three units from database A and one delivery slot from database B. Initially both are available.
  2. 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.
  3. A single Python try block around both calls therefore does not provide distributed atomicity. The program can catch an exception while the first resource has already changed.
  4. 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#

  1. 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.
  2. Atomic commitment and isolation are distinct. Agreement that all participants commit does not automatically give arbitrary external readers a single synchronized multi-database snapshot.
  3. 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#

  1. 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]
  2. 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]
  3. 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#

  1. 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#

  1. 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.
  2. 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.
  3. 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#

  1. 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.
  2. Once a clerk has made that promise, independently releasing the resource while others commit can break the all-or-nothing agreement.
  3. 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#

  1. 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.
  2. The coordinator durably records COMMIT under the protocol. A receives the decision and commits. B’s decision message is delayed.
  3. B remains prepared until it obtains the valid decision; it must not invent ABORT merely because the reply took longer than expected.
  4. 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#

  1. 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.
  2. 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.
  3. 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#

  1. PostgreSQL PREPARE TRANSACTION detaches the prepared transaction from the session and stores state needed for later COMMIT PREPARED or ROLLBACK PREPARED. The transaction remains a resource-management obligation. [S156]
  2. 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.
  3. 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#

  1. 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#

  1. 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.
  2. This is one cost of coordinated atomic commitment. A failed coordinator or missing decision path can leave otherwise healthy resources waiting.
  3. 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#

  1. 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.
  2. Releasing the slot because “surely the coordinator gave up” can leave a confirmed order without the promised delivery resource.
  3. 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#

  1. Continue TX-AB-1: A is committed, B is prepared and the coordinator is temporarily unreachable. B’s local timeout fires.
  2. The timeout establishes that B has not learned the outcome within its waiting policy. It does not prove the global decision was abort.
  3. 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.
  4. 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#

  1. Prepared work can hold locks and prevent cleanup. Its age and resource footprint need monitoring, especially if a transaction manager fails repeatedly.
  2. 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.
  3. 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#

  1. 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]
  2. 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]
  3. 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#

  1. 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#

  1. 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.
  2. 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.
  3. 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#

  1. Mira reserves three notebooks for a delivery that later becomes impossible. She records a cancellation and releases those three notebooks back to availability.
  2. She does not reset the entire shelf count to what it was yesterday, because other customers may have bought or reserved items since then.
  3. 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#

  1. Start a separate stock fixture at 5. Saga reservation S1 takes 3, leaving 2. Another legitimate reservation S2 takes 1, leaving 1.
  2. 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.
  3. 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.
  4. 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#

  1. 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.
  2. 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.
  3. 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#

  1. The original saga formulation decomposes long-lived work into transactions that can interleave with other work and pairs appropriate steps with compensating transactions. [S157]
  2. Compensation preserves application meaning rather than necessarily restoring identical prior bytes. Its correctness depends on concurrent work, allowed transitions and business policy.
  3. 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#

  1. 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#

  1. 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.
  2. Each retry must refer to the same intended step unless a new operation is explicitly being requested. Otherwise the recovery process can create duplicates.
  3. 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#

  1. 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.
  2. A new worker can resume from the recorded state rather than starting every task again from memory.
  3. 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#

  1. Define fictional saga states NEW, STOCK_RESERVED, DELIVERY_CONFIRMED, COMPENSATING, COMPENSATED and REVIEW_REQUIRED.
  2. After the stock step commits, the delivery request times out. The correct next state may be REVIEW_REQUIRED or a pending status query because the slot could already have been confirmed remotely.
  3. 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.
  4. 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#

  1. 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.
  2. 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.
  3. 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#

  1. 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.
  2. 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]
  3. 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#

  1. 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#

  1. 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.
  2. “Use microservices” does not answer whether a half-completed order is acceptable. “Use transactions” does not make an external partner support a prepare protocol.
  3. 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#

  1. 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.
  2. Staff know which interruptions can be retried, which need review and which require cancelling a specific prior commitment.
  3. 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#

  1. 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.
  2. 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.
  3. 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.
  4. 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#

  1. Write down the invariant, the transaction boundary and every external effect. Mark which resources can atomically participate and which can only acknowledge separate actions.
  2. For each interruption point, identify known committed state, unknown outcomes, recovery evidence and permitted next actions. Include compensation failures and stale responses.
  3. 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#

  1. 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]
  2. 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.
  3. 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#

  1. 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#

  1. Do two successful local transaction APIs automatically form one distributed transaction? No. Their decisions remain independent without a coordinating protocol.
  2. 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.
  3. Why monitor old prepared transactions? They retain resources and can interfere with cleanup; they require a responsible resolution mechanism.
  4. 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.
  5. What should a repeated S1 compensation do? Return or recognize the already recorded compensation, not release the quantity again.
  6. 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.
  7. Does a saga restore exactly the same historical bytes? Not necessarily. Compensation is a new semantic action and preserves the fact of prior work.
  8. 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#

  1. Wrong: one exception handler makes remote calls atomic. Right: an earlier resource may already have committed.
  2. Wrong: prepared means finished. Right: prepared work awaits a valid final decision and can retain locks.
  3. Wrong: a timeout proves global abort. Right: it may leave a participant or caller in doubt.
  4. Wrong: two-phase commit is the same as consensus. Right: atomic commitment and replicated agreement have different roles and assumptions.
  5. Wrong: compensation is ordinary rollback. Right: it is a new action after earlier work has committed.
  6. Wrong: restoring an old total is always a valid compensation. Right: it can erase legitimate intervening work.
  7. Wrong: sagas remove the need for idempotency. Right: original and compensating steps both need replay handling.
  8. Wrong: every failure should trigger automatic cancellation. Right: unknown outcomes and irreversible effects require explicit policy.

38.99 Chapter summary in 20 lines#

  1. Local transaction guarantees stop at their resource boundary.
  2. Independent resources need an explicit coordination or recovery design.
  3. Atomic commitment is distinct from global isolation and simultaneous visibility.
  4. Two-phase commit separates preparation from the final decision.
  5. Prepared participants retain a promise and relevant recovery state.
  6. A durable decision must be recoverable after interruption.
  7. Participants can learn the same outcome at different times.
  8. In-doubt work cannot be resolved safely by guessing.
  9. Prepared transactions require monitoring and recovery ownership.
  10. Compensation is a new business action, not time travel.
  11. Correct compensation preserves unrelated intervening work.
  12. Repeated compensation needs stable identity and duplicate handling.
  13. Sagas allow intermediate committed states under defined policy.
  14. Workflow state must distinguish pending, failed and unknown outcomes.
  15. Lost external replies require reconciliation rather than blind repetition.
  16. Outboxes connect local intent to later delivery without global atomicity.
  17. Compensation can fail and needs bounded operational handling.
  18. Architecture should follow invariants and participant capabilities.
  19. User-visible labels must match the actual workflow state.
  20. Educational transition tests do not certify a distributed transaction service.

Return to contents