Skip to main content
Reliability, async & real-time
PostgreSQL
信頼性
アーキテクチャ設計
Python
データベース

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
Reading time
20 min read
Author
友田 陽大
Share

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 — 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:

# ❌ 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.

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 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.


I can take on the implementation from this article as an engagement

Data-layer architecture: ORM selection, schema design, and zero-downtime migration

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).


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

   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

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

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

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:

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).
  • 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:

ClassExamplesHandlingAttempts
Success / terminal no-opSent / the receiver is already in the desired statedone—
Transient failure (technical retry)Timeout, 5xx, 429, connection lossBack to pending with jittered backoff; dead after the limitConsumed
Business deferA resource it depends on (a customer account, say) doesn't exist yetBack to pending after a delayNot consumed
Permanent failure4xx validation error, nonexistent destinationdead (a human reviews it)—
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.

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

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), writes to the same aggregate are serialized until commit and the precondition holds.

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

-- 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.

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:

TestExpected result
4 workers claim 200 rows concurrently200 IDs, no duplicates
Broken dispatcher (publish after COMMIT) with 2 workersThe same messages are published twice (failure reproduced)
4 workers run dispatch_once over 100 rowsEach of the 100 is published once; all done
Worker stops after claiming → lease expiresAnother worker re-claims with lease_version +1; the old worker's mark_done returns False
Lease expired but nobody took overThe original worker's mark_done succeeds
Transient failure three times (limit 3)pending (attempts=1) → … → dead (attempts=3)
Business deferStays pending, attempts=0, available_at pushed back
Permanent failureStraight 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 workersFinalized 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).

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).
  • 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, and detecting and repairing drift between your database and external services after delivery stalls for a long time is covered in Designing reconciliation.

Frequently asked questions

If I use FOR UPDATE SKIP LOCKED, can several workers process the outbox without double processing?
If you keep the transaction open until the external API call completes, the row lock persists, but you hold a connection the whole time and still face the API succeeding and then COMMIT failing. If you commit before the API call, the row lock is gone and another worker can take the same row. Either way, you need a lease written to the table to represent ownership after COMMIT.
What is the difference between a lease and a fencing token?
A lease is time-limited ownership — this worker owns the message until a given time — so if the worker dies, another can take over once the lease expires. A fencing token is a number that increases on every claim; at finalize time only the worker whose number still matches may write. If an old worker returns after its lease expired, its number is stale and its write is rejected.
Does a fencing token also prevent duplicate external API calls?
No. A fencing token protects the final write to the database, but if the old worker already called the external API around the time its lease expired, that side effect remains. Preventing it on the external side requires the external system to check the token or an idempotency key. With an outbox, you combine fencing with downstream de-duplication that uses the message ID as an idempotency key.
Is it OK to use PostgreSQL as a job queue?
The value of an outbox is that messages are written in the same database as the business transaction, so using PostgreSQL as the outbox's delivery queue is reasonable. For very high throughput, such as tens of thousands of messages per second, or for holding large numbers of long-running jobs, a dedicated queue or broker fits better. Monitor queue depth and the age of the oldest message, and move delivery to a broker if database load becomes a problem.

References

友田

友田 陽大

Developer of a METI Minister's Award–winning product. With TypeScript + Python + AWS, I deliver SaaS, industry DX, and production-grade generative AI (RAG) end to end — from requirements to infrastructure and operations — single-handedly.

I can take on the implementation from this article as an engagement

Data-layer architecture: ORM selection, schema design, and zero-downtime migration

The real question is not which ORM you pick — it is whether the data model survives five years of change. I handle the selection (Prisma / Drizzle / SQLAlchemy), where to normalise and where not to, N+1 and connection-pool design, and schema migrations that never take the service down. Having led the reliability layer of a payment platform — designing the idempotency and consistency that kept production double-charges at zero — I build data layers that fail loudly and roll back cleanly.

Available for both project-based (contract) and advisory engagements. Start with a free 30-minute consult.

Also worth reading