The outbox is not a queue
The constraint arrived before the design did. The service that writes user activity — posts, follows, playlist changes — runs as a Cloudflare Worker, and a Worker cannot open a raw TCP socket. Not with a library, not with a shim. The Kafka wire protocol and CQL are both permanently unreachable from that process.
The first instinct is to route around it: put an HTTP shim in front of Kafka, have the Worker POST to it, done. That works, and it is wrong, because it moves the interesting failure into a place where you cannot see it. The post is written to Postgres and the event is POSTed to the shim, and those are two operations that can disagree. The user's post exists and their followers never hear about it, or the reverse. You have built a distributed transaction and declined to admit it.
What the outbox actually buys
So the event goes into a Postgres table, inside the same transaction as the post itself:
BEGIN;
INSERT INTO posts (...) VALUES (...);
INSERT INTO outbox (topic, key, payload) VALUES (...);
COMMIT;
Either both rows exist or neither does. There is no window. That is the whole trick, and it is worth being precise about what it does and does not give you: it makes the decision to publish atomic with the write. It does not make the publish atomic with the write, and nothing can.
A host-resident Go relay polls the outbox, produces to a six-partition topic keyed by actor, and marks the row shipped. It runs on a real host, so it has real sockets. The Worker never talks to Kafka at all.
The part people skip
Here is where most write-ups stop, and where the actual engineering starts. The relay can produce to Kafka and die before it marks the row shipped. On restart it produces the same event again. This is not an edge case to be engineered away — it is the permanent, load-bearing condition of the system. The pipeline is at-least-once and there is no configuration that changes that.
Which means every consumer has to be able to see an event twice and do the right thing. Three consumer groups read this topic, and each pays that cost differently:
- The feed materializer writes to Cassandra. Cassandra
INSERTis an upsert, so writing the same feed row twice is genuinely idempotent — the second write produces the identical row. It pays nothing. This is luck arising from the data model, not a design achievement, and it should be described that way. - The notification worker writes to Postgres and must not send two push notifications for one follow. It carries a unique constraint on the event id and swallows the conflict. It pays a table constraint and an index.
- The realtime fanout pushes to connected WebSocket clients. It keeps a short-TTL set of recently-delivered event ids in Redis. It pays memory, and it pays correctness at the edges: if the TTL expires before a duplicate arrives, a user sees a duplicate. That window is a real defect with a known size, and knowing its size is the point.
Three sinks, three different idempotency mechanisms, three different costs. None of them is "exactly-once", and the system does not claim to be anywhere.
Why the shape matters more than the components
Swap Kafka for Redis Streams or NATS and almost none of the above changes. The outbox is not about Kafka. It is about refusing to have two sources of truth about whether something happened. Once the event is committed with the post, every downstream problem becomes a delivery problem, and delivery problems have known solutions.
The version people usually build instead — write the row, then publish, and handle the failure with a retry — is a system where "did this happen" has two answers depending on which service you ask. That is not a harder version of the same design. It is a different design, and a worse one.
The full case study, with the code paths and the commit, is the activity graph.