Implementing Sagas with Compensating Jobs

A saga is a sequence of local transactions where every step that changes the outside world has a compensating step that undoes it, and this guide builds one on an ordinary job queue as part of Job Chaining & Workflow Orchestration in Queue Fundamentals & Architecture. It is the practical answer to "what happens to steps one and two when step three fails?" when those steps touched a payment provider, an inventory service, and a shipping API that share no transaction.

Problem Statement

A checkout workflow runs four jobs: reserve stock, charge the card, create a shipment, and send a confirmation. Roughly one checkout in 2,000 fails at the shipment step because the carrier API rejects the address. Today that leaves the customer charged, the stock reserved, and no shipment — support refunds by hand. You want every failed checkout to end in a consistent state automatically: stock released and payment refunded, exactly once each, even if a worker crashes halfway through the undo.

Prerequisites

  • A job framework with retries and a dead-letter destination (Celery, Sidekiq, BullMQ, or plain SQS consumers all work).
  • A relational database for the saga log, in the same database as the order record if possible.
  • External APIs that accept idempotency keys, or local tables that let you make the calls idempotent yourself.
  • A clear list of which steps have side effects visible outside your system, since only those need compensation.

Step 1 — List Each Step with Its Compensation

Start on paper. For each step, write the forward action, the compensating action, and whether the compensation can fail. Order the steps so that the ones hardest to undo run last.

Step Forward action Compensation Notes
1 Reserve stock Release reservation Cheap, always possible
2 Authorize payment Void authorization Voiding is cheaper than refunding a capture
3 Create shipment Cancel shipment Possible only before pickup
4 Capture payment (pivot — no compensation) After this the saga must go forward
5 Send confirmation email None needed Retry until it succeeds

Two design choices come out of this table. First, authorize and capture are split so the irreversible money movement happens after the steps most likely to fail. Second, step 4 is the pivot: once it succeeds, the saga only moves forward and every later step is retried until it succeeds rather than compensated. Steps after the pivot must be retryable indefinitely, which is why the email step comes last.

Forward until the pivot, backward before it Forward steps run left to right: reserve stock, authorize payment, create shipment, capture payment, send email. Capture is the pivot. If create shipment fails, compensations run right to left: void the authorization, then release the stock. After capture succeeds, later steps are retried rather than compensated. Checkout saga reserve stock authorize shipment fails capture (pivot) email release stock void auth compensate in reverse order Steps after the pivot are never compensated; they retry until they succeed.

Step 2 — Create a Saga Log

The saga log is the durable record of which steps have completed. It is what lets compensation know what to undo, and what lets a crashed coordinator resume.

CREATE TABLE sagas (
  id          uuid PRIMARY KEY,
  kind        text NOT NULL,                 -- 'checkout'
  subject_id  text NOT NULL,                 -- order id
  state       text NOT NULL DEFAULT 'forward',  -- forward | compensating | done | failed
  created_at  timestamptz NOT NULL DEFAULT now(),
  updated_at  timestamptz NOT NULL DEFAULT now()
);

CREATE TABLE saga_steps (
  saga_id     uuid REFERENCES sagas(id),
  step        text NOT NULL,
  status      text NOT NULL,                 -- done | compensated
  result      jsonb,                         -- ids needed to undo: reservation_id, auth_id
  updated_at  timestamptz NOT NULL DEFAULT now(),
  PRIMARY KEY (saga_id, step)
);

Store in result whatever the compensation needs — the reservation id, the authorization id, the shipment id. Compensation must never have to rediscover what the forward step did.

Step 3 — Write Forward Steps That Record Themselves

Each forward job does its work with an idempotency key derived from the saga id, records completion in the saga log, and enqueues the next step. If the external call succeeds but the worker dies before recording, the retry repeats the call with the same key and gets the same result back.

# checkout_saga.py (Celery flavour; the shape is identical in Sidekiq or BullMQ)
STEPS = ["reserve_stock", "authorize_payment", "create_shipment", "capture_payment", "send_email"]
PIVOT = "capture_payment"

@app.task(bind=True, acks_late=True, max_retries=5, retry_backoff=True)
def run_step(self, saga_id: str, step: str) -> None:
    saga = sagas.get(saga_id)
    if saga.state != "forward":
        return                                            # saga already compensating or done
    if saga_steps.is_done(saga_id, step):
        return enqueue_next(saga_id, step)                # redelivery after completion
    try:
        result = HANDLERS[step](saga, idempotency_key=f"{saga_id}:{step}")
    except PermanentStepError as exc:
        if steps_before_pivot(step):
            return begin_compensation(saga_id, failed_step=step, error=str(exc))
        raise                                             # after the pivot: keep retrying
    except TransientStepError as exc:
        raise self.retry(exc=exc)
    saga_steps.mark_done(saga_id, step, result)           # INSERT ... ON CONFLICT DO NOTHING
    enqueue_next(saga_id, step)

def enqueue_next(saga_id: str, step: str) -> None:
    i = STEPS.index(step)
    if i + 1 < len(STEPS):
        run_step.delay(saga_id, STEPS[i + 1])
    else:
        sagas.set_state(saga_id, "done")

The distinction between PermanentStepError and TransientStepError is doing the real work: transient failures retry the same step; permanent ones trigger compensation. Classifying errors by type is covered in retrying only transient errors by exception type.

Step 4 — Run Compensations in Reverse, Idempotently

Compensation flips the saga state, then undoes every completed step in reverse order. Each compensation is its own job with its own retries, because undo actions call external APIs that can fail too.

def begin_compensation(saga_id: str, failed_step: str, error: str) -> None:
    with db.transaction() as tx:
        flipped = tx.execute(
            "UPDATE sagas SET state='compensating', updated_at=now() "
            "WHERE id=:id AND state='forward'", id=saga_id).rowcount
    if flipped:                                           # only the first failure starts it
        done = saga_steps.done_steps(saga_id)             # ordered by STEPS index
        compensate.delay(saga_id, list(reversed(done)))

@app.task(bind=True, acks_late=True, max_retries=None, retry_backoff=True, retry_backoff_max=600)
def compensate(self, saga_id: str, remaining: list[str]) -> None:
    if not remaining:
        sagas.set_state(saga_id, "failed")                # consistent, but not completed
        notify_customer_checkout_failed(saga_id)
        return
    step, rest = remaining[0], remaining[1:]
    record = saga_steps.get(saga_id, step)
    if record.status != "compensated":
        UNDO[step](record.result, idempotency_key=f"{saga_id}:undo:{step}")
        saga_steps.set_status(saga_id, step, "compensated")
    compensate.delay(saga_id, rest)

max_retries=None is intentional: a compensation that gives up leaves money or stock in limbo, which is worse than a compensation that keeps trying and alerts. Cap the backoff instead, and page someone if a compensation has been retrying for more than an hour.

The saga log drives the undo The saga log lists reserve_stock and authorize_payment as done with their result ids. After create_shipment fails, the compensation job reads the log, voids the authorization using its stored id, marks it compensated, then releases the reservation and marks that compensated, and finally sets the saga to failed. saga_steps for checkout 7f2c step result status reserve_stock reservation_id: r-551 compensated (2nd) authorize_payment auth_id: pi_3Nf... compensated (1st) create_shipment never recorded "done", so there is nothing to cancel for it.

Step 5 — Recover Sagas Stuck Mid-Flight

A saga can stall if a job is lost (a broker without persistence, a deleted queue) or dead-lettered. A sweeper finds sagas that have not moved recently and re-enqueues the step they should be on.

@app.task
def sweep_stuck_sagas() -> None:
    stuck = db.fetch_all(
        "SELECT id, state FROM sagas WHERE state IN ('forward','compensating') "
        "AND updated_at < now() - interval '15 minutes' LIMIT 500")
    for s in stuck:
        if s.state == "forward":
            next_step = first_not_done(s.id)
            run_step.delay(str(s.id), next_step)          # idempotent by design
        else:
            remaining = [st for st in reversed(saga_steps.done_steps(s.id))]
            compensate.delay(str(s.id), remaining)
        metrics.sagas_resumed.inc()

Run it every few minutes with a scheduler — see cron-style scheduling with Celery beat. Because every step and compensation is idempotent, re-enqueueing something that is actually still running is harmless.

Step 6 — Make Every External Call Safe to Repeat

Everything above assumes that repeating a step or a compensation is harmless. That is only true if the external systems cooperate. A worker can die at three points in a step — before the call, after the call but before recording it, or after recording — and only the middle case is dangerous: the side effect happened, but the saga log does not know. The retry will call again.

The fix depends on what the external API offers:

  • Idempotency keys. Payment providers (Stripe, Adyen, Braintree) and many modern APIs accept a client-supplied key and return the original result for a repeated request. Derive the key from the saga id and step name, as the code above does, and the retry becomes a lookup.
  • Natural uniqueness. Some operations have a natural unique field — a shipment reference, an order number sent to the warehouse. Send your saga id as that reference and treat "already exists" as success, fetching the existing object's id.
  • Read-before-write. When neither exists, query the external system for an object tagged with your saga id before creating one. It narrows the window rather than closing it, so pair it with reconciliation.
def create_shipment(saga, idempotency_key: str) -> dict:
    try:
        shp = carrier.create_shipment(
            reference=idempotency_key,                  # carrier rejects duplicate references
            address=saga.order.address, parcels=saga.order.parcels)
    except carrier.DuplicateReference:
        shp = carrier.find_shipment(reference=idempotency_key)   # the first attempt succeeded
    except carrier.AddressRejected as exc:
        raise PermanentStepError(str(exc))              # triggers compensation
    except (carrier.Timeout, carrier.ServerError) as exc:
        raise TransientStepError(str(exc))              # retry the same step
    return {"shipment_id": shp.id}

Compensations need the same treatment, with their own keys ("{saga_id}:undo:{step}"), because a void or a cancellation can also be repeated after a crash. Finally, schedule a reconciliation job that compares the saga log with the external systems once a day — authorizations with no matching saga, reservations older than any live checkout — and reports mismatches. It catches the rare case where every other safeguard missed.

The one dangerous crash window A step has three phases: call the external API, record completion in the saga log, enqueue the next step. A crash before the call is safe because nothing happened. A crash after recording is safe because the retry sees the step as done. A crash between the call and the record means the retry repeats the call, which is only safe with an idempotency key. Crash points inside one step call external API record in saga log enqueue next step crash before: nothing happened crash between: call repeats crash after: step seen as done An idempotency key turns the middle case into a lookup of the first attempt's result.

Verification

Test the unhappy paths, not just the happy one:

@pytest.mark.parametrize("fail_at", ["reserve_stock", "authorize_payment", "create_shipment"])
def test_failure_before_pivot_leaves_nothing_behind(fail_at, fakes, run_saga):
    fakes[fail_at].fail_permanently()
    saga_id = run_saga(order_id="o-1")
    assert sagas.get(saga_id).state == "failed"
    assert fakes.stock.reservations_for("o-1") == []
    assert fakes.payments.open_authorizations_for("o-1") == []

def test_crash_during_compensation_is_resumed(fakes, run_saga, kill_worker_after):
    fakes.create_shipment.fail_permanently()
    kill_worker_after("void_authorization")
    saga_id = run_saga(order_id="o-2")
    sweep_stuck_sagas()
    assert sagas.get(saga_id).state == "failed"
    assert fakes.payments.void_calls("o-2") == 1       # idempotency key held

In production, graph sagas by state and alert on any saga in compensating for more than an hour.

Gotchas & Edge Cases

Compensation is not rollback. Other systems saw the intermediate state. A customer may have received a "payment authorized" notification; stock was invisible to others for a few seconds. Design user-facing messages for this, and prefer steps whose intermediate state is private.

Concurrent sagas on the same entity. Two checkouts for the last unit of stock both reserve and one fails later. The reservation step, not the saga, must enforce the invariant (a conditional decrement).

Compensation that cannot succeed. A shipment already picked up cannot be cancelled. Model that as a business process — a return label, a human task — and make the compensation job create it rather than fail forever.

Forgetting the pivot. If every step is compensable, a failure in the email step refunds a completed order. Decide explicitly where the point of no return is.

FAQ

Do I need a workflow engine for sagas? Not for four or five steps. A saga log, idempotent jobs, and a sweeper are enough. When sagas run for days, involve human approval, or number in the dozens of kinds, a workflow engine starts to pay for itself — see when to move from job queues to Temporal.

Orchestrated or choreographed sagas? Orchestrated — one coordinator enqueuing steps from a log — is far easier to debug. Choreographed sagas, where services react to each other's events, spread the logic across teams and make "what state is order 42 in?" hard to answer.

What if the compensation itself fails permanently? Park the saga in a needs_attention state and alert. It is the one case where a human must decide, and the saga log gives them everything needed to act.

Related