Distributed systems

Dual write: the database committed but the event was never published

Written and reviewed by Sahil Srivastav

ConsistencyOutbox patternEvent-driven
2026-09-14T09:31:02.417Z INFO  OrderService  order ORD-90117 committed (tx 4412987)
2026-09-14T09:32:02.911Z ERROR OrderService  failed to publish order-events for ORD-90117
  org.apache.kafka.common.errors.TimeoutException: Expiring 1 record(s) for order-events-3:60000 ms has passed since batch creation

# the row exists
select id, status from orders where id = "ORD-90117";  ->  ORD-90117 | CONFIRMED

# the event does not: downstream never shipped it
$ kafka-console-consumer.sh --topic order-events --from-beginning | grep ORD-90117
(no output)

What this error actually means

Two independent systems, two independent commits, no shared transaction. Whatever order you choose, there is a window in which one has accepted the write and the other has not, and a crash or a timeout inside that window leaves the two permanently disagreeing. No amount of careful ordering removes the window; ordering only decides which kind of inconsistency you get.

Commit the database first and publish second, and a publish failure gives you an order that exists but was never announced — nothing ships, no email is sent, and the downstream read model simply lacks the record. Publish first and commit second, and a rollback gives you an event for an order that does not exist; consumers then act on a phantom, and the harder failure mode is that they may create real, non-reversible effects from it.

The reason the timeout in the log above is so dangerous is that a publish timeout is ambiguous. The broker may well have persisted the record and only the acknowledgement was lost. So a retry can double-publish and an abandonment can lose the event — from the producer’s position those cases are indistinguishable, which is why "just retry the publish" is not a fix.

The pattern that removes the window is to stop writing to two systems. Write the business row and a row describing the event into the same database transaction, so they commit or roll back together, then let a separate process move outbox rows to the broker. Publication becomes a retriable, at-least-once activity with a durable record of intent — and the only remaining requirement is that consumers tolerate duplicates, which they must anyway.

Causes, most common first

  1. 1Publish inside the transaction, commit afterwards. The `send()` happens while the transaction is still open, so any rollback — a constraint violation later in the method, an optimistic-lock failure, a retry by the framework — leaves an event describing a state that never existed. This ordering feels safer and is strictly worse.
  2. 2Publish after commit, with the failure logged and ignored. The common shape: `@Transactional` method commits, a listener publishes, the publish throws, the exception is caught and logged because throwing now would be pointless — the transaction is already committed. The event is lost by design.
  3. 3Publish in an after-commit hook that the process does not outlive. Spring’s `TransactionalEventListener(phase = AFTER_COMMIT)` and equivalents narrow the window but do not close it. A pod terminated between the commit and the asynchronous send loses the event, and rolling deploys make that a routine event rather than a rare one.
  4. 4Blind retry of an ambiguous publish. A timeout is retried and the original batch had in fact been persisted, so the event is now duplicated. Without producer idempotence and without a consumer-side key, you have converted a possible loss into a certain duplicate.
  5. 5Distributed transactions assumed rather than configured. Teams sometimes believe an XA or two-phase commit is in play because a transaction manager is present. Kafka does not participate in XA, and even where a broker does, the operational cost and the blocking-coordinator failure mode are why almost nobody runs it. If you did not deliberately set it up, you are doing a dual write.

When you see it

  • A record exists in the database with no corresponding downstream effect, and no error that anyone noticed
  • Publish failures logged and swallowed inside a `catch` block
  • Consumers occasionally receive an event whose primary key does not exist when they query back
  • Counts diverge slowly: the source table grows faster than the derived table
  • The gap appears in clusters during broker maintenance, network blips, or deploys
  • Manual "republish" scripts exist in the repository, which is the organisational symptom of this bug

How to diagnose it

Step 1

Quantify the divergence before changing anything

You need the size and the shape of the gap: which rows exist upstream with no downstream counterpart, and over what period. This query is also the backfill input, so write it once and keep it.

SELECT o.id, o.created_at
FROM orders o
LEFT JOIN order_projection p ON p.order_id = o.id
WHERE o.created_at < now() - interval '15 minutes'
  AND p.order_id IS NULL
ORDER BY o.created_at;

Step 2

Find the swallowed publish failures

They are almost always in the logs and almost never alerted on. Search for the producer exception types rather than for your own wording, because the log message varies and the exception class does not.

grep -E 'TimeoutException|NotLeaderOrFollowerException|RecordTooLargeException|failed to publish' app.log | wc -l

Step 3

Determine which direction the inconsistency runs

Rows without events means publish-after-commit is dropping them. Events without rows means publishing happens before or inside the transaction. The fix differs, and consumers that already acted on phantom events need a compensation plan, so establish this first.

Step 4

Check whether producer acknowledgement is even strict enough to notice

With `acks=0` or `acks=1` a publish can be reported successful and still be lost to a leader failover. If the producer configuration is lax, some of your missing events were never failures in your logs at all.

grep -E 'acks|enable.idempotence|delivery.timeout.ms' src/main/resources/application*.yml

The fix

Adopt the transactional outbox. In the same transaction as the business change, insert a row into an `outbox` table containing the aggregate id, event type, payload, and a creation timestamp. The transaction is the atomicity boundary, so there is no longer any window where one side committed and the other did not.

Publish from the outbox with a separate relay: either a poller that claims rows with `FOR UPDATE SKIP LOCKED`, publishes, and marks them sent, or change data capture reading the write-ahead log (Debezium is the usual choice) so no polling is needed at all. Either way publication is retriable because the intent is durable.

Accept that the relay is at-least-once and make consumers idempotent. A crash after publishing and before marking the row sent republishes the event. That is the correct trade: duplicates you can deduplicate, versus losses you cannot detect.

Preserve per-aggregate order by keying the published record on the aggregate id and having the relay process rows in id order per aggregate. An outbox drained by a multi-threaded relay with no key affinity reintroduces the out-of-order problem you may have just fixed elsewhere.

Set the producer to `acks=all` with `enable.idempotence=true` and a finite `delivery.timeout.ms` so the relay learns about failures reliably and its own retries do not reorder or duplicate within the broker. Then alert on outbox lag — the age of the oldest unsent row — which is the one metric that makes this whole mechanism observable.

// Dual write: two commits, no atomicity, and a catch that loses the event.
@Transactional
public void confirm(String orderId) {
    orders.updateStatus(orderId, CONFIRMED);
    try {
        kafka.send("order-events", orderId, payload(orderId)).get();
    } catch (Exception e) {
        log.error("failed to publish {}", orderId, e);   // silently divergent
    }
}

// Outbox: one transaction, one atomicity boundary.
@Transactional
public void confirm(String orderId) {
    orders.updateStatus(orderId, CONFIRMED);
    jdbc.update("""
            INSERT INTO outbox (id, aggregate_id, type, payload, created_at)
            VALUES (?, ?, "order.confirmed", ?::jsonb, now())
            """, UUID.randomUUID(), orderId, payload(orderId));
}

-- the relay: concurrent-safe claim, ordered per aggregate, retriable
WITH claimed AS (
  SELECT id, aggregate_id, type, payload
  FROM outbox
  WHERE sent_at IS NULL
  ORDER BY created_at
  FOR UPDATE SKIP LOCKED
  LIMIT 100
)
SELECT * FROM claimed;
-- publish, then:
UPDATE outbox SET sent_at = now() WHERE id = ANY($1);

How to stop it coming back

  • Make "no publish inside a service method that also writes the database" a review rule; the outbox insert is the only allowed form
  • Alert on the age of the oldest unsent outbox row rather than on publish exceptions — it catches a stopped relay, a stuck row, and a poison payload alike
  • Keep a permanent reconciliation job comparing source counts with derived counts; the outbox makes divergence rare, not impossible
  • Delete or archive sent outbox rows on a schedule, and index `(sent_at, created_at)` so the claim query stays cheap as the table ages
  • Test the relay’s crash window explicitly: kill it between publish and mark-sent and assert the consumer end state is unchanged

Practise this failure in a real repository

The Gronex challenge starts from a service that commits the row and then publishes, with tests that inject a broker failure at the worst moment and assert that no confirmed order is left unannounced. It also kills the relay mid-flight, so the solution has to be genuinely at-least-once rather than hopeful.

FAQ

Can I avoid the outbox by publishing first and only committing if the publish succeeds?

No. That ordering trades a lost event for a phantom event: if the commit then fails, consumers have already been told about a state that does not exist, and any real-world effect they produced is not reversible by your rollback. It is the strictly more dangerous of the two orderings.

Is change data capture better than a polling relay?

It is lower latency and puts no query load on the table, at the cost of operating a connector and its own failure modes. A poller with `SKIP LOCKED` is entirely respectable at moderate volumes and is far less to run. Both are correct; choose by operational appetite, not by fashion.

Do I still need idempotent consumers with an outbox?

Yes, and that is not a weakness of the pattern. The relay can crash after publishing and before marking the row sent, so republication is expected. The outbox converts an undetectable loss into a detectable duplicate, which is the whole point.

Will Kafka transactions solve this?

Only for Kafka-to-Kafka flows, where the consumed offsets and the produced records commit atomically inside Kafka. A transaction spanning Kafka and PostgreSQL does not exist, so the moment your source of truth is a database you are back to the outbox.

Related

Other errors engineers hit next to this one

Full error and symptom index →