# Running an Outbox Dispatcher Safely with Multiple Workers: Why FOR UPDATE SKIP LOCKED Isn't Enough, and Leases with Fencing Tokens

> When several workers process a PostgreSQL outbox table, doing SELECT ... FOR UPDATE SKIP LOCKED, committing, and then calling an external API makes two workers process the same message. This explains why, and how to claim and finalize correctly with a lease and a fencing token (lease_version) — plus backoff, dead-lettering, per-aggregate ordering and monitoring SQL — with Python code tested concurrently against PostgreSQL 18.

- Published: 2026-10-04
- Author: 友田 陽大
- Tags: PostgreSQL, 信頼性, アーキテクチャ設計, Python, データベース
- URL: https://tomodahinata.com/en/blog/outbox-dispatcher-skip-locked-lease-fencing-token-guide
- Category: Reliability, async & real-time
- Pillar guide: https://tomodahinata.com/en/blog/transactional-outbox-pattern-reliable-event-publishing-guide

## Key points

- FOR UPDATE row locks are released at COMMIT. 'SELECT ... FOR UPDATE SKIP LOCKED → COMMIT → external API' lets another worker take the same row while the API call is running, and the message is processed twice
- SKIP LOCKED only guarantees that two transactions don't grab the same row at the same instant. Ownership that survives COMMIT must be persisted in the table as status, lease_until and lease_version
- Claim increments lease_version, and finalize runs 'WHERE id = ? AND lease_version = ?'. When a lease expires and another worker takes over, the old worker's finalize updates zero rows and is rejected (a fencing token)
- Fencing protects only the database's state. If the old worker had already called the external API, that side effect can't be undone, so downstream idempotency (de-duplication by message ID) remains mandatory
- Classify failures: transient failures back off with jitter, business defers are postponed without consuming attempts, and permanent failures or exhausted retries go to dead for a human to review

---

**The short answer.** When several workers process an outbox, taking rows with `SELECT ... FOR UPDATE SKIP LOCKED`, **committing, and then calling an external API lets two workers process the same message.** Row locks are released at COMMIT, so while you're calling the API the row belongs to nobody. The correct design uses a short transaction to set the row to `processing` and **writes the owner, an expiry (the lease) and a generation number (the fencing token) to the table before committing.** At finalize time you write only if the generation still matches, which rejects writes from an old worker whose lease expired and was taken over. Even then, duplicate external side effects don't drop to zero, so **downstream idempotency remains mandatory.**

This article picks up after the [transactional outbox](/blog/transactional-outbox-pattern-reliable-event-publishing-guide) — where the business update and the event are recorded in the same transaction — and covers how to **parallelize the delivery side (the dispatcher, or relay) in production.** The code was **verified on PostgreSQL 18.6 / psycopg 3.3 with 11 tests that reproduce four concurrent workers, a worker crash, and an expired lease** (October 2026).

> **Terminology**: in this article, the process that takes outbox rows and sends them out (to a message broker or an external API) is the **dispatcher**. It plays the same role as microservices.io's polling publisher, often called a relay.

---

## 1. A broken dispatcher: calling the external API after COMMIT

This implementation is common:

```python
# ❌ Broken: external I/O after COMMIT has dropped the locks. Do not copy this
def dispatch_broken(conn: psycopg.Connection, publisher: Publisher, *, batch: int = 50) -> list[int]:
    with conn.transaction():
        rows = conn.execute(
            """
            SELECT id, aggregate_id, event_type, payload FROM outbox
             WHERE status = 'pending' ORDER BY id LIMIT %s
               FOR UPDATE SKIP LOCKED
            """,
            (batch,),
        ).fetchall()
    # ← COMMIT has happened. The row locks are gone, and another worker can SELECT the same rows
    for id_, agg, typ, payload in rows:
        publisher.publish(message_id=str(id_), key=agg, event_type=typ, payload=payload)
        with conn.transaction():
            conn.execute("UPDATE outbox SET status = 'done', done_at = now() WHERE id = %s", (id_,))
    return [r[0] for r in rows]
```

Because it uses `FOR UPDATE SKIP LOCKED`, it looks safe to run on several workers. But as the PostgreSQL documentation says, **row-level locks are released at transaction end.**

```text
time  Worker A                                   Worker B
t1    BEGIN; SELECT ... FOR UPDATE SKIP LOCKED → locks id=1..3
t2    COMMIT (locks released; status still pending)
t3    publish(id=1) … (waiting on the external API)
t4                                               BEGIN; SELECT ... SKIP LOCKED → gets id=1..3
t5                                               COMMIT; publish(id=1)
t6    UPDATE status='done'                       UPDATE status='done'
      → id=1..3 are each published twice
```

In my tests, pausing worker A in the middle of publishing and running worker B meanwhile let B take every row A had, and every message was published twice.

### "Keep the transaction open while publishing" doesn't solve it either

If instead you COMMIT after publishing, the row lock persists during publish (the polling relay in [the existing outbox article](/blog/transactional-outbox-pattern-reliable-event-publishing-guide) has this shape). That prevents double-taking the row, but other problems appear:

- **The transaction and DB connection are held for as long as the external API takes.** When the broker slows down, connections and row locks pile up.
- **If publish succeeds but COMMIT fails** (connection loss, failover), the row goes back to `pending` and the next worker sends it again. That alone is acceptable as at-least-once, but the longer the transaction, the wider that window.
- **Long-running transactions hold back VACUUM** and bloat the table.

At small scale with fast publishing, this shape is workable. If you add workers, the external API is slow, or processing takes more than a few seconds, switch to the lease approach in the next section.

---

## 2. What SKIP LOCKED is really for

The PostgreSQL documentation describes `SKIP LOCKED` like this (paraphrased):

> With `SKIP LOCKED`, any selected rows that cannot be immediately locked are skipped. Skipping locked rows provides an inconsistent view of the data, so this is **not suitable for general purpose work**, but can be used to **avoid lock contention with multiple consumers accessing a queue-like table.**

So all `SKIP LOCKED` guarantees is that **two transactions don't grab the same row at the same instant.** It says nothing about what happens after COMMIT.

> **Row lock ≠ durable ownership**: a row lock lives only as long as its transaction. If you want the fact "worker A is processing this row" to survive COMMIT, that fact has to be **persisted as columns in the table.**

The same applies to `pg_advisory_xact_lock`. Transaction-level advisory locks also vanish at COMMIT, so they can't represent ownership after COMMIT ([PostgreSQL advisory locks vs FOR UPDATE](/blog/postgresql-advisory-lock-pg-advisory-xact-lock-race-condition-guide)).

---

## 3. The correct design: claim → COMMIT → external API → fenced finalize

```text
   pending ──claim (short tx, lease_version+1)──▶ processing
      ▲                                              │
      │  transient failure: back off and return      │ external API (outside any transaction)
      │  business defer: refund the attempt, delay   ▼
      └──────────────── fenced update ◀──────── classify the result
                          │      │
                          ▼      ▼
                        done    dead (permanent failure / retries exhausted → human review)

   rows still processing after lease_until → re-claimed by the next claim (lease_version+1)
```

### 3.1 Schema

```sql
CREATE TABLE outbox (
    id             bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
    aggregate_type text        NOT NULL,
    aggregate_id   text        NOT NULL,
    event_type     text        NOT NULL,
    payload        jsonb       NOT NULL,
    created_at     timestamptz NOT NULL DEFAULT now(),

    -- Columns below are for the dispatcher
    status         text        NOT NULL DEFAULT 'pending'
                   CHECK (status IN ('pending', 'processing', 'done', 'dead')),
    available_at   timestamptz NOT NULL DEFAULT now(),  -- earliest time it may be taken (backoff)
    attempts       integer     NOT NULL DEFAULT 0,      -- times claimed (a crash counts as one)
    lease_owner    text,                                -- for observability; correctness rests on lease_version
    lease_until    timestamptz,                         -- processing rows past this time may be re-claimed
    lease_version  bigint      NOT NULL DEFAULT 0,      -- fencing token; +1 on every claim
    last_error     text,
    done_at        timestamptz
);

CREATE INDEX outbox_ready  ON outbox (available_at, id) WHERE status = 'pending';
CREATE INDEX outbox_leased ON outbox (lease_until)      WHERE status = 'processing';
```

This is the schema from the existing outbox article (where `published_at IS NULL` means unpublished) with columns added for claiming.

- **`lease_version` is the fencing token.** It increases on every claim and is checked at finalize. `lease_owner` is an observability column for seeing who holds a row; correctness rests on `lease_version` alone. Judging by owner would fail to distinguish old and new work when a worker restarts under the same name.
- **`attempts` is incremented at claim time**, so that a "poison message" that crashes the worker as soon as processing starts still counts as an attempt and isn't re-claimed forever.
- **All times are computed with the database's `now()`**, so clock skew between workers can't distort lease decisions.

### 3.2 Claim: write ownership in a short transaction

```python
from dataclasses import dataclass
from datetime import timedelta
from typing import Any

import psycopg
from psycopg.rows import class_row


@dataclass(frozen=True)
class Claimed:
    id: int
    lease_version: int
    attempts: int
    aggregate_id: str
    event_type: str
    payload: dict[str, Any]


CLAIM_SQL = """
WITH picked AS (
    SELECT id
      FROM outbox
     WHERE (status = 'pending'    AND available_at <= now())
        OR (status = 'processing' AND lease_until  <  now()     -- left behind by a dead worker
            AND attempts < %(max_attempts)s)
     ORDER BY id
     LIMIT %(batch)s
       FOR UPDATE SKIP LOCKED
)
UPDATE outbox AS o
   SET status        = 'processing',
       lease_owner   = %(owner)s,
       lease_until   = now() + %(lease)s,
       lease_version = o.lease_version + 1,
       attempts      = o.attempts + 1
  FROM picked
 WHERE o.id = picked.id
RETURNING o.id, o.lease_version, o.attempts, o.aggregate_id, o.event_type, o.payload
"""


# A poison message that kills the worker on every claim never reaches finalize. Rows whose lease expired with attempts exhausted are sent to dead
REAP_SQL = """
UPDATE outbox
   SET status = 'dead', last_error = 'lease expired after max attempts (worker crash?)',
       lease_owner = NULL, lease_until = NULL
 WHERE status = 'processing' AND lease_until < now() AND attempts >= %(max_attempts)s
"""


def claim(
    conn: psycopg.Connection, *, owner: str, batch: int, lease: timedelta, max_attempts: int = 12
) -> list[Claimed]:
    params = {"batch": batch, "owner": owner, "lease": lease, "max_attempts": max_attempts}
    with conn.transaction(), conn.cursor(row_factory=class_row(Claimed)) as cur:
        cur.execute(REAP_SQL, params)  # send rows that exhausted attempts and lost their lease to dead
        cur.execute(CLAIM_SQL, params)
        return sorted(cur.fetchall(), key=lambda c: c.id)
```

What this single statement guarantees:

1. **Workers claiming at the same time get different rows.** This is where `SKIP LOCKED` does its job. In my tests, four workers repeatedly claiming 7 rows each from 200 got exactly 200 distinct IDs.
2. **After COMMIT, "processing, owner, expiry, generation" remain in the table.** Other workers only claim `pending` rows or `processing` rows whose lease has expired, so they never touch a row under lease.
3. **A dead worker's rows are re-claimed automatically once the lease expires**, and `lease_version` goes up each time.
4. **A message that keeps killing the worker stops at the attempt limit.** A poison message never reaches any exception branch, so the claim itself sends "rows whose lease expired with attempts exhausted" to `dead` (`REAP_SQL`). In my tests, claiming without finalizing (crashing every time) three times left the row `dead` with `attempts=3`, and it was never taken again.

The order of rows from `RETURNING` isn't guaranteed, so they are sorted by `id`.

### 3.3 Call the external API outside the transaction

The claim transaction commits immediately. The external API is called afterwards, **with no transaction open.** No connection or lock is held, so however slow the API is, the database is unaffected.

Make the lease **comfortably longer than the external API's timeout.** If publish times out after 10 seconds, use a 60-second lease. For work that can outlast the lease you'd extend it mid-flight (a heartbeat), but for short work such as publishing an outbox message, it's simpler to set the timeout shorter than the lease.

### 3.4 Fenced finalize: only the worker with the matching generation may finalize

```python
def _fenced_update(conn: psycopg.Connection, msg: Claimed, set_clause: str, params: dict[str, Any]) -> bool:
    # set_clause only ever receives constants from this module (never concatenate external input)
    with conn.transaction():
        cur = conn.execute(
            f"""
            UPDATE outbox
               SET {set_clause},
                   lease_owner = NULL,
                   lease_until = NULL
             WHERE id = %(id)s
               AND status = 'processing'
               AND lease_version = %(version)s
            """,
            {"id": msg.id, "version": msg.lease_version, **params},
        )
        return cur.rowcount == 1


def mark_done(conn: psycopg.Connection, msg: Claimed) -> bool:
    """False = my lease expired and someone else took over."""
    return _fenced_update(conn, msg, "status = 'done', done_at = now(), last_error = NULL", {})
```

`WHERE lease_version = %(version)s` is the fencing. It matters in this timeline:

```text
time  Worker A                                  Worker B
t1    claim → id=1, lease_version=1, 60 s lease
t2    publish(id=1) … stalls 90 s (GC pause, network delay)
t3    (lease expires)
t4                                              claim → id=1, lease_version=2
t5                                              publish(id=1) → success
t6                                              mark_done(version=2) → 1 row updated
t7    publish finishes → mark_done(version=1) → 0 rows updated (rejected)
```

In my tests, claiming with a 200 ms lease as worker A and leaving it alone meant worker B's claim returned nothing during the lease; after expiry B claimed the same ID with `lease_version` +1 and `attempts` 2, A's `mark_done` returned `False`, and B's returned `True`.

**If the lease has expired but nobody has taken over, the original worker's finalize still succeeds** (because `lease_until` isn't in the WHERE clause). There's no reason to throw away successful work and resend just because time ran out. What decides correctness isn't the clock but whether the generation changed.

There is one exception: the last attempt. If the claim that took the row used up `max_attempts` and publish then outlives the lease, the next `claim()` by any worker sends the row to `dead` through `REAP_SQL`, even though nobody took it over. The original worker's `mark_done` then finds the row no longer `processing`, so it returns `False` and logs `lost lease`. The message was delivered, but the row is left in `dead` for a human to review. Keep the publish timeout well below the lease so that this doesn't happen on the final attempt.

### 3.5 What fencing does not protect

In the timeline above, **message id=1 was published twice** (A at t2 and B at t5). All the fencing token prevented was A's stale result overwriting the database's state. The message A had already sent can't be recalled.

As Martin Kleppmann points out in his discussion of distributed locking, for a fencing token to protect external side effects, **the system being written to must itself check the token and reject stale ones.** Message brokers and most external APIs don't, so with an outbox you combine two things instead:

- **Pass the message ID (`outbox.id`) downstream as the de-duplication key.** Consumers use it to reject duplicates ([receiving events idempotently with a transactional inbox](/blog/webhook-idempotency-transactional-inbox-postgresql-guide)).
- **When calling an external API directly, attach an idempotency key derived from `outbox.id`.** The receiver collapses retries with the same key into one effect.

**Leases and fencing reduce how often double processing happens and keep the database's state intact; they don't make double processing impossible.** That distinction is the most commonly misunderstood point in dispatcher design.

---

## 4. Classify failures

Retrying every external failure the same way lets unfixable failures clog the queue and makes quickly recoverable ones wait longer than necessary. Sort outcomes into four groups:

| Class | Examples | Handling | Attempts |
| --- | --- | --- | --- |
| Success / terminal no-op | Sent / the receiver is already in the desired state | `done` | — |
| Transient failure (technical retry) | Timeout, 5xx, 429, connection loss | Back to `pending` with jittered backoff; `dead` after the limit | Consumed |
| Business defer | A resource it depends on (a customer account, say) doesn't exist yet | Back to `pending` after a delay | **Not consumed** |
| Permanent failure | 4xx validation error, nonexistent destination | `dead` (a human reviews it) | — |

```python
import random


class RetryableError(Exception):
    """Transient failure (timeout, 5xx, 429). Back off and retry."""


class PermanentError(Exception):
    """Failure retries won't fix (4xx validation and the like). Mark dead for a human."""


class Defer(Exception):
    """Can't run yet for business reasons. Postpone without consuming an attempt."""

    def __init__(self, delay: timedelta, reason: str) -> None:
        super().__init__(reason)
        self.delay = delay


def reschedule(conn: psycopg.Connection, msg: Claimed, *, delay: timedelta, error: str, refund_attempt: bool) -> bool:
    return _fenced_update(
        conn,
        msg,
        "status = 'pending', available_at = now() + %(delay)s, last_error = %(error)s,"
        " attempts = attempts - %(refund)s",
        {"delay": delay, "error": error[:500], "refund": 1 if refund_attempt else 0},
    )


def mark_dead(conn: psycopg.Connection, msg: Claimed, *, error: str) -> bool:
    return _fenced_update(conn, msg, "status = 'dead', last_error = %(error)s", {"error": error[:500]})


def full_jitter_backoff(attempt: int, *, base: float = 1.0, cap: float = 300.0) -> timedelta:
    """'Full Jitter' from the AWS Architecture Blog: sleep = random(0, min(cap, base * 2^attempt))"""
    return timedelta(seconds=random.uniform(0, min(cap, base * 2**attempt)))
```

**Why jitter**: when an outage ends, if every backed-up message retries on the same schedule, the freshly recovered service takes the whole load at once. "Full Jitter" from the AWS Architecture Blog spreads retries out by making the wait a uniform random value between zero and the cap. Retry strategy in general is covered in [the guide to retries, exponential backoff with jitter, and circuit breakers](/blog/retry-backoff-circuit-breaker-resilience-patterns-guide).

**Why a business defer refunds the attempt**: claim increments `attempts`, so if each defer consumed one, a slow-to-appear dependency alone would send the message to `dead`. A defer isn't a failure, so the attempt is refunded. If a defer could go on forever, though, cap the time since `created_at` and mark it `dead` beyond that.

### 4.1 The dispatch loop

```python
import logging

log = logging.getLogger("outbox.dispatcher")


def dispatch_once(
    conn: psycopg.Connection,
    publisher: Publisher,
    *,
    owner: str,
    batch: int = 50,
    lease: timedelta = timedelta(seconds=60),
    max_attempts: int = 12,
) -> int:
    """Process one batch. External calls happen outside any transaction."""
    messages = claim(conn, owner=owner, batch=batch, lease=lease, max_attempts=max_attempts)
    for msg in messages:
        try:
            publisher.publish(
                message_id=str(msg.id),  # downstream de-dup key; unchanged across retries
                key=msg.aggregate_id,
                event_type=msg.event_type,
                payload=msg.payload,
            )
        except Defer as d:
            fenced = reschedule(conn, msg, delay=d.delay, error=f"deferred: {d}", refund_attempt=True)
        except RetryableError as exc:
            if msg.attempts >= max_attempts:
                fenced = mark_dead(conn, msg, error=f"retries exhausted: {exc!r}")
            else:
                fenced = reschedule(
                    conn, msg, delay=full_jitter_backoff(msg.attempts), error=repr(exc), refund_attempt=False
                )
        except PermanentError as exc:
            fenced = mark_dead(conn, msg, error=repr(exc))
        else:
            fenced = mark_done(conn, msg)

        if not fenced:
            # Publish may have happened. The worker that took over will publish too; downstream absorbs the duplicate
            log.warning("lost lease", extra={"outbox_id": msg.id, "lease_version": msg.lease_version})
    return len(messages)
```

`Publisher` is a thin boundary around the external broker or API. Translating SDK exceptions into `RetryableError` / `PermanentError` / `Defer` is that boundary's job; the loop knows nothing about SDK exceptions. **Exceptions that can't be classified are deliberately not caught.** If an unexpected exception stops the loop, the same rows are re-claimed after the process restarts and the leases expire, and once `attempts` hits the limit they go to `dead`. Surfacing an unexpected bug is safer than silently retrying it forever as a "transient failure".

Messages near the end of a batch can see their lease expire before processing starts if earlier publishes are slow. Size the batch so that **`batch × worst-case time per message` is comfortably shorter than the lease.**

---

## 5. When ordering within an aggregate matters

With parallel workers, `OrderCreated` (id=1) and `OrderPaid` (id=2) for the same order can be taken by different workers, and id=2 can be published first. `ORDER BY id` decides the order rows are taken, not the order publishing completes.

If you need ordering, **only make the head message of each aggregate eligible for claiming.**

> **Precondition**: this assumes that within an aggregate, `id` order matches commit order. `id` is assigned at INSERT time, so if two transactions write to the same aggregate concurrently, T2 with id=11 can commit before T1 with id=10, and id=11 can be taken as the "head" until T1 commits. If producers **lock the aggregate row with `FOR UPDATE` before inserting into the outbox** (the `with_for_update=True` example in the [transactional outbox article](/blog/transactional-outbox-pattern-reliable-event-publishing-guide)), writes to the same aggregate are serialized until commit and the precondition holds.

```sql
WITH picked AS (
    SELECT o.id
      FROM outbox AS o
     WHERE ((o.status = 'pending' AND o.available_at <= now())
         OR (o.status = 'processing' AND o.lease_until < now() AND o.attempts < %(max_attempts)s))
       AND NOT EXISTS (                       -- wait if the same aggregate has an unfinished earlier message
            SELECT 1 FROM outbox AS prev
             WHERE prev.aggregate_id = o.aggregate_id
               AND prev.id < o.id
               AND prev.status IN ('pending', 'processing', 'dead'))
     ORDER BY o.id
     LIMIT %(batch)s
       FOR UPDATE SKIP LOCKED
)
UPDATE outbox AS o
   SET status = 'processing', lease_owner = %(owner)s, lease_until = now() + %(lease)s,
       lease_version = o.lease_version + 1, attempts = o.attempts + 1
  FROM picked
 WHERE o.id = picked.id
RETURNING o.id, o.lease_version, o.attempts, o.aggregate_id, o.event_type, o.payload
```

- While id=1 of an aggregate is `pending` or `processing`, id=2 isn't a candidate. If two workers claim at once, `SKIP LOCKED` gives id=1 to only one of them, and the other excludes id=2 because its predecessor is unfinished.
- **`dead` also counts as an unfinished predecessor**, so if id=1 goes `dead`, everything after it in that aggregate stops. If you're preserving order, you can't skip a broken message and move on; processing stays stopped until a human fixes id=1 or deliberately skips it (marks it `done`).
- In my tests, four workers processing 30 messages for the same aggregate finalized them in exactly ascending id order.
- **What this guarantees is claim and finalize order, not the order messages reach the broker.** If a lease expires and a stale worker's in-flight publish of id=1 lands after the worker that took over has finished id=1 and moved on to id=2, the broker can see id=2 before id=1. Fencing only protects the database's state. Consumers still need idempotent state transitions (version or sequence checks) for the ordering to be safe end to end.

There's a cost. Each aggregate is serialized, so throughput drops when messages concentrate on one aggregate, and the `NOT EXISTS` needs an index on `(aggregate_id, id)`. In practice, use this claim **only for event types that need ordering**, not for ones like notifications. If downstream can absorb reordering (by writing idempotent state transitions), you don't need ordering guarantees at all.

---

## 6. Monitoring

```sql
-- Queue depth (count per status)
SELECT status, count(*) FROM outbox GROUP BY status;

-- Age of the oldest message awaiting delivery (the key metric; grows monotonically when the dispatcher stalls)
SELECT now() - min(created_at) AS oldest_pending_age
  FROM outbox
 WHERE status IN ('pending', 'processing');

-- processing rows left with an expired lease (a sign of crashing workers)
SELECT count(*) AS expired_leases
  FROM outbox
 WHERE status = 'processing' AND lease_until < now();

-- Distribution of attempts (a sign that retries are increasing)
SELECT attempts, count(*)
  FROM outbox
 WHERE status IN ('pending', 'processing')
 GROUP BY attempts ORDER BY attempts;

-- Messages that need a human
SELECT id, aggregate_id, event_type, attempts, last_error, created_at
  FROM outbox
 WHERE status = 'dead'
 ORDER BY id;
```

Alerting guidelines:

- **Oldest message awaiting delivery** exceeds the SLO (say, 5 minutes) → the dispatcher stopped, an external service is down, or throughput is insufficient.
- **Any increase in `dead`** → a human checks the cause. With ordered claiming, that aggregate's later messages are blocked.
- **`expired_leases` keeps growing** → workers are crashing, or the lease is too short for the processing time.
- **The count of `lost lease` warnings** as an application metric → revisit the lease length or batch size.

Delete `done` rows periodically, keeping them for as long as you want delivery history available during investigations.

```sql
DELETE FROM outbox WHERE status = 'done' AND done_at < now() - interval '7 days';
```

---

## 7. Tests: reproduce concurrency, crashes and expiry

These tests were run against a PostgreSQL 18 container and all pass:

| Test | Expected result |
| --- | --- |
| 4 workers claim 200 rows concurrently | 200 IDs, no duplicates |
| Broken dispatcher (publish after COMMIT) with 2 workers | The same messages are published twice (failure reproduced) |
| 4 workers run `dispatch_once` over 100 rows | Each of the 100 is published once; all `done` |
| Worker stops after claiming → lease expires | Another worker re-claims with `lease_version` +1; the old worker's `mark_done` returns `False` |
| Lease expired but nobody took over | The original worker's `mark_done` succeeds |
| Transient failure three times (limit 3) | `pending` (`attempts=1`) → … → `dead` (`attempts=3`) |
| Business defer | Stays `pending`, `attempts=0`, `available_at` pushed back |
| Permanent failure | Straight to `dead` |
| Claimed and never finalized three times (limit 3) | `dead` with `attempts=3`; never taken again |
| Ordered claim: aggregate X (3 rows) and Y (2 rows) | First claim gets only each aggregate's head; the next becomes available when the head is `done` |
| Ordered claim: 30 rows of one aggregate across 4 workers | Finalized in ascending id order |

Crash tests don't need to kill a process. **Claiming and then never finalizing** produces the same state as a crash (`processing` with an expiring lease).

```python
def test_crashed_worker_lease_expires_and_stale_finalize_is_fenced(dsn):
    seed(1)
    with psycopg.connect(dsn) as ca, psycopg.connect(dsn) as cb:
        [a] = claim(ca, owner="A", batch=10, lease=timedelta(milliseconds=200))
        assert claim(cb, owner="B", batch=10, lease=timedelta(seconds=30)) == []  # can't take it during the lease
        time.sleep(0.3)
        [b] = claim(cb, owner="B", batch=10, lease=timedelta(seconds=30))
        assert b.id == a.id and b.lease_version == a.lease_version + 1
        assert mark_done(ca, a) is False  # the expired worker A is fenced off
        assert mark_done(cb, b) is True
```

---

## 8. When not to use this, and the alternatives

- **Very high delivery volume, or retention and replay of messages are requirements**: let the outbox focus on reliably handing messages from the database to a broker, and leave delivery itself to Kafka or SQS. Reading the WAL with CDC (Debezium and the like) also avoids polling load ([polling vs CDC in the transactional outbox article](/blog/transactional-outbox-pattern-reliable-event-publishing-guide)).
- **A single worker is enough**: then you need neither leases nor fencing. But you then need something else to guarantee only one is running (including duplicate starts during deploys), so starting with leases is often the safer choice.
- **Long-running jobs (minutes to hours)**: you'll need lease heartbeats, cancellation and progress tracking, and a dedicated job platform (such as a workflow engine) usually fits better.
- **Polling latency is unacceptable**: you can use `LISTEN/NOTIFY` to tell workers about new rows, keeping polling as a safety net. Note that `LISTEN` doesn't work through PgBouncer's transaction mode, so it needs a dedicated connection.

---

## Summary

- Row locks vanish at COMMIT. "Take with SKIP LOCKED, COMMIT, then call the external API" double-processes.
- SKIP LOCKED only stops two transactions grabbing the same row at the same instant. Ownership after COMMIT goes in the table as status, lease_until and lease_version.
- Claim increments lease_version, and finalize uses `WHERE id = ? AND lease_version = ?`. An old worker whose lease expired and was taken over updates zero rows.
- Fencing protects only the database's state. Duplicate external side effects don't go away, so combine it with downstream de-duplication by message ID and idempotency keys for external APIs.
- Treat failures in four classes (success / transient / business defer / permanent), with jittered backoff and dead-lettering.
- The key monitoring metric is the age of the oldest message awaiting delivery.

How the receiving side turns a message the dispatcher delivered at least once into exactly one effect is covered in [What exactly-once actually guarantees](/blog/exactly-once-at-least-once-idempotency-effectively-once-guide), and detecting and repairing drift between your database and external services after delivery stalls for a long time is covered in [Designing reconciliation](/blog/reconciliation-architecture-drift-detection-repair-guide).
