System Design
T7

Distributed transactions and sagas

An order touches stock, payment and shipping, and each service owns its own database. No single transaction can cover them. Pick a protocol, and make every step safe to retry.

Not startedSaved in this browser only.
  1. 1A transaction ends at one database. A change across services needs a protocol.
  2. 2Two-phase commit is atomic, but a prepared participant waits for the coordinator with its locks held.
  3. 3A saga commits each step locally and compensates finished steps on failure. Others can see the steps in between.
  4. 4Write the event in the same transaction as the data (outbox). Make every consumer idempotent (inbox).
T7
    A

    The dual write

    commit here, call there
    Order serviceowns ordersorders databasePostgresPayment serviceown databaseinsert order 42, commitcommitted: durable↯ process diescharge order 42: never sentno transaction covers both✕ Order 42 exists and is never charged. No error anywhere.
    • Each service commits in its own database.
    • A crash, a timeout or a deploy can fall between any two steps.
    • Calling payments first only reverses the problem: charged, with no order.

    Two writes in two systems are never atomic by themselves. A crash between them leaves one without the other.

    B

    Which protocol, where

    for a change across services
    methodwherestatus
    Commit, then call the next serviceServiceNot approved
    Two-phase commit, PREPARE TRANSACTIONPostgresFew databases, one team
    XA across databases and a brokerXABlocks on a crash
    Saga, orchestrated, with a saga logServiceApproved
    Saga, choreographed by eventsStream2 to 4 steps
    Transactional outbox and a relayPostgresApproved
    Idempotent consumer with an inboxPostgresApproved
    Try, confirm, cancel (TCC)ServicePartners can reserve
    Dedupe in Redis, effect in PostgresRedisNot approved

    The last row is a dual write again: a crash between the two stores loses or repeats the effect.

    I use a saga with an outbox for business flows, and keep two-phase commit for a few databases that one team runs.

    C

    Two-phase commit, saga, outbox, TCC

    what each one guarantees
    methodconsistencyothers see partial statelatency and costwhen something diesuse for
    Two-phase commitAtomicNo. Locks stay until phase 2.2 round trips and 2 durable writes per participant, plus the decision. 4,200 a second against 13,000 local pairs.Blocks Prepared participants wait for the coordinator.Moving rows between 2 shards; one team, short transactions.
    SagaEventual Done, or compensated.Yes, between steps. Add semantic locks.One local commit per step. No lock across services.Recovers Resume from the saga log.Business flows across services: orders, bookings, signups.
    OutboxData and event commit together.Consumers see the event later.One more insert per change: 31,000 to 17,000 orders a second.Duplicates The relay retries; consumers dedupe.Every change that other services must hear about.
    TCCEventual Confirmed, or cancelled.Reservations show as held, never as final.2 calls per participant.Expiry An unconfirmed try cancels itself.Partners that can reserve: cards, stock, seats.

    A saga and an outbox are used together: each saga step publishes its result through the outbox of its own service.

    Two-phase commit trades availability for atomicity. A saga trades isolation for availability, and it needs compensations.

    D

    Two-phase commit

    figure, then pseudo code
    phase 1: prepare, votephase 2: commitCoordinatorInventoryPaymentsin doubt: row locks heldin doubt: row locks heldPREPAREvotes: yes, yeslog commitdecision point↯ crash here?COMMIT PREPARED✓ done✓ doneA crash after the votes leaves both participants in doubt.They wait, with locks held, until the coordinator returns.
    • A participant that voted yes has promised to commit. It cannot change its mind or decide alone.
    • Postgres sets max_prepared_transactions to 0 by default, which turns PREPARE TRANSACTION off.
    • A forgotten prepared transaction holds its locks and stops VACUUM from removing dead rows. Alert on its age.
    coordinatorpseudo code
    commit(txid, work):                     // the coordinator
      FOR EACH participant p:               // phase 1: prepare, vote
        p: BEGIN; do the work; PREPARE TRANSACTION txid1
        IF p fails: log abort; ROLLBACK PREPARED on the rest; RETURN
      log "commit txid"2                     // the decision point
      FOR EACH participant p:               // phase 2: retried until done3
        p: COMMIT PREPARED txid
    
    recover():                              // after the coordinator restarts
      FOR EACH prepared transaction on each participant:
        IF the log says commit: COMMIT PREPARED
        ELSE: ROLLBACK PREPARED             // presumed abort4
    1. 1Writes the transaction to disk and detaches it from the session. It survives a disconnect and a restart, with its locks.
    2. 2The commit point of the whole transaction. Once this row is durable, every participant must commit.
    3. 3Phase 2 may not fail. A participant that is down receives COMMIT PREPARED when it returns.
    4. 4No decision in the log means nobody was told to commit, so rolling back is safe.
    Tested source Go: the coordinator · SQL: decision log and in-doubt list
    Go: the coordinatorgo
    // Commit runs work on every participant as one global transaction named txid.
    func (c *Coordinator) Commit(ctx context.Context, txid string, work map[string]Work) error {
      for i, p := range c.Participants { // phase 1: every participant prepares, or nobody commits
        if err := Prepare(ctx, p, txid, work[p.Name]); err != nil {
          return errors.Join(fmt.Errorf("%s voted no: %w", p.Name, err), c.abort(ctx, txid, c.Participants[:i]))
        }
        c.step(p.Name, "PREPARE TRANSACTION '"+p.GID(txid)+"'")
      }
      if c.CrashAt == "before decision" {
        return ErrCrashed // every participant waits; recovery finds no decision and rolls back
      }
      // The decision point. Once this row is durable, the transaction is committed.
      if _, err := c.Log.Exec(ctx, twopc["decide"], txid, "commit"); err != nil {
        return errors.Join(fmt.Errorf("log decision: %w", err), c.abort(ctx, txid, c.Participants))
      }
      c.step("coordinator", "log decision: commit")
      if c.CrashAt == "after decision" {
        return ErrCrashed // every participant waits, holding its locks
      }
      for _, p := range c.Participants { // phase 2: may not fail; Recover retries it
        if err := Finish(ctx, p.Conn, p.GID(txid), true); err != nil {
          return fmt.Errorf("commit on %s, retry with Recover: %w", p.Name, err)
        }
      }
      return nil
    }
    
    SQL: decision log and in-doubt listsql
    -- The coordinator's log. It must reach disk before any participant hears "commit".
    CREATE TABLE twopc_decisions (
      gid      text PRIMARY KEY,
      decision text NOT NULL CHECK (decision IN ('commit', 'abort'))
    );
    
    -- Prepared transactions in this database: they hold their locks until someone decides.
    SELECT gid FROM pg_prepared_xacts WHERE database = current_database() ORDER BY gid;

    In phase 1 every participant prepares and votes. The coordinator logs the decision, and phase 2 applies it. A crash in between blocks everyone.

    E

    See it: the coordinator dies

    recorded on a real Postgres server
    Coordinatorlog: no decisionrunningsessions to both databasesinventory · Postgresprepared: o42:inventoryrow lock held, waiting for a decisionothers read: stock reserved = 0payments · Postgresnothing in doubtno lock heldothers read: card spent = $0.00Another ordersame stock rownot startedlock_timeout 200 mslocks held byprepared txs: 3A prepared transaction belongs to no session. Only COMMIT PREPARED or ROLLBACK PREPARED ends it.
    step 1 of 9

    After this step

    in doubt
    o42:inventory
    locks
    3 entries in pg_locks with no session
    stock
    reserved = 0
    card
    spent = $0.00

    A prepared transaction outlives its session and even a restart. Until the coordinator returns, every order for that row waits.

    F

    The saga path

    click a step; its path lights up
    Customerplaces an orderOrder serviceorchestrator, N copiesOrders DBorders, saga logInventoryservice, own databasePaymentsservice, own databaseShippingservice, own database

    Step 1: Create order

    • Insert the order with status pending, in the order service database.
    • Pending is the semantic lock: other requests see that the order is in flight.
    • The order id is the idempotency key for every later step.

    If it fails

    The insert fails: return an error. Nothing else has happened yet.

    Each step is a local transaction in one service. The orchestrator logs every step, so a new copy can finish the saga after a crash.

    G

    Capabilities used

    what each tool gives you
    toolcapabilitywhat it gives this designalso used for
    PostgresLocal transactionsEach saga step is all or nothing inside its own service.Every multi-row change
    PostgresPREPARE TRANSACTION, COMMIT PREPAREDPhase 1 and phase 2 of two-phase commit. Any session can finish a prepared transaction.XA transaction managers
    Postgrespg_prepared_xactsLists the transactions in doubt, for recovery and alerts.Monitoring
    Postgresmax_prepared_transactionsDefault 0 Turns prepared transactions on. Changing it needs a restart.
    PostgresUnique key and ON CONFLICT DO NOTHINGSteps, compensations and the inbox become idempotent: a repeat finds the row.Idempotency keys, dedupe
    PostgresConditional UPDATE and its row countTake stock or spend within a limit in one statement. 0 rows is a business failure.Quotas, seats, balances
    PostgresForeign keyShipping refuses an unknown region, a failure the saga must compensate.Reference data
    PostgresFOR UPDATE SKIP LOCKEDSeveral relays drain one outbox. Each takes rows the others have not locked.Job queues
    PostgresPartial index WHERE published_at IS NULLThe relay finds unsent events without scanning sent ones.Queues, soft deletes
    RedisStreams: XADD, XREADGROUP, XACKThe broker. A consumer group shares the events, and each consumer acknowledges what it handled.Job queues, activity feeds
    RedisPending entries listA consumer that crashes before XACK receives the message again.Retry of failed jobs
    KafkaPartitioned log, consumer offsetsThe usual broker at scale. Delivery is at least once by default, and order holds per partition.Change data capture

    Postgres gives me local transactions, prepared transactions, unique keys for idempotency and SKIP LOCKED for the relay. A Redis stream gives me at-least-once delivery.

    H

    Orchestrate a saga

    pseudo code
    orchestratorpseudo code
    run(order):
      create order, status pending          // the semantic lock1
      FOR EACH step IN [reserve, charge, ship, confirm]:
        log step started
        result = step(order)                // one local transaction, keyed by order id2
        IF result is a timeout: retry the same step3
        IF result is a business failure:
          log step failed
          FOR EACH earlier step, newest first4:
            compensate(step)                // logged and idempotent
          set order cancelled; RETURN
        log step done
    
    resume(order):                          // a new orchestrator, after a crash
      read the saga log
      skip the steps marked done5; run the rest
    1. 1Pending tells other requests that the order is in flight. A cancel request on a pending order waits or is refused.
    2. 2Each service stores the order id under a unique key. A retried step finds it and changes nothing.
    3. 3A timeout says nothing about the outcome. Retry with the same key until the service answers.
    4. 4Undo in reverse order, as a stack. Only steps that finished have anything to undo.
    5. 5A step marked started may have committed. It runs again, which is safe because it is idempotent.
    Tested source Go: run and compensate
    Go: run and compensatego
    func (s *Saga) resume(ctx context.Context, o Order, logged map[string]string) error {
      for i, st := range steps {
        switch logged[st.name+"/forward"] {
        case "done":
          s.emit(Event{Actor: "orchestrator", Step: st.name, Phase: "recover", Do: "look up " + st.name + " in the log", Got: "done: skip it", Kind: "skip"})
          continue
        case "failed":
          return s.compensate(ctx, o, i, errors.New("failed before the crash"))
        }
        if err := s.log(ctx, o.ID, st.name, "forward", "started"); err != nil {
          return err
        }
        err := s.attempt(ctx, st, o)
        if permanent(err) {
          if err := s.log(ctx, o.ID, st.name, "forward", "failed"); err != nil {
            return err
          }
          return s.compensate(ctx, o, i, err)
        }
        if err != nil {
          return err
        }
        if s.CrashAfter == st.name {
          s.CrashAfter = ""
          s.emit(Event{Actor: "orchestrator", Step: st.name, Phase: "forward", Do: "crash before logging done", Got: "the step committed; the log says started", Kind: "crash"})
          return ErrCrashed
        }
        if err := s.log(ctx, o.ID, st.name, "forward", "done"); err != nil {
          return err
        }
      }
      return nil
    }
    
    // compensate undoes every step before failed, newest first, then cancels the order.
    func (s *Saga) compensate(ctx context.Context, o Order, failed int, cause error) error {
      for j := failed - 1; j >= 0; j-- {
        st := steps[j]
        if st.undo == nil {
          continue
        }
        if err := s.log(ctx, o.ID, st.name, "compensate", "started"); err != nil {
          return err
        }
        got, err := st.undo(ctx, s, o)
        if err != nil {
          return fmt.Errorf("compensate %s: %w", st.name, err) // retried later by Resume
        }
        s.emit(Event{Actor: st.actor, Step: st.name, Phase: "compensate", Do: undoName[st.name], Got: got, Kind: "comp"})
        if err := s.log(ctx, o.ID, st.name, "compensate", "done"); err != nil {
          return err
        }
      }
      if _, err := s.DB.Orders.Exec(ctx, q["finish_order"], o.ID, "cancelled"); err != nil {
        return fmt.Errorf("cancel order %d: %w", o.ID, err)
      }
      s.emit(Event{Actor: "orders", Phase: "compensate", Do: "cancel order", Got: "order is cancelled", Kind: "comp"})
      return fmt.Errorf("%w: %w", ErrCancelled, cause)
    }
    

    The orchestrator logs each step before and after it runs. On a business failure it compensates the finished steps, newest first.

    I

    Kinds of step

    order the steps by risk
    kindrulehere
    CompensatableHas an undo step.reserve stock, charge card
    PivotThe point of no return. If it commits, the saga goes forward.create shipment
    RetriableCannot fail for a business reason. Retry until done.confirm order
    • A compensation is a new transaction, not a rollback. A refund is visible on the card statement.
    • A compensation must succeed. Retry it, and alert a person after many failures.

    I put every step that can fail for a business reason before the pivot. After the pivot, steps only retry.

    J

    See it: pick the step that fails

    recorded against four Postgres databases

    Order 42: 1 lamp, $49.00, ship to AQ. Fault: create shipment: region AQ is not served.

    Order serviceorchestrator+ saga logorder: pendingrunningforward: one local transaction per stepcompensate: newest firstreserve stockinventorycompensatablenot runrelease stockif a later step failsreserved 0 of 5hold: nonecharge cardpaymentscompensatablenot runrefund cardif a later step failsspent $0.00charge: nonecreate shipmentshippingpivotnot runshipment: noneconfirm orderordersretriablenot runorder: pendingeach service's own database, read after this stepNo transaction spans two boxes. Between steps, other requests can see the partial order.
    step 1 of 7

    After this step

    order
    pending
    stock
    0 reserved, hold none
    card
    $0.00 spent, charge none
    shipment
    none

    When the pivot fails, the saga refunds the card and then releases the stock. When the orchestrator crashes, a new one resumes from the log.

    K

    Orchestration or choreography

    who knows the flow
    OrchestrationOrchestratorsaga log, the flowinventorypaymentsshippingcommands, repliesOne place knows every step.Easy to read, test and recover.The orchestrator is one more service.Choreographyevent streamordersinventorypaymentsshippingorder placed → stock reserved → card chargedEach service reacts to the last event.No central service, loose coupling.The flow exists only in the subscriptions.
    questionorchestrationchoreography
    Where is the flow?In one serviceSpread over subscriptions
    Who compensates?The orchestratorEach service, on a failure event
    Cyclic dependenciesNoneEasy to create

    I orchestrate when the flow has more than a few steps or needs compensations. Choreography suits short flows with few owners.

    L

    Idempotent steps

    pseudo code, then SQL
    a step and its compensationpseudo code
    reserve(order):                        // a forward step
      BEGIN
        insert hold for order id, skip if it exists1
        IF it existed: COMMIT; RETURN done  // a retry, or already released
        take qty WHERE reserved + qty <= on_hand2
        IF 0 rows: ROLLBACK; RETURN out of stock
      COMMIT
    
    release(order):                        // its compensation
      BEGIN
        insert hold as released, skip if it exists
        IF it was inserted: COMMIT; RETURN  // reserve never ran: block it3
        set hold released WHERE status = held
        IF a row changed: give the qty back
      COMMIT
    1. 1INSERT ... ON CONFLICT DO NOTHING on the order id. The row count says whether this call is the first.
    2. 2Check and take in one statement. Two sagas cannot both take the last lamp.
    3. 3The compensation leaves a released row. A delayed reserve for the same order then finds it and does nothing.
    Tested source SQL: reserve · SQL: release
    SQL: reservesql
    -- Forward step 1 (inventory). The order id is the idempotency key.
    INSERT INTO stock_holds (order_id, sku, qty) VALUES ($1, $2, $3)
    ON CONFLICT (order_id) DO NOTHING;
    
    -- Take stock only if enough is free. 0 rows: out of stock.
    UPDATE stock SET reserved = reserved + $2
    WHERE sku = $1 AND reserved + $2 <= on_hand;
    SQL: releasesql
    -- Compensation 1. If the reserve never ran, leave a 'released' row so a late reserve does nothing.
    INSERT INTO stock_holds (order_id, sku, qty, status) VALUES ($1, $2, $3, 'released')
    ON CONFLICT (order_id) DO NOTHING;
    
    UPDATE stock_holds SET status = 'released'
    WHERE order_id = $1 AND status = 'held'
    RETURNING sku, qty;
    
    UPDATE stock SET reserved = reserved - $2 WHERE sku = $1;

    Each step and each compensation claims the order id first. A repeat, or a compensation that came too early, then changes nothing.

    M

    Transactional outbox

    pseudo code, then SQL
    place, relay, handlepseudo code
    place(order):                          // the order service
      BEGIN
        insert the order
        insert the event INTO outbox        // same transaction1
      COMMIT
    
    relay():                               // a loop; run 1 or more copies
      BEGIN
        batch = oldest unsent events, FOR UPDATE SKIP LOCKED2
        publish the batch to the stream
        mark the batch sent
      COMMIT                               // crash before this: sent again3
    
    handle(message):                       // every consumer
      BEGIN
        insert message id INTO inbox4, skip if it exists
        IF it existed: COMMIT; RETURN       // a duplicate
        do the effect
      COMMIT
      acknowledge the message5
    1. 1The order and its event commit together. No crash can leave one without the other.
    2. 2Several relays can run. Each skips rows another relay holds, so no row is sent by two relays at once.
    3. 3The rows stay unsent and go out again. Delivery is at least once, never at most once.
    4. 4The inbox row and the effect commit together, so a message id takes effect once.
    5. 5Acknowledge after the commit. A crash before it means a redelivery, which the inbox absorbs.
    Tested source Go: place the order · Go: the relay · Go: the consumer · SQL: outbox and inbox
    Go: place the ordergo
    // PlaceOrder writes the order and its event in one local transaction. Both exist or neither does.
    // The service never talks to the broker; the relay does.
    func PlaceOrder(ctx context.Context, db *pgxpool.Pool, o Order) error {
      payload, err := eventFor(o)
      if err != nil {
        return err
      }
      return pgx.BeginFunc(ctx, db, func(tx pgx.Tx) error {
        if _, err := tx.Exec(ctx, ob["place_order"], o.ID, o.SKU, o.Qty, o.CardID, o.AmountCents, o.Region); err != nil {
          return fmt.Errorf("place order %d: %w", o.ID, err)
        }
        if _, err := tx.Exec(ctx, ob["write_event"], "orders", payload); err != nil {
          return fmt.Errorf("write event for order %d: %w", o.ID, err)
        }
        return nil
      })
    }
    
    Go: the relaygo
    // PublishBatch publishes up to n of the oldest unpublished events and returns how many it sent.
    // Delivery is at least once: a crash before the commit sends the same events again later.
    func (r *Relay) PublishBatch(ctx context.Context, n int) (sent int, err error) {
      err = pgx.BeginFunc(ctx, r.DB, func(tx pgx.Tx) error {
        rows, err := tx.Query(ctx, ob["next_batch"], n)
        if err != nil {
          return fmt.Errorf("read outbox: %w", err)
        }
        var id int64
        var payload string
        var ids []int64
        pipe := r.RDB.Pipeline()
        if _, err := pgx.ForEachRow(rows, []any{&id, &payload}, func() error {
          ids = append(ids, id)
          pipe.XAdd(ctx, &redis.XAddArgs{Stream: r.Stream, Values: []any{"msg_id", id, "payload", payload}})
          return nil
        }); err != nil {
          return fmt.Errorf("read outbox: %w", err)
        }
        if len(ids) == 0 {
          return nil
        }
        if _, err := pipe.Exec(ctx); err != nil {
          return fmt.Errorf("publish %d events: %w", len(ids), err)
        }
        if r.CrashAfterPublish {
          r.CrashAfterPublish = false
          return ErrRelayCrashed // the transaction rolls back: the rows stay unpublished
        }
        if _, err := tx.Exec(ctx, ob["mark_published"], ids); err != nil {
          return fmt.Errorf("mark published: %w", err)
        }
        sent = len(ids)
        return nil
      })
      return sent, err
    }
    
    Go: the consumergo
    // handle applies one message. The inbox row and the effect commit together, so a message id is
    // applied at most once, however many times the broker delivers it.
    func (c *Consumer) handle(ctx context.Context, m redis.XMessage) (applied bool, err error) {
      msgID, err := strconv.ParseInt(fmt.Sprint(m.Values["msg_id"]), 10, 64)
      if err != nil {
        return false, fmt.Errorf("message %s: bad msg_id: %w", m.ID, err)
      }
      var ev orderPlaced
      if err := json.Unmarshal([]byte(fmt.Sprint(m.Values["payload"])), &ev); err != nil {
        return false, fmt.Errorf("message %s: bad payload: %w", m.ID, err)
      }
      err = pgx.BeginFunc(ctx, c.DB, func(tx pgx.Tx) error {
        if c.Inbox {
          tag, err := tx.Exec(ctx, ob["inbox_claim"], msgID)
          if err != nil {
            return fmt.Errorf("claim message %d: %w", msgID, err)
          }
          if tag.RowsAffected() == 0 {
            return nil // seen before: skip the effect
          }
        }
        if _, err := tx.Exec(ctx, ob["send_email"], ev.OrderID); err != nil {
          return fmt.Errorf("send email for order %d: %w", ev.OrderID, err)
        }
        applied = true
        return nil
      })
      return applied, err
    }
    
    SQL: outbox and inboxsql
    -- The relay takes the oldest unpublished events. SKIP LOCKED lets several relays share the work.
    SELECT id, payload::text FROM outbox
    WHERE published_at IS NULL
    ORDER BY id
    LIMIT $1
    FOR UPDATE SKIP LOCKED;
    
    UPDATE outbox SET published_at = now() WHERE id = ANY($1);
    
    -- The consumer records the message id in the same transaction as the effect.
    INSERT INTO inbox (message_id) VALUES ($1) ON CONFLICT (message_id) DO NOTHING;
    
    -- The effect. Without the inbox, a redelivery runs it twice: sent becomes 2.
    INSERT INTO order_emails (order_id) VALUES ($1)
    ON CONFLICT (order_id) DO UPDATE SET sent = order_emails.sent + 1;
    the outbox tablesql
    -- The outbox: an event written in the same transaction as the order.
    CREATE TABLE outbox (
      id           bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
      topic        text        NOT NULL,
      payload      jsonb       NOT NULL,
      published_at timestamptz1
    );
    CREATE INDEX outbox_unpublished ON outbox (id) WHERE published_at IS NULL2;
    1. 1NULL until the relay sends the event.
    2. 2A partial index: it holds only unsent rows, so it stays small.
    relaystatuswhy
    Poll the table, batch of 100ApprovedSimple. 39,000 events a second here.
    Read the WAL (change data capture)ApprovedNo polling. Needs a connector such as Debezium.
    Publish inside the transactionNot approvedA rollback after the publish sends an event for nothing.

    The service writes the event into an outbox table in the same transaction. A relay publishes it, and each consumer dedupes by message id in an inbox.

    N

    See it: crash between commit and publish

    recorded from Postgres and a Redis stream
    Order serviceone transactionorders databaseorders: 1outbox, unsent: 1both commit togetherRelaySKIP LOCKEDRedis streamentries: 0at least onceemail serviceinbox rows: 0emails sent: 0dedupe by message id
    step 1 of 5

    After this step

    orders
    1
    outbox
    1 unsent
    stream
    0 entries
    inbox
    0 rows
    emails
    0 sent

    The dual write lost an event. The outbox lost nothing but sent duplicates, and the inbox made the duplicates harmless.

    O

    Try, confirm, cancel

    reservations with an expiry
    calldoesfor this order
    TryReserve, with an expiry. Nothing is final.Hold 1 lamp. Authorize $49.00.
    ConfirmMake each reservation final. Retried until it succeeds.Ship the lamp. Capture $49.00.
    CancelDrop each reservation. The expiry also cancels.Release the hold. Void the authorization.
    • Card networks already work this way: authorize, then capture or void.
    • Each participant must offer all three calls, each idempotent.

    TCC is a saga whose first step only reserves. Confirm makes the reservation final, and cancel or the expiry drops it.

    P

    Sagas have no isolation

    anomalies and countermeasures
    anomalyexamplecountermeasure
    Dirty readA report counts stock that a saga later gives back.Semantic lock: read pending as not final.
    Lost updateTwo sagas each set reserved to a value they computed.Commutative update: reserved + qty.
    Fuzzy readA step reads the price, and it changes before the charge.Reread with a version check before the write.
    Early cancelThe customer cancels while the saga runs.Refuse or queue it while the order is pending.

    Other requests see a saga's steps before it ends. I use semantic locks and commutative updates to keep that safe.

    Q

    Failure cases

    what breaks, and what keeps the data right
    eventresultwhy it is safesaved by
    The coordinator dies after the votes, before the decisionBoth participants wait with locks held.Recovery finds no decision and rolls both back.Decision log
    The coordinator dies after it logs commitSame wait, even across a restart.Recovery reads commit and runs COMMIT PREPARED on each.Prepared tx
    A prepared transaction is forgottenIts locks stay. VACUUM cannot clean behind it.Not safe by itself. Alert on the age of rows in pg_prepared_xacts.Monitoring
    The orchestrator dies after a step commitsThe log says started, not done.A new orchestrator runs the step again. The unique key makes it a no-op.Unique key
    A step times outThe outcome is unknown.Retry with the same order id until the service answers.Unique key
    A compensation arrives before its stepThe step has not run yet.The compensation writes a marker. The late step finds it and does nothing.Marker row
    A compensation keeps failingThe order stays half undone.Retry from the log. After many tries, alert a person with the saga log.Saga log
    The service dies between commit and publishWith an outbox, the event waits in the table.The relay sends it later.Outbox
    The relay dies after it publishesThe same events go out again.The consumer's inbox skips message ids it has seen.Inbox
    A consumer dies before XACKThe stream delivers the message again.Same inbox check, in the same transaction as the effect.Pending list
    Two relays publish events of one order out of orderA consumer sees "shipped" before "paid".Partition by order id, so one relay or partition owns each order.Partition key
    R

    Scale ladder

    start simple; climb only on a signal
    Each step adds one component1One database2+ outbox, relay3+ orchestrator4+ relays or CDC5+ partitionsmore load →
    Measured rates on one laptop1001k10k100k1MOrder insert, alone: 31,009 per secondOrder insert, alone31,009Order and outbox row: 17,378 per secondOrder and outbox row17,378Relay, 4 copies, batch 100: 66,733 per secondRelay, 4 copies, batch 10066,733Relay, batch 1: 1,151 per secondRelay, batch 11,151Two local commits: 13,091 per secondTwo local commits13,091Two-phase commit, 2 DBs: 4,151 per secondTwo-phase commit, 2 DBs4,151Saga, 4 steps, 4 DBs: 1,759 per secondSaga, 4 steps, 4 DBs1,759per second, log scale
    stepaddit handlesmove up when you see
    1One database, one transaction. Orders, stock and payments in one Postgres, as modules.About 31,000 single-row inserts a second here. No protocol needed.Separate teams must own and deploy their own data.
    2Outbox, relay and inbox. One relay polls in batches into a stream.About 17,000 orders a second with the event row; 39,000 events a second from one relay.A business flow crosses services and needs an undo.
    3An orchestrated saga with a saga log. Run N stateless copies.About 1,800 sagas a second on a laptop: 13 commits each.The relay falls behind: events wait in the outbox.
    4More relays with SKIP LOCKED, or change data capture from the WAL.About 67,000 events a second from 4 relays.One database or one stream is at its write limit.
    5Partitions: saga logs and stream partitions by order id.Each partition adds its own rate. Order holds per order.Top of the ladder.

    Two-phase commit is not a step on this ladder. Keep it for a few databases that one team runs, such as moving rows between 2 shards.

    If the data can live in one database, I use one transaction. I add an outbox when other services must hear about changes, and a saga when a flow crosses services.

    S

    Drill

    predict, then reveal

    0 of 10 known

    1. The order service writes to its database and then calls payments. Why can it not wrap both in one transaction?

    2. Both participants voted yes, then the coordinator died. What can a participant decide on its own?

    3. What does PREPARE TRANSACTION give that an open transaction does not?

    4. The card charge fails. Which compensations run?

    5. The orchestrator crashed after the charge committed but before it logged the step. Why is the card not charged twice on resume?

    6. A compensation arrives before its forward step, because the forward request was delayed. What stops the late step?

    7. What is the pivot step, and where do you put it?

    8. Why write the event to an outbox table instead of publishing after the commit?

    9. The relay published a batch, then crashed before it marked the rows. What happens?

    10. A saga holds stock for a pending order. Another request counts available stock. What can go wrong?

    T

    Numbers to say

    measured in the lab
    2PC
    About 4,200 global transactions a second over 2 databases, against 13,000 pairs of local commits.
    saga
    About 1,800 sagas a second: 4 steps, 4 databases, 13 commits each.
    outbox
    One more insert per order: 31,000 became 17,000 orders a second.
    relay
    Batch 1: 1,200 events a second. Batch 100: 39,000. 4 relays: 67,000.
    default
    max_prepared_transactions is 0 in Postgres, so PREPARE TRANSACTION is off.

    Postgres 16 and Redis 8 on an 8-core laptop, 32 clients (16 for 2PC), 3 s per case, median of 3 runs. Other workloads shared the machine, so use these as orders of magnitude.