Work on a distributed system long enough and you'll know the particular cold sweat of the dual write. You insert an order into the database and the transaction commits. Then you have to tell something outside the database, a service over HTTP or a Kafka topic, and the connection drops before you can. The order exists and nobody downstream knows. Or it goes the other way: the message goes out first, the transaction rolls back, and a consumer is now holding an event for an order that was never created, a ghost.
No amount of try/catch will make those two writes atomic.
Think of paying at a restaurant. The payment is your database transaction and the receipt is the message to the outside world. Walk out without the receipt and you can't prove the meal happened. The transactional outbox pattern is the industry's answer to that gap. You don't collect the receipt at the door, you put it in your pocket the moment you pay: the message goes into an outbox table in the same transaction as the business row, so both commit or neither does.
That solves the dual write, and creates the problem this post is about.
The relay is a second system
Someone still has to take the receipt out of your pocket and show it to the world. In the usual design that's a relay, a separate service outside the database that polls the outbox table, delivers each row and marks it done.
The relay looks small on the whiteboard and turns out large in production. It has its own deployment, configuration, connection string and alerts. It has to claim rows without two instances grabbing the same one, retry, back off, give up, and put whatever it gave up on somewhere a person can look. And since it lives outside the database, its view of the world can drift away from the database's, which is how relays end up polling a replica for hours while reporting an empty queue.
I wrote ulak because I kept building that relay, and every time it was the least interesting and most fragile part of the system. ulak is a C extension for PostgreSQL 14 through 18, under Apache 2.0, that keeps the outbox pattern and gets rid of the relay. The queue is a table you write to in your own transaction, and PostgreSQL background workers do the delivery.
What it is not
ulak isn't an event streaming platform, and it doesn't replace change data capture. If you're moving hundreds of thousands of log events a second, use Kafka. If you need every change in a database mirrored into a log, use Debezium. ulak's job is narrower. For a system whose core already lives in PostgreSQL, it delivers outbox messages locally and transactionally, without a broker.
Enqueue is your transaction
On the application side it's one function call inside the transaction you already have:
BEGIN;
INSERT INTO orders (customer_id, total_cents) VALUES (42, 12990);
SELECT ulak.send(
'billing-webhook',
jsonb_build_object('order_id', 17, 'total_cents', 12990)
);
COMMIT;
ulak.send inserts a row into ulak.queue, and that row shares its fate with your order. If the next line of your code throws and the transaction rolls back, the message doesn't float off somewhere. It's gone, as if it never existed. If the transaction commits, the message is durable before anything has tried to deliver it. The receipt is in your pocket.
Delivery is a PostgreSQL process
This is where the design gets unusual, and a bit bold. The workers that drain the queue are PostgreSQL background workers, registered through shared_preload_libraries when the server starts, so they live and die with the server. There's no relay to deploy, and when the database fails over the workers go with it, because they're part of it.
Most engineers' first reaction is the right one. PostgreSQL is the heart of the system and the only guardian of the data, and putting a library inside it that runs its own loops and makes its own network calls sounds risky. What if a bug in a worker locks the database or eats its memory?
That worry is the central trade-off of the whole design, and I'm not going to argue it away. You drop the operational cost of a separate service and the network hop between it and the database, and in exchange the delivery pipeline lives and dies with the database. PostgreSQL's background worker infrastructure has matured a lot. Workers are isolated processes with their own memory contexts under PostgreSQL's rules, so a leak is reclaimed when the transaction ends, and a crash doesn't take the postmaster down. It's still a decision that needs a deliberate yes from whoever runs your database.
How workers share one table
Say you run ten workers. The obvious question is what stops all ten charging at the first row in the queue, and ulak answers it in two layers.
The first is arithmetic. Each worker gets a deterministic sequential id at startup, and its poll query has one extra condition: the row's id modulo the worker count has to equal the worker's id. One big table becomes ten separate logical slices, and normally nobody eats off anyone else's plate.
Real life is messier than a formula. Change the worker count while the system is live and the slices shift. Let an operator lock a few rows by hand in a console and a worker could get stuck on them. So the arithmetic is only there to cut contention, and the second layer is what makes it safe: FOR UPDATE SKIP LOCKED. When a worker hits a row someone else holds, it skips it and moves on instead of waiting.
SELECT id, endpoint, payload
FROM ulak.queue
WHERE status = 'pending'
AND scheduled_at <= now()
AND id % 10 = 3 -- this worker's slice
ORDER BY id
LIMIT 200
FOR UPDATE SKIP LOCKED;
▶ Three workers sharing one queue table: an animation that plays in the original post.
There's a third decision underneath. ulak requires the workers to run at READ COMMITTED, PostgreSQL's default, and refuses REPEATABLE READ. The instinct is to reach for the stricter level. But under REPEATABLE READ PostgreSQL uses snapshot isolation, and when several workers use SKIP LOCKED on overlapping rows, a worker that sees a row someone else updated after its snapshot was taken gets serialization failure 40001. The workers would spend their whole lives aborting and retrying. At READ COMMITTED every statement sees the latest committed data, so a locked row is skipped, an unlocked one is claimed, and there's no conflict to raise.
What "delivered" means
Can ulak promise exactly-once delivery? No, and I want to be clear about that, because it's the golden rule of distributed systems: nothing that crosses a network can guarantee exactly-once by itself.
| Step | Guarantee | Why |
|---|---|---|
Writing to ulak.queue
|
Exactly once | It is your transaction |
| Delivering to HTTP, Kafka, etc. | At least once | The network can lose the acknowledgement |
Picture the worst case. ulak sends a message to your payment API. The API processes it and writes it to its own database, and the connection dies the instant before it can say "got it". ulak never sees the acknowledgement, so it does the only thing it can and sends the message again. That's why your consumer has to be idempotent: getting the same message twice mustn't corrupt its data. Correctness doesn't stop at the database. It reaches whoever consumes the message.
The same problem exists on the way in. Your code calls ulak.send, and the database connection drops just as the transaction is finishing. Your code reasonably decides the write failed and retries, and a blind retry would enqueue the same order twice. The fix is an idempotency key you supply yourself:
SELECT ulak.send_with_options(
'billing-webhook',
jsonb_build_object('order_id', 17),
idempotency_key => 'order:17:created'
);
ulak stores the MD5 of the key, not of the payload, and enforces it with a partial unique index on ulak.queue. "Partial" matters here. Over time the queue piles up millions of delivered rows, and an index over all of them would get enormous and slow down every insert. So the index only covers rows whose status is pending or processing. A second send with the same key while the first is still active hits the index, ulak compares the hashes, ignores the new row, and returns the id of the existing message. There's no trip to Redis for a distributed lock, and it gets settled where the data lives.
The trade-off that should bother you
This bothered me most while I was designing it, and it's written up in the 0.0.3 architecture notes so nobody finds it out in production.
Claiming a batch, making the network call and updating the status all happen inside one open PostgreSQL transaction. If the API you're delivering to takes two seconds to answer, the worker's transaction stays open for two seconds. Open transactions hold locks and tie up a connection, and a system doing thousands of deliveries a second against a slow downstream can wedge itself fast.1
Two things keep it in check. The first is at the storage layer: worker transactions run with synchronous_commit = off. Normally PostgreSQL won't treat a transaction as finished until the WAL record is flushed to disk, and after the network the disk is the slowest thing around. With the flush deferred, a worker marks a message delivered and goes straight on to the next one.
Skipping the flush would be reckless for business data. For the queue it's safe, because of the at-least-once model from the previous section. If the server loses power in the millisecond after a worker marks a message delivered and that WAL record is lost, the row comes back as pending on restart, as if it had never been sent. A worker wakes up and sends it again, and since the consumer is idempotent, the duplicate does no harm.
The second is the circuit breaker, in the next section. Together they limit how long a transaction can stay open behind a bad downstream. They don't make the cost go away, though. Timeouts on your endpoints and a modest batch size are your side of the bargain.
One more cost sits in the same place.2 Kafka, MQTT, Redis Streams, AMQP and NATS each need their C client library compiled into the extension, which means optional external C dependencies on your database server.
When the other side is down
If a downstream API dies outright, retrying every message against it at full speed burns connections and achieves nothing. ulak keeps an in-memory circuit breaker per endpoint, with three states. Closed is normal. When consecutive failures cross the threshold the breaker opens, and workers skip every row for that endpoint. Nothing gets sent, the downstream gets a breather, and other endpoints keep flowing.
The breaker shouldn't stay open forever. After a cooldown it goes half-open, and one probe is let through. The subtle part is deciding which worker sends it. Fifty workers scanning the queue can all notice the half-open state at the same moment, and without coordination they'd all probe, which is exactly the stampede the breaker is there to stop.
ulak settles this with a compare-and-swap on the breaker state. Think of one microphone in a room. Fifty workers want to ask whether the other side is back, and only the one that wins the swap gets to speak. It sends the single probe, and the rest see the microphone is taken and keep deferring their rows. If the probe gets a 200, the breaker closes and things carry on. If it fails, the breaker opens again and the cooldown restarts.
▶ A circuit breaker opening, and closing on one probe: an animation that plays in the original post.
Not every failure deserves a retry, so ulak sorts them first. A transient error, like a timeout or a connection reset, is marked retryable and rescheduled with growing backoff. A permanent error, like an HTTP 400 because the payload failed the other side's validation, isn't retried at all, because a million attempts wouldn't change the answer. A permanent error and an exhausted retry budget both move the message to ulak.dlq, the dead letter queue, where it waits for a person.3
When the worker itself dies
There's one more failure, and it's the one people forget. A worker claims a message, marks it processing, and then the process dies mid-request, from a bug or a hardware fault. Now the row is processing forever. Other workers won't touch it, because as far as they can tell someone's working on it.
So ulak gives worker 0 an extra job on top of its normal batches: a periodic crash recovery pass that looks for rows stuck in processing past a timeout and puts them back to pending. The workers that are still alive pick the row up and finish the job.
The bet
Put together, ulak is the outbox pattern with PostgreSQL's own machinery doing the relay's job from inside the database. In return it asks you to accept one deliberate trade: network calls now run inside your database process, and the transaction stays open while they run.
For a decade the microservices consensus has been that databases are dumb storage and the logic belongs in external tools: in Kafka, in RabbitMQ, in a relay. That got taught as a rule. ulak bets that for a database-centric system the rule is backwards, and that the database, the one component that already knows exactly what's been committed, is the right thing to announce it.
Whether that's a bet you should make depends on whether you run your own PostgreSQL, whether your core already lives there, and whether you're willing to hand an old friend that much responsibility again. The code is at github.com/zeybek/ulak. If you run it and it breaks, open an issue. I'd rather know.
Originally published at zeybek.dev.
-
Making HTTP requests from inside a database is playing with fire, and I'm the one who lit it. ↩
-
HTTP is built in, since libcurl is the only hard dependency. ↩
-
And when an HTTP server answers with a
Retry-Afterheader, ulak drops its own backoff calculation and waits exactly as long as it was told to. ↩
Top comments (0)