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.
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.
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.
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
- Job Chaining & Workflow Orchestration — where sagas fit among chains and workflow records.
- Tracking Progress of Multi-Step Jobs — expose saga state to users and support.
- Preventing Duplicate Job Execution with Idempotency — the idempotency keys every step relies on.
- Reliable Webhook Delivery with a Job Queue — retry-until-success steps after the pivot.