Building a Postgres Job Queue with SKIP LOCKED

This guide builds a complete job queue on plain Postgres, the pattern introduced in Database-Backed Job Queues within Backend Frameworks & Worker Scaling. By the end you have transactional enqueue, concurrent claiming, retries with backoff, crash recovery, and cleanup — roughly 150 lines of SQL and Python, and every piece is something a library would otherwise hide from you.

Problem Statement

A Python service already runs on Postgres and needs background jobs for a few hundred tasks per minute: webhook deliveries, PDF generation, and nightly syncs. The team does not want to operate Redis just for this, and it has been bitten by jobs enqueued for rows that were later rolled back. Library options exist (Procrastinate, Dramatiq with a Postgres broker), but the team wants to understand and own the mechanics first. The requirements are: enqueue in the same transaction as the business write, several worker processes claiming without duplicates or blocking, bounded retries with exponential backoff, automatic recovery of jobs held by crashed workers, and a table that does not grow without bound.

Prerequisites

  • Postgres 12 or later (SKIP LOCKED exists from 9.5; 12+ gives better partial-index planning and generated columns).
  • psycopg 3 and a connection pool (psycopg_pool) in the worker.
  • Permission to create tables, indexes, and functions in the application database.
  • Idempotent job handlers — any queue that survives crashes delivers at least once.

Step 1 — Create the Table and the Claim Index

Keep the job table narrow. Arguments go in jsonb; everything the claim query filters or sorts on is a real column so the index can serve it.

CREATE TABLE jobs (
  id           bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
  queue        text        NOT NULL DEFAULT 'default',
  kind         text        NOT NULL,                 -- handler name
  args         jsonb       NOT NULL DEFAULT '{}',
  priority     smallint    NOT NULL DEFAULT 0,       -- higher runs first
  run_at       timestamptz NOT NULL DEFAULT now(),
  state        text        NOT NULL DEFAULT 'available'
               CHECK (state IN ('available','running','completed','discarded')),
  attempts     int         NOT NULL DEFAULT 0,
  max_attempts int         NOT NULL DEFAULT 10,
  timeout_s    int         NOT NULL DEFAULT 300,     -- rescue after this long without heartbeat
  locked_by    text,
  heartbeat_at timestamptz,
  last_error   text,
  created_at   timestamptz NOT NULL DEFAULT now(),
  finished_at  timestamptz
) WITH (fillfactor = 70, autovacuum_vacuum_scale_factor = 0.01);

-- Serves the claim query and nothing else: only candidate rows are indexed
CREATE INDEX jobs_claim ON jobs (queue, priority DESC, run_at, id) WHERE state = 'available';
-- Serves the rescuer
CREATE INDEX jobs_running ON jobs (heartbeat_at) WHERE state = 'running';

fillfactor = 70 leaves free space in each page so state changes can often be written as heap-only tuple (HOT) updates that do not touch indexes. The low autovacuum_vacuum_scale_factor makes vacuum run after 1% of rows are dead rather than the default 20%, which matters for a table whose rows die constantly.

The partial index stays small The jobs table heap holds available, running, completed, and discarded rows. The claim index includes only available rows, so its size tracks the backlog rather than the table's history, and the claim query reads a handful of index entries regardless of how many completed jobs remain. Heap size vs index size jobs heap: 4.2M rows available 1.8k running 40 completed 4.19M discarded 310 jobs_claim index 1.8k entries only Deleting completed rows still matters for heap bloat and vacuum, but not for claim speed.

Step 2 — Enqueue Inside the Caller's Transaction

The enqueue function takes an open connection or cursor rather than opening its own. That is the whole trick: whatever transaction the caller is in, the job row joins it.

# queue.py
import json
from psycopg import Connection

def enqueue(conn: Connection, kind: str, args: dict, *, queue: str = "default",
            priority: int = 0, delay_s: int = 0, max_attempts: int = 10) -> int:
    row = conn.execute(
        """INSERT INTO jobs (queue, kind, args, priority, run_at, max_attempts)
           VALUES (%s, %s, %s, %s, now() + make_interval(secs => %s), %s)
           RETURNING id""",
        (queue, kind, json.dumps(args), priority, delay_s, max_attempts),
    ).fetchone()
    return row[0]

# usage: the webhook job exists if and only if the order commits
with pool.connection() as conn, conn.transaction():
    order_id = create_order(conn, payload)
    enqueue(conn, "deliver_webhook", {"order_id": order_id, "event": "order.created"})

Resist the temptation to enqueue on a separate autocommit connection "for speed". That reintroduces the dual-write problem this design exists to remove.

Step 3 — Claim a Batch Without Blocking

The claim runs as its own short transaction: select candidates with FOR UPDATE SKIP LOCKED, mark them running, commit. The handler then runs outside any transaction, so a slow job never holds a row lock or an open transaction that would block vacuum.

CLAIM_SQL = """
WITH next AS (
  SELECT id FROM jobs
  WHERE state = 'available' AND queue = ANY(%(queues)s) AND run_at <= now()
  ORDER BY priority DESC, run_at, id
  LIMIT %(limit)s
  FOR UPDATE SKIP LOCKED
)
UPDATE jobs j
SET state = 'running', locked_by = %(worker)s, heartbeat_at = now(), attempts = j.attempts + 1
FROM next
WHERE j.id = next.id
RETURNING j.id, j.kind, j.args, j.attempts, j.max_attempts, j.timeout_s
"""

def claim(pool, worker_id: str, queues: list[str], limit: int = 10) -> list[tuple]:
    with pool.connection() as conn:                       # autocommit off: one short txn
        rows = conn.execute(CLAIM_SQL, {"queues": queues, "limit": limit,
                                        "worker": worker_id}).fetchall()
        conn.commit()
    return rows

attempts is incremented at claim time, not at failure time. If the worker dies mid-job, the attempt still counts, so a job that crashes its worker every time (an OOM on a huge input) eventually reaches max_attempts instead of looping forever.

Keep transactions short around a long job The claim is a transaction lasting a few milliseconds that marks rows as running and commits. The handler then runs for seconds or minutes with no transaction open. A second short transaction records completion or failure. No lock or snapshot is held while the handler runs. Two short transactions, no long one claim txn ~3 ms handler runs: no transaction open heartbeat updates every 20 s finish txn ~2 ms anti-pattern: claim, run, and finish in one transaction A transaction held open for the job's duration pins vacuum and holds row locks for minutes.

Step 4 — Run the Handler and Record the Outcome

Completion deletes the row (the cheapest possible cleanup). Failure either reschedules with backoff or discards the job once attempts are exhausted. Both are guarded by locked_by, so a worker whose job was rescued and re-claimed elsewhere cannot overwrite the new owner's state.

import random, traceback

HANDLERS = {"deliver_webhook": deliver_webhook, "render_pdf": render_pdf}

def finish_ok(pool, job_id: int, worker_id: str) -> None:
    with pool.connection() as conn:
        conn.execute("DELETE FROM jobs WHERE id=%s AND locked_by=%s AND state='running'",
                     (job_id, worker_id))

def finish_err(pool, job_id: int, worker_id: str, attempts: int, max_attempts: int, err: str) -> None:
    backoff = min(2 ** attempts, 3600) * random.uniform(0.5, 1.0)     # full-ish jitter, 1 h cap
    with pool.connection() as conn:
        conn.execute(
            """UPDATE jobs SET
                 state = CASE WHEN %(a)s >= %(m)s THEN 'discarded' ELSE 'available' END,
                 run_at = now() + make_interval(secs => %(b)s),
                 last_error = %(e)s, locked_by = NULL, heartbeat_at = NULL,
                 finished_at = CASE WHEN %(a)s >= %(m)s THEN now() END
               WHERE id = %(id)s AND locked_by = %(w)s AND state = 'running'""",
            {"a": attempts, "m": max_attempts, "b": backoff, "e": err[:4000],
             "id": job_id, "w": worker_id})

def run_one(pool, worker_id, job) -> None:
    job_id, kind, args, attempts, max_attempts, _ = job
    try:
        HANDLERS[kind](**args)
    except Exception:
        finish_err(pool, job_id, worker_id, attempts, max_attempts, traceback.format_exc())
    else:
        finish_ok(pool, job_id, worker_id)

Discarded rows are the dead-letter queue: they keep their arguments and last error, and replay is an UPDATE ... SET state='available', attempts=0. The backoff formula mirrors the one in exponential backoff with jitter in Celery.

Step 5 — Heartbeat Long Jobs and Rescue Crashed Ones

A worker that dies leaves its rows in running. A rescuer returns rows with a stale heartbeat to available. Long jobs must heartbeat so they are not rescued while still working.

import threading

def heartbeat_loop(pool, worker_id: str, job_ids: set[int], stop: threading.Event) -> None:
    while not stop.wait(20):                              # every 20 s
        if job_ids:
            with pool.connection() as conn:
                conn.execute("UPDATE jobs SET heartbeat_at = now() "
                             "WHERE id = ANY(%s) AND locked_by = %s",
                             (list(job_ids), worker_id))

RESCUE_SQL = """
UPDATE jobs SET state = 'available', locked_by = NULL, heartbeat_at = NULL,
                last_error = 'rescued: worker stopped heartbeating'
WHERE state = 'running' AND heartbeat_at < now() - make_interval(secs => timeout_s)
RETURNING id
"""

def rescue(pool) -> int:
    with pool.connection() as conn:
        return len(conn.execute(RESCUE_SQL).fetchall())

Run rescue every minute from any one worker, guarded by pg_try_advisory_lock(42) so only one process does it at a time. Because attempts was incremented at claim time, a job that keeps killing its worker is eventually discarded rather than rescued forever.

Heartbeat age decides who gets rescued Worker A updates the heartbeat of its long job every 20 seconds, so the heartbeat never becomes older than the timeout and the job is left alone. Worker B crashes after its last heartbeat; once the heartbeat is older than the job's timeout, the rescuer resets the row to available and another worker claims it. Live worker vs crashed worker worker A never stale worker B crashed: heartbeat ages rescued Rescue fires when heartbeat_at is older than timeout_s, which is set per job kind.

Step 6 — Wire Up the Worker Loop

The loop claims up to concurrency jobs, runs them in a thread pool, and sleeps briefly only when the queue is empty.

from concurrent.futures import ThreadPoolExecutor
import os, socket, time

def worker_main(queues=("default",), concurrency=8):
    worker_id = f"{socket.gethostname()}:{os.getpid()}"
    in_flight: set[int] = set()
    stop = threading.Event()
    threading.Thread(target=heartbeat_loop, args=(pool, worker_id, in_flight, stop), daemon=True).start()
    with ThreadPoolExecutor(max_workers=concurrency) as ex:
        while not stop.is_set():
            free = concurrency - len(in_flight)
            jobs = claim(pool, worker_id, list(queues), limit=free) if free else []
            for job in jobs:
                in_flight.add(job[0])
                ex.submit(run_one, pool, worker_id, job).add_done_callback(
                    lambda _f, jid=job[0]: in_flight.discard(jid))
            if not jobs:
                time.sleep(1.0)                          # or wait on LISTEN, see below

Size concurrency against the connection pool: each running handler may need its own connection for business queries, plus the claim, heartbeat, and finish calls. Replace the one-second sleep with LISTEN for near-instant pickup, as shown in Postgres LISTEN/NOTIFY for job wake-ups. On SIGTERM, set stop, let in-flight jobs finish within your grace period, and exit — the pattern from graceful shutdown and deployments.

Step 7 — Keep the Table Small

Deleting on success keeps the heap from growing with history, but discarded rows accumulate, and heavy delete traffic still needs vacuum. Schedule a cleanup of old discarded jobs and verify autovacuum keeps up.

-- Nightly: drop dead letters older than 14 days
DELETE FROM jobs WHERE state = 'discarded' AND finished_at < now() - interval '14 days';

-- Is vacuum keeping up? dead tuples should stay far below live tuples
SELECT n_live_tup, n_dead_tup, last_autovacuum, autovacuum_count
FROM pg_stat_user_tables WHERE relname = 'jobs';

If you need an audit trail of completed jobs, insert a compact record into an append-only, time-partitioned job_history table on completion rather than keeping rows in jobs.

Verification

Prove the two properties that matter — no duplicate claims, and recovery after a crash:

def test_concurrent_claims_are_disjoint(pool):
    with pool.connection() as c, c.transaction():
        for i in range(1000):
            enqueue(c, "noop", {"i": i})
    claimed = []
    def grab():
        while rows := claim(pool, f"w{threading.get_ident()}", ["default"], 25):
            claimed.extend(r[0] for r in rows)
    threads = [threading.Thread(target=grab) for _ in range(8)]
    [t.start() for t in threads]; [t.join() for t in threads]
    assert len(claimed) == 1000 and len(set(claimed)) == 1000

def test_crashed_job_is_rescued(pool):
    with pool.connection() as c, c.transaction():
        jid = enqueue(c, "noop", {})
    claim(pool, "dead-worker", ["default"], 1)
    with pool.connection() as c:
        c.execute("UPDATE jobs SET heartbeat_at = now() - interval '1 hour' WHERE id=%s", (jid,))
    assert rescue(pool) == 1

In staging, run EXPLAIN (ANALYZE, BUFFERS) on the claim query with a realistic backlog; it should use jobs_claim with an index scan and touch a few buffers.

Gotchas & Edge Cases

The claim query sorts by the wrong thing. If the ORDER BY does not match the index column order, Postgres sorts every candidate row. Keep the index and the query in lockstep, including DESC.

Priority starvation. Strict priority DESC ordering starves low-priority jobs under sustained high-priority load. Use separate queues with weighted polling if that matters — the same issue discussed in preventing tenant starvation with weighted queues.

Clock assumptions. run_at <= now() uses the database clock, which is the right choice; never compare against application-server time.

Serializable isolation. Under SERIALIZABLE, SKIP LOCKED claims can raise serialization failures. Run the claim at READ COMMITTED.

FAQ

How many workers can claim from one table? Dozens of processes are routine. Contention appears when many workers repeatedly lock the same few head rows; batch claiming and a small random OFFSET jitter for very large fleets both help.

Should I use a library instead? For most teams, yes — once you understand these mechanics. Procrastinate (Python), River (Go), Oban (Elixir), Solid Queue and GoodJob (Ruby) implement everything above with years of edge cases fixed.

Can I use advisory locks instead of SKIP LOCKED? You can, but advisory-lock queues are easy to get subtly wrong (locks leaking across pooled connections, planner-dependent lock ordering). SKIP LOCKED is simpler and is what modern libraries use.

Related