When you need background jobs, the usual answer is Redis plus a queue library. But if you
already run Postgres, you may not need another service at all. Postgres has everything a
reliable job queue needs: row locks, transactions, SKIP LOCKED and LISTEN/NOTIFY.
I use one on my course platform to send hundreds of videos through a rate-limited video API. It's one table and a small worker. This post covers the pattern in general: the table, how workers claim jobs safely, retries, crash recovery, deduplication, waking workers without polling, and when you should reach for something else.
Why a queue in Postgres?
- One less service. No Redis to run, back up, secure and monitor.
- Transactional enqueue. Insert the job in the same transaction as the data it's about. If the transaction rolls back, the job never existed. With a separate queue, you can save an order and lose its "send receipt" job, or send a receipt for an order that rolled back.
- You can query it. It's a table. "Which jobs failed today, and why?" is a
SELECT. - It's fast enough for most apps: thousands of jobs a minute on ordinary hardware.
The table
create table jobs (
id bigint generated always as identity primary key,
kind text not null, -- what to run: 'send_email', 'resize_image'...
payload jsonb not null default '{}', -- its arguments
status text not null default 'queued'
check (status in ('queued', 'running', 'done', 'failed')),
run_at timestamptz not null default now(), -- not before this
attempts int not null default 0,
max_attempts int not null default 5,
locked_by uuid, -- the worker holding it
last_error text,
created_at timestamptz not null default now()
);
-- Workers only ever look for jobs that could run now.
create index jobs_ready_idx on jobs (run_at) where status in ('queued', 'running');
The column doing the most work is run_at. It means "not before this time", and it serves
three purposes:
- A new job runs at
now(), or later if you schedule it. - A failed job waiting to retry gets a
run_atin the future (backoff). - A running job gets a
run_atat the end of its lease. If the worker dies, the job becomes claimable again when the lease runs out (more on that below).
The index is partial: finished jobs don't bloat it, so it stays small however long the table gets.
Claiming a job with FOR UPDATE SKIP LOCKED
The core problem: several workers poll the same table. Two of them must never run the same job, and they shouldn't wait on each other either.
SELECT ... FOR UPDATE locks the rows it returns. SKIP LOCKED (Postgres 9.5+) tells the
query to skip rows another transaction has already locked, instead of waiting for them. Put
that in a subquery, and claim the job in one statement:
update jobs
set status = 'running',
attempts = attempts + 1,
locked_by = $1, -- this worker's id
run_at = now() + interval '5 minutes' -- the lease
where id = (
select id
from jobs
where status in ('queued', 'running')
and run_at <= now()
order by run_at
limit 1
for update skip locked
)
returning *;
One round trip, one atomic statement. If two workers run it at the same moment, the second one skips the row the first one locked and takes the next job. If nothing is due, it returns no rows.
To claim a batch, change limit 1 to limit 10 and where id = (...) to where id in (...).
Why status in ('queued', 'running')? Because a running job whose lease has expired is
claimable again. That's the crash recovery, built into the same query.
Leases: surviving crashed workers
A worker can die mid-job: a deploy, an out-of-memory kill, a lost connection. A queue that
only marks jobs running leaks those jobs forever. My own first version had exactly this
gap: a deploy in the middle of a job would have left it running, with nothing running it.
The lease fixes it. Claiming a job sets run_at to "now plus five minutes". If the worker
finishes, it marks the job done. If it dies, nothing marks it, the lease expires, and the
claim query picks the job up again.
Two rules make this safe:
- The lease must outlast the job. If a job can take ten minutes, a five-minute lease
lets a second worker start it while the first is still going. For long jobs, have the
worker extend its lease as it goes (
update jobs set run_at = now() + interval '5 minutes' where id = $1 and locked_by = $2). - Only the lease holder may finish the job. That's what
locked_byis for. If a slow worker's lease expired and another worker took the job over, the slow one must not mark it done:
update jobs
set status = 'done', locked_by = null
where id = $1 and locked_by = $2; -- 0 rows: we lost the lease; someone else has it
The worker loop
The worker is a loop: claim, run, record the outcome. Here it is in TypeScript with the
postgres driver. Any language looks the same.
import postgres from 'postgres';
const sql = postgres(process.env.DATABASE_URL!);
const workerId = crypto.randomUUID();
type Job = { id: number; kind: string; payload: unknown; attempts: number; max_attempts: number };
const handlers: Record<string, (payload: any) => Promise<void>> = {
send_email: async (payload) => { /* ... */ },
};
async function claim(): Promise<Job | undefined> {
const [job] = await sql<Job[]>`
update jobs
set status = 'running', attempts = attempts + 1, locked_by = ${workerId},
run_at = now() + interval '5 minutes'
where id = (
select id from jobs
where status in ('queued', 'running') and run_at <= now()
order by run_at limit 1
for update skip locked
)
returning *`;
return job;
}
async function work() {
for (;;) {
const job = await claim();
if (!job) {
await Bun.sleep(1000); // nothing due: wait a moment (or LISTEN, below)
continue;
}
try {
await handlers[job.kind]!(job.payload);
await sql`update jobs set status = 'done', locked_by = null
where id = ${job.id} and locked_by = ${workerId}`;
} catch (error) {
await fail(job, error as Error);
}
}
}
Run as many of these as you like, in one process or many. SKIP LOCKED keeps them out of
each other's way.
Retries with backoff (and when not to retry)
Not every failure deserves another try. A timeout, a 429 or a 5xx from an API will probably pass. A validation error won't, and retrying it only wastes time and rate limit. Decide which is which, and retry only the first kind:
function isRetryable(error: unknown): boolean {
const status = (error as { status?: number }).status;
if (status !== undefined) return status === 429 || status >= 500;
return /timed? ?out|ECONNRESET|temporar/i.test(String(error));
}
async function fail(job: Job, error: Error) {
const again = isRetryable(error) && job.attempts < job.max_attempts;
// Exponential backoff with jitter: about 2s, 4s, 8s... capped at 10 minutes.
const delay = Math.min(2 ** job.attempts, 600) * (0.5 + Math.random());
await sql`
update jobs
set status = ${again ? 'queued' : 'failed'},
run_at = now() + ${delay} * interval '1 second',
last_error = ${error.message},
locked_by = null
where id = ${job.id} and locked_by = ${workerId}`;
}
The jitter (0.5 + Math.random()) matters when many jobs fail at once, for example when an
API goes down. Without it they all retry at the same instant and knock it over again.
Jobs that run out of attempts stay in the table as failed, with their last error. That's
your dead-letter queue, and you can query it:
select kind, last_error, count(*)
from jobs
where status = 'failed' and created_at > now() - interval '1 day'
group by 1, 2
order by 3 desc;
To retry one by hand: update jobs set status = 'queued', run_at = now(), attempts = 0 where id = $1.
Enqueueing, and not enqueueing twice
Enqueueing is an insert, ideally inside the transaction that made the job necessary:
await sql.begin(async (tx) => {
const [order] = await tx`insert into orders ${tx(order)} returning id`;
await tx`insert into jobs (kind, payload) values ('send_receipt', ${{ orderId: order.id }})`;
});
Often you also want at most one pending job per thing: one "reindex this course", not five because someone clicked five times. Add a dedupe key and a partial unique index over the jobs that are still pending:
alter table jobs add column dedupe_key text;
create unique index jobs_pending_dedupe_idx
on jobs (kind, dedupe_key)
where status in ('queued', 'running');
insert into jobs (kind, payload, dedupe_key)
values ('reindex_course', '{"courseId": 42}', 'course:42')
on conflict (kind, dedupe_key) where status in ('queued', 'running') do nothing;
The second click is a no-op while the first job is pending. Once it's done, a new one can be queued.
Waking workers with LISTEN/NOTIFY
Polling every second is simple and fine for most apps, but it adds up to a second of delay and a steady trickle of queries. Postgres can tell workers when there's work instead:
create function jobs_notify() returns trigger as $$
begin
perform pg_notify('jobs', new.kind);
return new;
end;
$$ language plpgsql;
create trigger jobs_notify after insert on jobs
for each row execute function jobs_notify();
let wake = () => {};
await sql.listen('jobs', () => wake());
// In the worker loop, instead of a fixed sleep:
await Promise.race([
new Promise<void>((resolve) => (wake = resolve)),
Bun.sleep(30_000), // still poll now and then: retries and scheduled jobs don't notify
]);
Three things to know:
- Notifications are sent on commit. A job inserted in a transaction that rolls back never wakes anyone.
- Keep polling as a fallback. Notifications aren't stored. A worker that was
reconnecting misses them, and jobs whose
run_atarrives (retries, scheduled jobs) don't send one. LISTENneeds a dedicated connection. It doesn't work through PgBouncer in transaction pooling mode. Connect the listener directly to Postgres.
Limiting concurrency and rate
Sometimes the limit isn't your workers, it's what they call. My video jobs call an API that allows about one request a second and gets overloaded with too many videos processing at once. Two cheap ways to respect limits like that:
- Global concurrency: count before claiming, and only claim while there's room:
`select count(*) from jobs where kind = 'process_video' and status = 'running' and run_at
now()
. With several workers, take a transaction-scoped advisory lock first (select pg_advisory_xact_lock(hashtext('process_video'))`) so they don't all see room at once. - Rate: for a single worker, sleep between calls. For several, store the time of the
last call in a one-row table and update it with
select ... for update, so one worker at a time decides who may go next.
And before an expensive call, check whether it's still needed. A job that waited ten minutes for a retry may find the work already done by something else in the meantime.
Keeping the table healthy
A queue table sees a lot of updates, and every Postgres update writes a new row version. The old versions are dead tuples until vacuum cleans them up. Keep it in check:
- Delete or archive finished jobs. A nightly
delete from jobs where status = 'done' and created_at < now() - interval '7 days'keeps the table small. Keep failed ones longer. - Let autovacuum keep up. For a busy queue, make it run more often on this table:
alter table jobs set (autovacuum_vacuum_scale_factor = 0.01). - Keep the index partial, as above. The workers' index then only holds pending jobs, however many finished ones the table holds.
Things to get right in your handlers
- At least once, not exactly once. A worker can finish the work and die before marking the job done; the lease brings it back and it runs again. Write handlers that are safe to run twice: check before acting, use idempotency keys with external APIs, upsert instead of insert.
- Keep jobs small. One job per email, per image, per video. Big jobs hold leases for long, retry expensively and make failures harder to read.
- Store ids, not data, in the payload. Load the current state when the job runs; the data may have changed since it was queued.
When not to use Postgres as a queue
Postgres is a good queue when:
- jobs number in the thousands a minute, not hundreds of thousands a second;
- each job does real work (an HTTP call, a file, an email), so the queue isn't the bottleneck;
- you already run Postgres, and want enqueueing to be part of your transactions.
Use a dedicated broker (Redis/BullMQ, RabbitMQ, SQS, Kafka) when you need very high throughput, fan-out to many consumers, or millisecond latency at scale, or when the queue's load would compete with your main database.
Libraries that do this for you
If you'd rather not write it yourself, these use exactly this pattern:
- Node.js: pg-boss, Graphile Worker
- Go: River
- Elixir: Oban
- Ruby on Rails: Solid Queue
- Postgres extension: pgmq
Writing your own is still worth it when your needs are small and specific, like mine: about
a hundred lines, no new dependency, and a table you can read in the admin page. Either way,
knowing how SKIP LOCKED, leases and backoff fit together makes any of them easier to run.

Comments
Sign in to join the conversation. It's free, with just your email.