Using pgmq for Postgres Message Queues

pgmq is a Postgres extension that packages an SQS-like message queue as SQL functions, and this guide shows how to use it for background work as part of Database-Backed Job Queues in Backend Frameworks & Worker Scaling. Where a job framework gives you workers, retries, and handlers, pgmq gives you queue primitives — send, read with a visibility timeout, delete, archive — callable from any language that can run SQL.

Problem Statement

A platform team runs services in Python, TypeScript, and Go against one Postgres cluster. Each service needs a small work queue — thumbnail requests, webhook fan-out, cache invalidation — at a few hundred messages per second in total. A job framework per language would mean three different queue schemas and three sets of operational tooling. SQS would add a cloud dependency and the dual-write problem. The team wants a single, language-neutral queue API with visibility timeouts and at-least-once delivery, that participates in the application's transactions and is observable with SQL.

Prerequisites

  • Postgres 14–17 with permission to install extensions (pgmq is available on Supabase, Tembo, and many managed providers; self-hosted installs use the prebuilt packages or pgxn).
  • A role that can execute functions in the pgmq schema.
  • A consumer design that tolerates redelivery — pgmq is at-least-once, like SQS.
  • A rough message volume estimate, to decide between standard and partitioned queues (Step 6).

Step 1 — Install the Extension and Create a Queue

Each queue is a pair of tables inside the pgmq schema: q_<name> for live messages and a_<name> for archived ones.

CREATE EXTENSION IF NOT EXISTS pgmq;

SELECT pgmq.create('thumbnails');          -- standard (logged) queue
SELECT pgmq.create_unlogged('cache_bust'); -- faster, but lost on crash: only for disposable work

-- Inspect what was created
\dt pgmq.*
--  pgmq.q_thumbnails   live messages: msg_id, read_ct, enqueued_at, vt, message (jsonb)
--  pgmq.a_thumbnails   archive: same columns plus archived_at

The vt (visibility time) column is the heart of the design: a message is readable only when vt <= now(). Reading a message pushes its vt into the future, which hides it from other consumers — the same mechanism as an SQS visibility timeout, implemented as a column.

A visibility timeout stored as a column pgmq.send inserts a message with vt equal to now, making it visible. pgmq.read returns visible messages and sets vt to now plus the timeout, hiding them. If the consumer deletes or archives the message before vt passes, it is gone. If not, vt passes and the message becomes visible again with read_ct incremented on the next read. send, read, then delete or archive visible vt <= now() hidden vt = now() + 30s deleted / archived vt passes without delete: visible again, read_ct + 1 pgmq.read pgmq.delete / archive

Step 2 — Send Messages Inside Business Transactions

pgmq.send is a normal function call, so it participates in whatever transaction the caller is in. A message sent alongside a business write is committed with it or not at all.

BEGIN;
INSERT INTO uploads (id, owner_id, s3_key) VALUES ('u-981', 'acct-7', 'raw/u-981.png');
SELECT pgmq.send(
  queue_name => 'thumbnails',
  msg        => '{"upload_id": "u-981", "sizes": [64, 256, 1024]}'::jsonb,
  delay      => 0                                   -- seconds before it becomes visible
);
COMMIT;

-- Batch send returns one msg_id per message
SELECT * FROM pgmq.send_batch('thumbnails', ARRAY[
  '{"upload_id": "u-982"}'::jsonb, '{"upload_id": "u-983"}'::jsonb ]);

From application code it is just SQL, which is the point: the TypeScript service and the Go service use the same queue without a shared client library.

// Node: pg driver, inside the caller's transaction
await client.query("BEGIN");
await client.query("INSERT INTO uploads (id, owner_id, s3_key) VALUES ($1, $2, $3)", [id, owner, key]);
await client.query("SELECT pgmq.send($1, $2::jsonb)", ["thumbnails", JSON.stringify({ upload_id: id })]);
await client.query("COMMIT");

Step 3 — Read with a Visibility Timeout and Delete on Success

A consumer reads a batch, processes each message, and deletes (or archives) it. If the consumer crashes, the visibility timeout expires and the message is delivered again.

# consumer.py — psycopg 3
import json, time
import psycopg

VT_SECONDS = 60          # must exceed worst-case processing time for one message
BATCH = 20

def consume(conn: psycopg.Connection) -> None:
    while True:
        rows = conn.execute(
            "SELECT msg_id, read_ct, message FROM pgmq.read(%s, %s, %s)",
            ("thumbnails", VT_SECONDS, BATCH)).fetchall()
        conn.commit()
        if not rows:
            # read_with_poll blocks server-side up to max_poll_seconds instead of sleeping here
            time.sleep(1)
            continue
        for msg_id, read_ct, message in rows:
            try:
                make_thumbnails(message)
                conn.execute("SELECT pgmq.archive('thumbnails', %s)", (msg_id,))
            except Exception:
                if read_ct >= 5:
                    dead_letter(conn, msg_id, message)       # see Step 5
                # otherwise do nothing: it reappears after VT_SECONDS
            conn.commit()

pgmq.read_with_poll(queue, vt, qty, max_poll_seconds, poll_interval_ms) replaces the client-side sleep with a server-side wait that returns as soon as a message is visible — useful for keeping latency low without a tight client loop. For finer control over retry spacing, call pgmq.set_vt(queue, msg_id, seconds) on failure to push the next attempt out by an exponential backoff instead of waiting the full VT_SECONDS.

Step 4 — Choose Archive or Delete

pgmq.delete removes the row. pgmq.archive moves it to the a_ table with an archived_at timestamp, giving you an audit trail of processed messages at the cost of storage and write volume.

-- Archive for an audit trail; purge old archive rows on a schedule
SELECT pgmq.archive('thumbnails', 4211);
DELETE FROM pgmq.a_thumbnails WHERE archived_at < now() - interval '7 days';

-- Or delete when nobody needs history (half the writes)
SELECT pgmq.delete('cache_bust', 90112);

-- Pop = read + delete in one call: at-most-once, the message is gone even if you crash
SELECT * FROM pgmq.pop('cache_bust');

pop is worth calling out: it trades at-least-once for at-most-once. It is correct for disposable work like cache invalidation, where losing a message is harmless, and wrong for anything that must happen.

Three ways to finish a message Read then delete is at-least-once with the fewest writes and no history. Read then archive is at-least-once with an audit trail but more writes and storage. Pop removes the message as it is read, which is at-most-once and suitable only for disposable work. Finishing a message read + delete at-least-once fewest writes, no history read + archive at-least-once audit trail, purge it pop at-most-once disposable work only

Step 5 — Dead-Letter Messages by Read Count

pgmq has no built-in redrive policy, but read_ct gives you everything needed to build one. A message whose read_ct exceeds a threshold is moved to a separate dead-letter queue in the same transaction that removes it from the source.

def dead_letter(conn, msg_id: int, message: dict) -> None:
    with conn.transaction():
        conn.execute("SELECT pgmq.send('thumbnails_dlq', %s::jsonb)",
                     (json.dumps({"original": message, "source_msg_id": msg_id,
                                  "failed_at": time.time()}),))
        conn.execute("SELECT pgmq.delete('thumbnails', %s)", (msg_id,))

Because both statements run in one transaction, a message is never lost between queues and never exists in both. Replaying is the reverse: read from thumbnails_dlq, send the original back to thumbnails, delete from the DLQ. The alerting side mirrors alerting on dead-letter queue growth — alert on any non-zero DLQ depth.

Step 6 — Partition High-Volume Queues

A standard pgmq queue is one table. At sustained high volume, dead tuples from reads and deletes accumulate faster than autovacuum clears them. Partitioned queues (backed by pg_partman) split the table by message id or time, so old partitions are dropped instead of vacuumed.

CREATE EXTENSION IF NOT EXISTS pg_partman;

-- New partition every 100k messages; keep 1M messages of retention in the archive
SELECT pgmq.create_partitioned(
  queue_name           => 'webhooks',
  partition_interval   => '100000',
  retention_interval   => '1000000'
);
-- pg_partman's background worker (or a cron call to run_maintenance) creates and drops partitions

The difference is in how space is reclaimed. On a single table, every read updates a row (new vt, incremented read_ct) and every delete leaves a dead tuple; vacuum must scan and clean them while consumers keep adding more. On a partitioned queue, new messages land in the newest partition, consumers drain the oldest ones, and once a partition is fully processed pg_partman detaches and drops it — reclaiming its space instantly with no vacuum work at all.

Vacuum a table, or drop a partition A standard queue is one table in which reads and deletes leave dead tuples that vacuum must clean continuously. A partitioned queue writes new messages into the newest partition while consumers drain older ones, and fully drained partitions are dropped, which frees their space without vacuum. Standard vs partitioned queue standard q_webhooks: live rows mixed with dead tuples vacuum never stops partitioned p1: dropped p2: draining p3: receiving sends space freed by DROP Partitioning pays off only at volumes where autovacuum on the single table cannot keep up.

Stay on standard queues until metrics tell you otherwise; partitioning adds a dependency and moving parts. The trigger to switch is autovacuum running continuously on pgmq.q_<name> while its dead-tuple count still climbs.

Step 7 — Monitor with the Metrics Function

pgmq exposes queue length, the age of the oldest and newest messages, and total messages ever sent — the numbers you would otherwise pull from CloudWatch for SQS.

SELECT queue_name, queue_length, oldest_msg_age_sec, newest_msg_age_sec, total_messages
FROM pgmq.metrics_all();

Expose this through postgres_exporter custom queries and alert on oldest_msg_age_sec per queue — the queue-time signal behind a latency SLO, as described in defining SLOs for job latency. queue_length includes invisible (in-flight) messages; compare it with a count(*) WHERE vt <= now() if you need visible depth alone.

Verification

-- Send, read, and let the vt expire without deleting: the message must come back
SELECT pgmq.send('thumbnails', '{"test": true}');
SELECT msg_id, read_ct FROM pgmq.read('thumbnails', 2, 1);   -- read_ct = 1
SELECT pg_sleep(3);
SELECT msg_id, read_ct FROM pgmq.read('thumbnails', 2, 1);   -- same msg_id, read_ct = 2

-- Transactional send: nothing is visible after a rollback
BEGIN; SELECT pgmq.send('thumbnails', '{"rolled_back": true}'); ROLLBACK;
SELECT count(*) FROM pgmq.q_thumbnails WHERE message ? 'rolled_back';   -- 0

Gotchas & Edge Cases

No ordering guarantee under concurrency. pgmq.read returns the oldest visible messages, but with several consumers and redeliveries, completion order is arbitrary. For per-key order, use FIFO read functions where your pgmq version provides them, or the techniques in Message Ordering Guarantees.

Unlogged queues vanish on crash. create_unlogged skips WAL, so the queue is emptied after a crash and is not replicated. Use it only for work that can be regenerated.

Visibility timeout too short. Exactly as with SQS, a vt shorter than processing time delivers the same message to a second consumer while the first is still working. Extend with pgmq.set_vt for long work.

Extension upgrades. pgmq's SQL API has changed between major versions. Pin the version, and read the upgrade notes before ALTER EXTENSION pgmq UPDATE.

FAQ

How is pgmq different from a job framework like Solid Queue or River? pgmq is a queue, not a job system: it has no workers, handler registry, or retry scheduler. You build those in your consumer. In exchange, any language can use it with plain SQL, and the semantics match SQS closely.

Can pgmq replace SQS? For moderate volumes inside one Postgres cluster, often yes, with the bonus of transactional sends. SQS scales further and costs nothing in database load, so very high volume or cross-account fan-out still favours it.

Does read use SKIP LOCKED? Yes. pgmq.read selects visible rows with FOR UPDATE SKIP LOCKED and updates their vt in one statement, the same technique described in building a Postgres job queue with SKIP LOCKED.

Related