A saga replaces one distributed transaction with a sequence of local transactions, each in one service, plus compensating actions that semantically undo earlier steps if a later one fails. In the orchestrated style, one component owns the sequence: it tells each participant what to do, waits for the reply, and decides what happens next. The idea is simple. The hard part is the orchestrator itself, because it must survive crashes, duplicate and lost messages, slow participants and concurrent sagas touching the same data.

This page assumes you know what a saga is (see the saga pattern) and focuses on building and operating the orchestrator. We use one worked example throughout: booking a trip that holds a flight seat, holds a hotel room, charges a card, then issues the ticket and confirms the room.

Advertisement

The orchestrator is a persisted state machine

Strip away frameworks and an orchestrator is a state machine whose state lives in a database row per saga: which step it is on, whether it is running forward or compensating, and which command it is waiting for. Each incoming reply or timer is an event; the orchestrator loads the row, decides the transition, writes the new state and emits the next command. It keeps no in-memory state that matters, so any instance can process any event, and a restarted process picks up exactly where the database says it was.

Two rules make that work. First, the state change and the command it causes must commit atomically, otherwise a crash between them either loses the command (the saga hangs) or sends it without recording it (the saga repeats work it does not know about). The standard answer is a transactional outbox: write the command into an outbox table in the same transaction, and let a relay publish it. Second, events for one saga must be processed one at a time, via a row lock or an optimistic version column, so two replies cannot both advance the same step.

A saga orchestrator: state, outbox and timers in one databaseOrchestratorloads saga, decides next stepOrchestrator databasesaga: id, state, step, awaitingoutbox: commands to sendtimer: saga, step, dueall written in ONE transactionOutbox relaypublishes rowsBrokercommands / repliescommandsFlightsHold / Release / IssueHotelsHold / Release / ConfirmPaymentsCharge (pivot)each participant: dedupe table +own outbox for repliesreplies (saga_id, command_id)Every state change and the command it causes commit together, so a crash never loses or invents a step.Delivery is at least once in both directions; correlation ids and dedupe make it effectively once.
The orchestrator writes saga state, outgoing commands and timers in one transaction; participants deduplicate commands and reply through their own outbox.

Step kinds and the pivot

Not every step can be undone, so ordering is a design decision, not an accident. Classify each step:

  • Compensatable: has a compensating action. Holding a seat is undone by releasing it.
  • Pivot: the go/no-go point. If it succeeds the saga must run to completion; if it fails, earlier steps are compensated. Charging the card is the pivot here, because refunds are a business process, not an automatic undo.
  • Retriable: after the pivot, guaranteed to succeed eventually, so they are retried rather than compensated. Issuing the ticket and confirming the room qualify only because the earlier holds reserved the inventory.

The ordering rule follows: compensatable steps first, then one pivot, then retriable steps. Put the steps most likely to fail early, where failing is cheap, and make the pre-pivot steps reservations rather than final effects. If a post-pivot step can fail permanently, it is not retriable, and your design has a hole that will surface as a manual refund queue.

Advertisement

Worked example: the trip orchestrator

The sketch below is framework-neutral Python. db is any store with transactions and row locks. The step list is the saga definition; on_reply is the whole transition function. Command ids are derived from the saga, step and command name, so a retry of the same step sends the same id and participants can deduplicate it.

import uuid
from dataclasses import dataclass

@dataclass(frozen=True)
class Step:
    name: str
    action: str               # command sent to a participant
    compensation: str | None  # None: nothing to undo
    kind: str                 # "compensatable" | "pivot" | "retriable"

@dataclass
class Saga:
    id: str
    state: str                # RUNNING | COMPENSATING | COMPLETED | COMPENSATED
    step: int
    awaiting: str | None      # command_id of the reply we are waiting for
    data: dict

TRIP = [
    Step("hold_flight",  "flights.Hold",    "flights.ReleaseHold", "compensatable"),
    Step("hold_hotel",   "hotels.Hold",     "hotels.ReleaseHold",  "compensatable"),
    Step("charge_card",  "payments.Charge", None,                  "pivot"),
    Step("issue_ticket", "flights.Issue",   None,                  "retriable"),
    Step("confirm_room", "hotels.Confirm",  None,                  "retriable"),
]

def start(db, request):
    saga_id = str(uuid.uuid4())
    s = Saga(id=saga_id, state="RUNNING", step=0, awaiting=None, data=request)
    with db.transaction() as tx:
        dispatch(tx, s, TRIP[0].action, step=0)          # sets s.awaiting
        tx.insert("saga", s)
    return saga_id

def dispatch(tx, s, command, step):
    command_id = f"{s.id}:{step}:{command}"            # stable across retries
    tx.insert("outbox", topic=command.split(".")[0],
              payload={"command": command, "command_id": command_id, "saga_id": s.id})
    tx.upsert("timer", saga_id=s.id, due=now() + TIMEOUT[command])
    s.awaiting = command_id                             # saved by the caller's update

def on_reply(db, reply):
    with db.transaction() as tx:
        s = tx.select_for_update("saga", reply.saga_id)
        if s.awaiting != reply.command_id:
            return                                      # duplicate or stale: ignore
        step = TRIP[s.step]
        if s.state == "RUNNING":
            if reply.status == "OK":
                s.step += 1
                if s.step == len(TRIP):
                    s.state, s.awaiting = "COMPLETED", None
                    tx.delete("timer", saga_id=s.id)
                else:
                    dispatch(tx, s, TRIP[s.step].action, s.step)
            elif reply.status == "RETRY" or step.kind == "retriable":
                retry_later(tx, s, step.action)          # same command_id, backoff
            else:                                       # definitive failure at or before the pivot
                s.state = "COMPENSATING"
                compensate_from(tx, s, s.step - 1)
        elif s.state == "COMPENSATING":
            if reply.status == "OK":
                compensate_from(tx, s, s.step - 1)
            else:
                retry_later(tx, s, step.compensation)    # compensations must eventually succeed
        tx.update("saga", s)

def compensate_from(tx, s, i):
    while i >= 0 and TRIP[i].compensation is None:
        i -= 1
    if i < 0:
        s.state, s.awaiting = "COMPENSATED", None
        tx.delete("timer", saga_id=s.id)
        return
    s.step = i
    dispatch(tx, s, TRIP[i].compensation, i)

Trace a failure. The flight hold succeeds, then the hotel replies FAILED because the property is full. The saga flips to COMPENSATING and dispatches flights.ReleaseHold; when that replies OK, compensate_from walks further back, finds nothing left, and ends in COMPENSATED. Now trace a decline at the card: both holds are released in reverse order. A failure after the charge never compensates; it retries with backoff until the participant succeeds or an operator intervenes.

Participants: idempotent, order-tolerant, honest

Brokers deliver at least once, the orchestrator retries on timeout, and a compensation can overtake the command it undoes when the original was slow. Every participant therefore needs three properties. It records each command_id with its reply and replays that reply on a duplicate. Its compensation succeeds when there is nothing to undo. And when a compensation arrives first, it leaves a tombstone so the late original is refused instead of creating an orphan hold. Replies go out through the participant's own outbox, in the same transaction as the work.

def handle_command(db, msg):
    with db.transaction() as tx:
        seen = tx.get("processed_commands", msg.command_id)
        if seen:
            reply = seen.reply                          # replay the original answer
        elif msg.command == "hotels.ReleaseHold":
            hold = tx.get("holds", msg.saga_id)
            if hold is None:                            # release overtook the hold
                tx.insert("holds", id=msg.saga_id, status="CANCELLED")  # tombstone
            elif hold.status == "HELD":
                tx.update("holds", msg.saga_id, status="RELEASED")
            reply = {"status": "OK"}                    # releasing nothing is success
        elif msg.command == "hotels.Hold":
            hold = tx.get("holds", msg.saga_id)
            if hold and hold.status in ("CANCELLED", "RELEASED"):
                reply = {"status": "FAILED", "reason": "cancelled"}
            else:
                reply = reserve_room(tx, msg)           # OK or FAILED, never half-done
        else:
            reply = handle_other(tx, msg)               # Confirm, and so on
        if not seen:
            tx.insert("processed_commands", id=msg.command_id, reply=reply)
            tx.insert("outbox", topic="saga-replies",
                      payload={**reply, "saga_id": msg.saga_id, "command_id": msg.command_id})

Honest means a participant answers OK or FAILED only when the outcome is final, and RETRY for a transient problem. A participant that times out internally and replies FAILED while its work still commits causes the worst saga bug: the orchestrator compensates something that later happens. The broader dedupe patterns are covered in idempotency architecture.

Timeouts and crash recovery

Silence is the common failure: the command or reply was lost, or the participant is slow. The orchestrator keeps a timer row per saga, written with each dispatch, and a sweeper fires due timers as events. The safe response to a timeout depends on the step:

  • Before the pivot: resend the same command id a few times, then compensate. Because compensations tolerate missing work and leave tombstones, compensating an unknown outcome is safe.
  • At the pivot: never assume. Resend the same idempotent charge, or query the payment provider by idempotency key, and decide only on a definite answer.
  • After the pivot: keep retrying with capped exponential backoff and jitter; past a threshold, page a human but do not abandon the saga.

Crash recovery needs no extra machinery: on restart, the sweeper finds overdue timers and the relay finds unpublished outbox rows. Test it deliberately by killing the orchestrator between every pair of steps in a staging run.

Isolation: what sagas give up and how to compensate

A saga is atomic in the eventual sense and has no isolation: other transactions can see a hold that will later be released, or overwrite data a saga is about to compensate. The usual countermeasures are:

  • Semantic lock: mark records with a pending state (HELD, PENDING_PAYMENT) that other operations respect, and clear it on completion or compensation.
  • Commutative updates: design operations such as credit and debit so order does not matter, which makes compensation safe under concurrency.
  • Reread value: before an update, check the record has not changed since the saga read it, and fail the step if it has.
  • By value: route high-risk requests to stricter flows, for example a distributed lock or manual approval for large transfers.

Workflow engines or hand-rolled?

Durable-execution engines such as Temporal, AWS Step Functions and Camunda provide the persisted state, timers, retries and history that the code above builds by hand. Step Functions, for example, declares retry and catch policies on each task state, and a catch can route to compensation states. Engines buy you visibility, versioned definitions and tested recovery. You pay with a new runtime to operate, determinism or definition-language constraints, and lock-in.

ChooseWhenWatch out for
Hand-rolled state machineFew sagas, simple steps, strong in-house database skillsYou own timers, versioning, tooling and dashboards
Workflow engineMany sagas, long waits, human steps, audit needsWorkflow versioning, engine availability, cost per transition
Choreography insteadTwo or three steps, no pivot subtletiesFlow logic scattered across services, hard to see

Whichever you choose, the participant rules do not change: idempotent commands, order-tolerant compensations and honest replies.

Operating sagas in production

Track counts of sagas by state and definition, the age of the oldest running saga, compensation rate per step and time spent in each step. A rising compensation rate is often the first sign that a downstream service is degraded. Keep a stuck-saga view listing each saga past its deadline with its last command and error, and give operators audited actions: retry step, force compensate (before the pivot only) and mark resolved.

Changing a definition while sagas are in flight is a migration. Store the definition version on each saga row and keep old versions runnable until their sagas drain; never renumber steps under running sagas.

Failure modes

  • State and command written separately: lost or phantom steps after a crash. Use the outbox.
  • Non-idempotent participants: duplicate holds or double charges on redelivery.
  • Compensation that fails on missing work: a saga stuck forever in COMPENSATING.
  • No tombstone: a late hold lands after its release and leaks inventory.
  • Blind compensation at the pivot: refunding or releasing after a charge that actually succeeded.
  • Retriable steps that can fail permanently: manual refunds become a standing process.
  • Unversioned definitions: a deploy reorders steps under running sagas.

What to do next

  1. Write your saga as a step table and label every step compensatable, pivot or retriable; reorder until it fits the rule.
  2. Put saga state, outbox and timers in one database and commit them together.
  3. Derive stable command ids and add a processed-commands table to every participant.
  4. Make every compensation succeed on missing work and leave a tombstone.
  5. Define per-step timeouts and the action on expiry, with special handling at the pivot.
  6. Add dashboards for sagas by state, oldest running saga and compensation rate, plus a stuck-saga runbook.
  7. Kill the orchestrator between each pair of steps in staging and confirm every saga finishes. Then compare with the checkout design in designing an event-driven order system.
Key takeaway: An orchestrated saga is a persisted state machine: each reply or timer loads the saga row, decides one transition and writes the new state and the next command in a single transaction through an outbox. Order steps as compensatable reservations first, one pivot, then retriable steps that earlier holds guarantee. Participants must deduplicate by command id, succeed on compensations with nothing to undo, leave tombstones, and reply honestly. Handle timeouts per step, never compensate blindly at the pivot, add isolation countermeasures, and monitor stuck sagas and compensation rates.