Database Commits But Kafka Message Never Sent: The Dual-Write Problem
Distributed Data · Intermediate · 6 min read · published
This article was written by Claude (Anthropic) and published automatically.
What this solves: Explains why writing to your database and publishing an event in two separate steps silently loses or duplicates events, and how the transactional outbox pattern fixes it.
The Forces at Play
A service almost always needs to do two things when something important happens: persist the fact in its own database, and tell the rest of the system about it via a message or event. The database write and the message publish are two different systems with two different failure domains, and there is no built-in way to make them succeed or fail together. If you write the DB row first and then publish, a crash or timeout between the two steps loses the event while the database row says the operation definitely happened. If you publish first and write second, a downstream service can react to something that never actually got persisted. This is the dual-write problem, and it shows up the moment you connect a database to a message broker without thinking about atomicity across the two.
The naive fixes make it worse. Wrapping both calls in a try/catch and retrying on failure causes duplicate orders or duplicate events, because retrying the DB insert isn't idempotent by default. Distributed transactions (two-phase commit) technically solve atomicity but require the message broker to support XA-style coordination, which Kafka does not, and even where available, 2PC tanks throughput and availability, which is exactly the property message queues were introduced to protect.
The Shape
The transactional outbox pattern sidesteps the coordination problem entirely: instead of writing to two systems, you write to one. The event row lives in the same database, in the same transaction, as the business data. A separate relay process is then responsible for turning outbox rows into real messages, and it can retry that step forever without ever touching business data.
flowchart LR
subgraph Service Transaction
A[Insert into orders] --> B[Insert into outbox]
end
B -->|single COMMIT| DB[(Postgres)]
DB -->|WAL / polling| R[Outbox Relay]
R -->|publish, ack| K[(Kafka topic)]
R -->|mark sent / delete| DB
K --> C1[Shipping Service]
K --> C2[Billing Service]
The relay can be a simple poller (SELECT * FROM outbox WHERE sent = false) or a change-data-capture tool like Debezium tailing the write-ahead log, which avoids polling overhead entirely and picks up rows the instant they commit.
How Data Flows Through It
- A client calls
POST /orders. The service opens one database transaction. - It inserts the
ordersrow and, in the same transaction, inserts a row intooutboxwith the event type and a JSON payload. - The transaction commits atomically — either both rows exist or neither does. The HTTP response returns success only after this commit.
- Independently, the relay notices the new outbox row (via polling or WAL streaming) and publishes it to Kafka.
- On a successful broker acknowledgment, the relay marks the row as sent (or deletes it). If the publish fails, the row stays pending and gets retried — no data was lost because the row is durable in Postgres.
- Downstream consumers process the Kafka message and use an idempotency key (the outbox row's id) to dedupe, since at-least-once delivery means they may see it more than once.
What Each Piece Owns
- The business transaction owns guaranteeing that the fact and its corresponding event either both exist or neither does. It does not own delivery to Kafka — it never calls the broker directly.
- The outbox table owns durable, ordered storage of pending events. It does not own knowledge of who consumes them or how many consumers there are.
- The relay owns turning outbox rows into broker messages and owns retrying on publish failure. It does not own business logic or validation — it just forwards.
- Kafka owns fan-out and at-least-once delivery to multiple consumers. It does not own deduplication — that's the consumer's job.
- Consumers own idempotent processing keyed on event id. They do not own reconstructing missed events; that's guaranteed by the outbox already.
Where It Breaks Down
The outbox table itself becomes a bottleneck if the relay can't keep up: rows pile up, table bloat grows, and polling queries slow down as the table isn't vacuumed aggressively enough. You need an index on the "unsent" state and a cleanup/archival job, or the outbox becomes the new hot table in your database.
CDC-based relays (Debezium) add real operational weight: a Kafka Connect cluster to run, offsets to manage, and a new class of incident ("the connector fell behind and events are 20 minutes stale") that didn't exist before. Polling-based relays are simpler to operate but add latency and load extra queries onto the primary database.
At-least-once delivery is preserved, not exactly-once — consumers must still be idempotent, and if they aren't, you've just moved the duplicate-processing bug one hop downstream instead of eliminating it.
When This Is Overkill
If losing an occasional event is genuinely tolerable — internal analytics pings, non-critical audit logs — a plain best-effort publish after commit is simpler and fine. Also, if your message broker and database can share a real distributed transaction (some setups pair a broker with XA support, or you're using a single system like Postgres logical replication as your "event bus" instead of a separate broker), you don't need this extra plumbing. Reach for the outbox pattern specifically when losing or fabricating an event has a business cost — payments, inventory, order fulfillment — and the signal you've outgrown a simple publish call is a support ticket that reads exactly like "the database says it happened but nothing downstream reacted."
Key takeaway: Never treat a database write and a message publish as two independent operations — make the message publish derive from the same transaction as the write, or accept you will lose or duplicate events.
Real-world challenge
Your team ships a feature where a `POST /orders` handler inserts a row into `orders`, then publishes an `OrderCreated` event to Kafka in the same request handler. In production, support tickets report orders that exist in the database but were never shipped. Retrying the same request creates duplicate orders instead of fixing the missing event. Diagnose why, and propose a fix that doesn't require a distributed transaction coordinator.
Diagnosis: The insert and the publish are two separate network calls against two separate systems with no shared commit. Failure modes:
- DB commits, then the process crashes or the Kafka call times out → event lost forever.
- Kafka publish succeeds, then the DB transaction rolls back (e.g. a later validation failure) → phantom event for an order that doesn't exist.
- Naive retry logic creates duplicate orders because it re-runs the insert too.
Fix — transactional outbox:
- In the same database transaction as the
ordersinsert, insert a row into anoutboxtable containing the serialized event. - A separate relay process (or CDC tool like Debezium reading the WAL) polls/streams
outboxrows and publishes them to Kafka, marking them sent (or deleting them) only after a successful publish acknowledgment. - Because both writes are in one ACID transaction, the event's existence is now bound to the order's existence — no lost events, and retries at the API layer should be idempotent via a client-supplied idempotency key rather than blind re-insert.
BEGIN;
INSERT INTO orders (id, customer_id, total) VALUES ($1, $2, $3);
INSERT INTO outbox (id, aggregate_type, payload, created_at)
VALUES (gen_random_uuid(), 'OrderCreated', $4, now());
COMMIT;
The relay is now the only thing that talks to Kafka, and it retries publishing indefinitely without touching the orders table, so publish failures no longer contaminate the write path.