Postgres LISTEN/NOTIFY for Job Wake-Ups

Polling is the reason a naive database queue feels sluggish, and LISTEN/NOTIFY is the Postgres-native fix covered here as part of Database-Backed Job Queues in Backend Frameworks & Worker Scaling. Workers sleep on a notification channel and wake the instant a job commits, which drops pickup latency from up to a second to a few milliseconds and removes the constant empty-poll query load from the primary.

Problem Statement

A Postgres-backed queue (built as in building a Postgres job queue with SKIP LOCKED) runs 40 worker processes that each poll once per second. The queue is idle most of the day, yet those polls account for 40 queries per second on the primary — 3.4 million a day that return nothing. Meanwhile users notice that an "export ready" email arrives up to a second after the button click, and product wants it faster. Lowering the poll interval to 100 ms would make latency acceptable and multiply the wasted queries by ten. You want workers to wake on demand, keep a safety-net poll for missed notifications, and survive connection poolers and failovers.

Prerequisites

  • A job table and claim query already working with polling.
  • Postgres 11+ (any supported version has LISTEN/NOTIFY; newer versions improve notification queue handling).
  • A driver with asynchronous notification support: psycopg 3, asyncpg, pgx for Go, pg for Node.
  • Knowledge of how your connections reach Postgres — directly, through PgBouncer, or through a cloud proxy — because it decides where the listener can connect.

Step 1 — Understand the Delivery Semantics

NOTIFY is transactional on the sending side and fire-and-forget on the receiving side. A notification raised inside a transaction is delivered only when that transaction commits, and never if it rolls back — exactly right for a queue, since a worker is never woken for a job that does not exist. But a listener that is disconnected at that moment never receives it; there is no replay. Notifications with identical channel and payload raised in the same transaction are collapsed into one.

Those properties give the design rule for the whole guide: notifications are hints, the table is the truth. A notification means "something may be claimable now", and the worker responds by running its normal claim query. Losing a notification must cost latency, never correctness.

Commit, wake, claim The producer inserts a job row and raises a notification inside one transaction. On commit, Postgres delivers the notification to every listening session. An idle worker waiting on its listener connection wakes and runs the usual claim query to take the job. If the transaction rolls back, no notification is delivered. Notification is delivered at COMMIT producer txn INSERT job; NOTIFY COMMIT Postgres delivers to listeners idle worker wakes runs normal claim query Rollback: no notification. Disconnected listener: notification lost, fallback poll finds the job.

Step 2 — Raise Notifications from a Trigger

A trigger keeps the notification in lockstep with the insert, so no producer can forget it. Notify per queue so workers only wake for queues they serve, and keep the payload empty — the worker re-queries anyway, and payloads are limited to 8,000 bytes.

CREATE OR REPLACE FUNCTION jobs_notify() RETURNS trigger AS $$
BEGIN
  -- one notification per queue per transaction: identical notifies are collapsed
  PERFORM pg_notify('jobs:' || NEW.queue, '');
  RETURN NULL;
END $$ LANGUAGE plpgsql;

CREATE TRIGGER jobs_notify_insert
AFTER INSERT ON jobs
FOR EACH ROW
WHEN (NEW.run_at <= now())            -- scheduled jobs are found by polling when due
EXECUTE FUNCTION jobs_notify();

The WHEN clause skips notifications for jobs scheduled in the future; the worker's fallback poll picks them up when they become due. A bulk insert of 10,000 jobs in one transaction raises 10,000 pg_notify calls that Postgres collapses into one delivered notification per channel — the collapse is what keeps bulk enqueues from causing a storm.

Step 3 — Hold a Dedicated Listener Connection

A LISTEN belongs to a session. Each worker process opens one direct connection used only for listening, separate from the pool it uses for claims and handlers.

# listener.py — psycopg 3
import threading
import psycopg

class JobListener:
    def __init__(self, dsn: str, queues: list[str]) -> None:
        self.dsn, self.queues = dsn, queues
        self.wake = threading.Event()

    def run(self, stop: threading.Event) -> None:
        while not stop.is_set():
            try:
                with psycopg.connect(self.dsn, autocommit=True) as conn:
                    for q in self.queues:
                        conn.execute(f'LISTEN "jobs:{q}"')
                    self.wake.set()                      # re-poll after (re)connecting
                    for _ in conn.notifies(timeout=30):  # yields notifications as they arrive
                        self.wake.set()
                        if stop.is_set():
                            return
            except psycopg.OperationalError:
                self.wake.set()                          # disconnected: fall back to polling
                stop.wait(2)                             # brief backoff, then reconnect

Setting wake after every (re)connect matters: notifications raised while the listener was down are gone, so the worker must poll once immediately to catch up.

Step 4 — Replace the Sleep with a Wait

The worker loop from the polling version slept for a fixed second when the queue was empty. Now it waits on the event with a long timeout, which becomes the safety-net poll interval.

FALLBACK_POLL_S = 15            # catches missed notifications and newly-due scheduled jobs

def worker_main(pool, listener: JobListener, queues, concurrency=8):
    while True:
        jobs = claim(pool, WORKER_ID, queues, limit=free_slots())
        dispatch(jobs)
        if jobs and len(jobs) == free_slots():
            continue                                # queue likely has more: claim again now
        listener.wake.wait(timeout=FALLBACK_POLL_S)
        listener.wake.clear()

The continue when a full batch was claimed is important: a notification wakes the worker once, but a burst of 500 jobs needs many claims. Keep claiming while batches come back full, and only wait when the queue is drained.

Fewer queries, lower latency Forty workers polling every second issue about 3.4 million empty claim queries a day and pick up jobs up to one second late. With notifications and a fifteen-second fallback poll, the same fleet issues about 230 thousand fallback queries a day and picks up jobs within a few milliseconds of commit. 40 workers, idle queue, one day poll every 1 s 3.4M empty queries, pickup up to 1,000 ms NOTIFY + 15 s poll 230k fallback queries, pickup ~5 ms Bars to scale. The fallback interval only bounds latency for missed notifications and scheduled jobs.

Step 5 — Work Around Connection Poolers

PgBouncer in transaction pooling mode — the most common setup — hands a server connection to a client only for the duration of a transaction. A LISTEN issued through it registers on whichever server connection happened to be assigned and is effectively lost when that connection returns to the pool. The notifications then go to some other client, or nowhere.

; pgbouncer.ini — two pools: transaction mode for app traffic, session mode for listeners
[databases]
app       = host=pg-primary port=5432 dbname=app pool_mode=transaction
app_listen = host=pg-primary port=5432 dbname=app pool_mode=session pool_size=60

Point the listener DSN at a session-mode pool (or directly at the primary) and keep claim and handler traffic on the transaction-mode pool. Each worker process holds one listener connection permanently, so budget for it: 40 worker processes means 40 extra long-lived connections. Cloud proxies (RDS Proxy, some managed pgbouncers) may "pin" or reject LISTEN; check their documentation and connect listeners directly if needed.

Step 6 — Handle Failover and Replicas

LISTEN/NOTIFY works only on the primary. Hot-standby replicas do not deliver notifications, and a listener connected to a replica waits forever. After a failover, every listener's connection breaks; the reconnect loop in Step 3 handles it, provided the DSN follows the primary (a DNS name updated on failover, or target_session_attrs=read-write in a multi-host connection string).

LISTEN_DSN = (
    "host=pg-a.internal,pg-b.internal port=5432 dbname=app user=worker "
    "target_session_attrs=read-write connect_timeout=5 "     # always land on the primary
    "keepalives=1 keepalives_idle=30 keepalives_interval=10 keepalives_count=3"
)

TCP keepalives detect a silently dead connection (a network partition where no RST arrives) within about a minute instead of the kernel default of hours. Until detection, the fallback poll keeps jobs flowing.

Step 7 — Avoid Notification Storms and Thundering Herds

With 40 workers listening on one channel, every committed job wakes all 40, and all 40 run the claim query at once. With SKIP LOCKED this is safe — one gets the job and the rest get nothing — but at high enqueue rates the wasted claims add up. Three mitigations, in order of preference:

  • Batch claiming means one wake-up takes many jobs, so the herd is small relative to useful work.
  • Per-queue channels (done in Step 2) mean workers only wake for queues they serve.
  • Wake a subset: have each process listen, but let only one thread per process claim on wake, with others waiting on the local event. For very large fleets, shard queues into sub-channels (jobs:default:0 .. jobs:default:7) chosen by a hash of the job id, and have each worker listen to one or two shards.
-- Monitor the notification queue: should stay near 0
SELECT pg_notification_queue_usage();     -- fraction of the 8 GB queue in use

If a listener stops reading (a hung process holding its connection), undelivered notifications accumulate in a shared queue; when it fills, NOTIFY fails with an error in the producer's transaction, which would break enqueue. Alert well before that point, and terminate stuck listener sessions.

Shard channels to shrink the herd With one channel, a single committed job wakes all forty workers and thirty-nine claim queries return nothing. With eight sharded channels, the same job wakes only the five workers listening on its shard. Workers woken per committed job one channel 40 woken, 39 empty claims 8 shard channels 5 woken, 4 empty claims Only worth doing at high enqueue rates; batch claiming already absorbs most of the herd.

Verification

Measure pickup latency before and after, and prove that lost notifications do not lose jobs:

-- Pickup latency: time from insert to claim, recorded by the worker on claim
SELECT percentile_cont(ARRAY[0.5, 0.95, 0.99]) WITHIN GROUP (ORDER BY claimed_at - created_at)
FROM job_audit WHERE created_at > now() - interval '1 hour';
def test_job_found_after_listener_outage(pool, worker, kill_listener):
    kill_listener()                                    # drop the LISTEN connection
    with pool.connection() as c, c.transaction():
        jid = enqueue(c, "noop", {})                   # notification is lost
    assert wait_until_done(jid, timeout=FALLBACK_POLL_S + 5)

Expect p95 pickup in single-digit milliseconds on an idle queue and the outage test to pass within the fallback interval.

Gotchas & Edge Cases

Notifying from inside long transactions. Notifications wait for commit. A batch import that enqueues jobs inside a 10-minute transaction wakes nobody until it commits — by design, but surprising if you expected streaming.

Large payloads. Payloads over 8,000 bytes raise an error. Keep payloads empty or tiny; never ship job arguments through the channel.

Listener connections and idle_session_timeout. A server-side idle timeout will kill an idle listener. Exclude the worker role, or rely on the reconnect loop and accept a brief fallback-poll window.

Serverless drivers. HTTP-based Postgres drivers (used from edge functions) cannot LISTEN. Keep listeners in long-running worker processes.

FAQ

Can I skip polling entirely? No. Notifications are lost during disconnects and are not raised for jobs scheduled in the future. A slow fallback poll is what makes the design correct; the notification only makes it fast.

Does LISTEN/NOTIFY scale to many channels? Hundreds of channels are fine. The shared notification queue and the per-commit delivery work are the costs, and they scale with notification volume, not channel count.

Is this what Oban, River, and Solid Queue use? Oban and River use LISTEN/NOTIFY for wake-ups with polling as a fallback. Solid Queue polls by default with a configurable interval, trading latency for fewer moving parts.

Related