An AI workflow usually needs a queue before it needs anything clever. A form submission has to be enriched, classified and written to a CRM; a document has to be extracted and checked; a batch of prompts has to run overnight. The reflex is to add a message broker. If the application already runs on Postgres, the database can often do the job, and it closes a gap that a separate broker opens.
The gap between saving and scheduling
The common first design saves a record, then publishes a message to a broker:
submission = Submission.objects.create(**data) # 1. commit to Postgres
broker.publish("process", submission.id) # 2. tell the worker
return thank_you()
If the process dies between step 1 and step 2, the submission exists and nobody will ever process it. If you swap the order, a worker can pick up an ID whose row hasn't committed yet. Either way there is a window where the database and the queue disagree, and under load you'll hit it.
The fix is to make "there is work to do" part of the same transaction as the work itself.
A jobs table is a Postgres queue
CREATE TABLE jobs (
id bigserial PRIMARY KEY,
kind text NOT NULL,
payload jsonb NOT NULL,
state text NOT NULL DEFAULT 'pending',
attempts int NOT NULL DEFAULT 0,
run_after timestamptz NOT NULL DEFAULT now(),
last_error text
);
CREATE INDEX jobs_ready ON jobs (run_after) WHERE state = 'pending';
The request handler inserts the submission and its job in one transaction. Both commit or neither does, so there's no window to fall into.
Workers claim jobs like this:
UPDATE jobs SET state = 'running', attempts = attempts + 1
WHERE id = (
SELECT id FROM jobs
WHERE state = 'pending' AND run_after <= now()
ORDER BY run_after
FOR UPDATE SKIP LOCKED
LIMIT 1
)
RETURNING id, kind, payload;
FOR UPDATE SKIP LOCKED is what makes this safe with several workers. A row another worker has locked is skipped instead of waited on, so workers don't block each other or take the same job. The PostgreSQL documentation describes exactly this use: skipping locked rows gives an inconsistent view, unsuitable for general queries, but useful for avoiding contention between consumers of a queue-like table (SELECT, the locking clause).
When the job succeeds, the worker marks it done. When it fails with a transient error, it sets state = 'pending' and pushes run_after into the future with backoff. When attempts run out or the error is permanent, it marks the job failed and records why. A failed job is a row you can query, count and retry, which beats a message that vanished from a dead-letter topic nobody watches.
The outbox: talking to the outside world
Jobs cover internal work. For messages to other systems, the same idea is called the transactional outbox: write the outgoing event to an outbox table in the same transaction as the business change, then have a relay read the table and deliver the events. The business change and the intent to notify can't drift apart.
The relay still delivers at least once, because it can crash after sending and before marking the row sent. That's why every outgoing call needs an idempotency key, usually the outbox row's ID or a stable business key, and the receiving side, or your adapter, has to recognize a repeat. When the destination offers no idempotency, read back before retrying: search for the record you meant to create before creating it again.
What an AI workflow gets from this
- Model calls are slow and fail in bursts. Backoff with
run_afterhandles rate limits without a separate scheduler. - Every attempt can be a row in an
attemptstable with the prompt version, model and outcome. That's your audit trail and your evaluation data. - A job waiting on a person is just a state. Review queues are a
WHERE state = 'held'away. - Replays are inserts. To reprocess yesterday's failures with a new prompt, insert new jobs for them.
Where it stops being enough
Postgres queues handle more than most teams expect, but they have limits. Very high job rates add write and vacuum pressure to your main database. Long-running jobs that hold a transaction open are a problem, so claim the job, commit, do the work, then update the row in a new transaction. Fan-out to many independent consumers, cross-service event streams and replayable logs are what brokers such as Kafka are built for. When you get there, keep the outbox and point its relay at the broker; you still want the transactional write.
Libraries implement all of this with more care than a weekend version: Procrastinate for Python, pg-boss for Node and River for Go. Read one of them before writing your own, even if you end up writing your own.
The rule I use
If the work starts from a database write, the job that finishes it should be written in the same transaction. The lead intake design study walks through this design with a CRM on the other end, and the durable execution pattern covers the general shape. Everything else, including which queue technology to use, is a question of scale you can answer later with measurements.
The practical next step
Map one real execution and one failure.
That will reveal more about the right architecture than a tool comparison or model demo.
Let's build something real