Transactional Outbox: DB + Kafka Without Losing Messages
Six weeks after our checkout service went live — the same CheckoutService.placeOrder we built in this track's capstone, "Building a Complete Spring Boot API" — the warehouse team asked why a paid order had been sitting in the database for three days with nobody picking it. The order row was perfect: customer email, line items, total, status CONFIRMED. The problem was upstream of the warehouse: the OrderPlaced event that should have told fulfillment about it had never been sent to Kafka. The database commit succeeded; the producer.send() after it never ran, because the pod was killed in the milliseconds between the two.
We had written the code every tutorial teaches: save the order, then publish the event. It looks obviously correct, and it has a hole in the middle you can drive a forklift through. Two weeks later we saw the mirror image — an event reached Kafka, the database transaction rolled back, and fulfillment started preparing an order that did not exist. There is no retry logic that fixes this. The database and the broker never share a transaction, so any sequence of "write here, then write there" can die between the two writes. The fix is to stop treating the event as a second write: make it part of the first one.
The dual-write problem: two systems, no shared transaction
A dual write is any operation that must change two systems that don't share a transaction — here, PostgreSQL and Kafka. Your code can order the two writes either way, and both orders lose:
| Order of writes | Where it dies | What you get |
|---|---|---|
| Save row, then send event | Crash between commit and send() | Order exists in the DB, invisible to every downstream system |
| Send event, then save row | DB rolls back after the send | Fulfillment acts on a ghost order the DB says never existed |
| Send event inside the DB transaction | Still no shared transaction | Worst of both: event sent, then DB rolls back — and you can't unsend |
Decision rule: if a business operation must change the database and notify the outside world, never let application code perform the two writes — collapse them into one transaction, and let a separate process move the notification. That separate process is the relay, and the table it reads is the outbox.
The outbox: the event becomes just another row
The transactional outbox pattern is disarmingly simple: add an outbox table, and write the event into it inside the same local transaction as the business row. The event is no longer a network call that can fail independently — it's data, committed or rolled back with everything else. A relay process then reads unpublished outbox rows and publishes them to Kafka.
Here is the schema. The event_id is a UUID the consumer will use for dedup; published is the relay's cursor; processed_events is the consumer's memory:
CREATE TABLE orders (
id BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY,
customer_email VARCHAR(255) NOT NULL,
total DECIMAL(10,2) NOT NULL,
status VARCHAR(32) NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
CREATE TABLE outbox (
id BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY,
aggregate_type VARCHAR(64) NOT NULL, -- 'Order'
aggregate_id BIGINT NOT NULL, -- the order id
event_type VARCHAR(64) NOT NULL, -- 'OrderPlaced'
event_id VARCHAR(36) NOT NULL UNIQUE, -- UUID, the dedup key
topic VARCHAR(128) NOT NULL, -- 'orders'
payload VARCHAR(4000) NOT NULL, -- the event JSON
published BOOLEAN DEFAULT FALSE, -- the relay's cursor
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
CREATE TABLE processed_events (
event_id VARCHAR(36) PRIMARY KEY, -- consumer-side dedup
processed_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
Three small support types before the main code — the event record, a broker abstraction so the demo runs with zero infrastructure, and its in-memory implementation:
package com.javamakeuse.orders;
/** A published domain event, as read from the outbox table. */
public record OutboxEvent(long outboxId, long aggregateId, String topic, String eventId, String payload) {}
package com.javamakeuse.orders;
import java.util.List;
/** Abstraction over the message broker. The runnable demo uses the in-memory
* implementation; production wires in a real KafkaProducer (see below). */
public interface EventSink {
void send(OutboxEvent event);
List<OutboxEvent> sentEvents();
}
package com.javamakeuse.orders;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
/** Fake broker for the demo: records every published event in memory,
* standing in for a Kafka producer. */
public final class InMemorySink implements EventSink {
private final List<OutboxEvent> sent = new ArrayList<>();
@Override
public synchronized void send(OutboxEvent event) {
sent.add(event);
}
@Override
public synchronized List<OutboxEvent> sentEvents() {
return Collections.unmodifiableList(new ArrayList<>(sent));
}
}
And here is the atomic write, in plain JDBC against H2 (the connection and transaction handling follows the fundamentals from this track's "JDBC First: Connections, Pools & Transactions" post). In the Spring service this same insert rides inside the @Transactional method from "Transactions: @Transactional, Isolation & Propagation" — one annotation, two inserts, one commit:
package com.javamakeuse.orders;
import java.math.BigDecimal;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Statement;
import java.util.UUID;
public final class CheckoutService {
private CheckoutService() {}
public static long placeOrder(Connection conn, String email, BigDecimal total,
boolean failAfterWrites) throws SQLException {
conn.setAutoCommit(false);
try {
long orderId;
try (PreparedStatement ps = conn.prepareStatement(
"INSERT INTO orders (customer_email, total, status) VALUES (?, ?, 'CONFIRMED')",
Statement.RETURN_GENERATED_KEYS)) {
ps.setString(1, email);
ps.setBigDecimal(2, total);
ps.executeUpdate();
try (ResultSet rs = ps.getGeneratedKeys()) {
rs.next();
orderId = rs.getLong(1);
}
}
// THE PATTERN: the event is a row, written by the same transaction.
String eventId = UUID.randomUUID().toString();
String payload = "{\"orderId\":" + orderId
+ ",\"email\":\"" + email
+ "\",\"total\":" + total + "}";
try (PreparedStatement ps = conn.prepareStatement(
"INSERT INTO outbox (aggregate_type, aggregate_id, event_type, event_id, topic, payload, published)"
+ " VALUES ('Order', ?, 'OrderPlaced', ?, 'orders', ?, FALSE)")) {
ps.setLong(1, orderId);
ps.setString(2, eventId);
ps.setString(3, payload);
ps.executeUpdate();
}
if (failAfterWrites) {
// Simulate a downstream failure AFTER both writes, BEFORE commit.
throw new RuntimeException("simulated payment failure AFTER both writes");
}
conn.commit();
return orderId;
} catch (Exception e) {
conn.rollback(); // order row AND outbox row roll back together
if (e instanceof SQLException se) throw se;
throw new RuntimeException(e);
} finally {
conn.setAutoCommit(true);
}
}
}
}The driver that wires it all together and runs the four scenarios — this is the program whose output follows:
package com.javamakeuse.orders;
import java.math.BigDecimal;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Statement;
/** Runnable proof of the transactional outbox pattern:
* (a) order + outbox commit atomically
* (b) on failure both roll back
* (c) the polling relay publishes each row exactly once
* (d) a simulated redelivery is deduped by the consumer */
public final class OutboxDemo {
private static final String DDL = """
CREATE TABLE orders (
id BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY,
customer_email VARCHAR(255) NOT NULL,
total DECIMAL(10,2) NOT NULL,
status VARCHAR(32) NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
CREATE TABLE outbox (
id BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY,
aggregate_type VARCHAR(64) NOT NULL,
aggregate_id BIGINT NOT NULL,
event_type VARCHAR(64) NOT NULL,
event_id VARCHAR(36) NOT NULL UNIQUE,
topic VARCHAR(128) NOT NULL,
payload VARCHAR(4000) NOT NULL,
published BOOLEAN DEFAULT FALSE,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
CREATE TABLE processed_events (
event_id VARCHAR(36) PRIMARY KEY,
processed_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
""";
public static void main(String[] args) throws Exception {
Connection conn = DriverManager.getConnection("jdbc:h2:mem:outboxdb;DB_CLOSE_DELAY=-1");
try (Statement st = conn.createStatement()) {
for (String ddl : DDL.split(";")) {
String trimmed = ddl.trim();
if (!trimmed.isEmpty()) st.execute(trimmed);
}
}
InMemorySink sink = new InMemorySink();
PollingRelay relay = new PollingRelay(conn, sink);
System.out.println("=== (a) Happy path: order + outbox commit atomically ===");
long orderId = CheckoutService.placeOrder(conn, "priya@example.com", new BigDecimal("149.99"), false);
System.out.println("placed order id=" + orderId);
System.out.println("orders rows = " + count(conn, "orders"));
System.out.println("outbox rows = " + count(conn, "outbox"));
System.out.println("outbox unpublished = " + count(conn, "outbox WHERE published = FALSE"));
System.out.println("=> one commit made the business row AND its event visible together");
System.out.println();
System.out.println("=== (b) Failure: downstream throws after both writes ===");
try {
CheckoutService.placeOrder(conn, "arjun@example.com", new BigDecimal("59.99"), true);
System.out.println("ERROR: expected the simulated failure");
} catch (RuntimeException e) {
System.out.println("caught: " + e.getMessage());
}
System.out.println("orders rows = " + count(conn, "orders") + " (still 1 — no ghost order)");
System.out.println("outbox rows = " + count(conn, "outbox") + " (still 1 — no orphan event)");
System.out.println("=> the exception rolled back BOTH writes in one unit of work");
System.out.println();
System.out.println("=== (c) Relay: each unpublished row published exactly once ===");
int first = relay.pollAndPublish(10);
System.out.println("relay run #1 published = " + first);
System.out.println("sink received = " + sink.sentEvents().size());
System.out.println("outbox unpublished = " + count(conn, "outbox WHERE published = FALSE"));
int second = relay.pollAndPublish(10);
System.out.println("relay run #2 published = " + second + " (nothing left to publish)");
System.out.println("sink received = " + sink.sentEvents().size() + " (no duplicates from the relay)");
System.out.println("=> the relay never re-publishes a row it marked");
System.out.println();
System.out.println("=== (d) Redelivery: relay crashes before marking, consumer dedups ===");
// Simulate the crash window: flip the published row back to FALSE,
// as if the relay had sent it but died before the UPDATE committed.
try (Statement st = conn.createStatement()) {
st.executeUpdate("UPDATE outbox SET published = FALSE WHERE id = 1");
}
System.out.println("simulated crash: row 1 reset to unpublished");
int resent = relay.pollAndPublish(10);
System.out.println("relay run #3 published = " + resent + " (at-least-once: the duplicate WAS re-sent)");
System.out.println("sink received = " + sink.sentEvents().size() + " (sink now holds the same event twice)");
int fulfilled = 0, skipped = 0;
for (OutboxEvent e : sink.sentEvents()) {
if (Consumer.handleOrderPlaced(conn, e)) {
fulfilled++;
System.out.println("consumer: fulfilled order from event " + e.eventId().substring(0, 8) + "...");
} else {
skipped++;
System.out.println("consumer: duplicate skipped for event " + e.eventId().substring(0, 8) + "...");
}
}
System.out.println("fulfilled = " + fulfilled);
System.out.println("duplicates skipped = " + skipped);
System.out.println("processed_events rows = " + count(conn, "processed_events"));
System.out.println("=> at-least-once relay + idempotent consumer = exactly-once EFFECT");
}
private static long count(Connection conn, String fromWhere) throws SQLException {
try (Statement st = conn.createStatement();
ResultSet rs = st.executeQuery("SELECT COUNT(*) FROM " + fromWhere)) {
rs.next();
return rs.getLong(1);
}
}
}
I ran this against H2 — Java 21, javac, no Spring, no broker. Scenario (a) places an order normally; scenario (b) throws after both writes but before commit. Genuine output:
=== (a) Happy path: order + outbox commit atomically ===
placed order id=1
orders rows = 1
outbox rows = 1
outbox unpublished = 1
=> one commit made the business row AND its event visible together
=== (b) Failure: downstream throws after both writes ===
caught: java.lang.RuntimeException: simulated payment failure AFTER both writes
orders rows = 1 (still 1 — no ghost order)
outbox rows = 1 (still 1 — no orphan event)
=> the exception rolled back BOTH writes in one unit of work
The failure case is the one that matters. A downstream exception after both writes leaves the database exactly as it was: no ghost order a support agent has to cancel by hand, and no orphan event waiting in the outbox to confuse a relay. Atomicity now covers the intent to notify, not just the business row.
Decision rule: every business operation that must notify the outside world writes its outbox row in the same transaction as the business row — no exceptions, no "we'll add the event later." If the event isn't in the outbox, it didn't happen.
The relay: a polling publisher that never loses a row
The relay is a loop: select unpublished outbox rows oldest-first, publish each to the sink, then mark them published. The demo's sink is an in-memory stand-in for a Kafka producer (a List that records every published event); the production Kafka wiring is shown after the relay. Notice the deliberate ordering inside pollAndPublish: publish first, then mark. If the process dies between the two, the row is still unpublished and the next poll resends it. That makes the relay at-least-once by design — and the consumer's dedup, covered below, is what turns the pair into exactly-once in effect:
package com.javamakeuse.orders;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.util.ArrayList;
import java.util.List;
public final class PollingRelay {
private final Connection conn;
private final EventSink sink;
public PollingRelay(Connection conn, EventSink sink) {
this.conn = conn;
this.sink = sink;
}
/** One polling cycle. Returns the number of events published. */
public int pollAndPublish(int batchSize) throws SQLException {
List<OutboxEvent> batch = new ArrayList<>();
conn.setAutoCommit(true);
try (PreparedStatement ps = conn.prepareStatement(
"SELECT id, aggregate_id, topic, event_id, payload FROM outbox"
+ " WHERE published = FALSE ORDER BY id LIMIT " + batchSize,
ResultSet.TYPE_FORWARD_ONLY, ResultSet.CONCUR_READ_ONLY)) {
try (ResultSet rs = ps.executeQuery()) {
while (rs.next()) {
batch.add(new OutboxEvent(rs.getLong("id"), rs.getLong("aggregate_id"),
rs.getString("topic"), rs.getString("event_id"),
rs.getString("payload")));
}
}
}
if (batch.isEmpty()) {
return 0;
}
// Publish first, then mark. Crash between the two and the row
// stays unpublished: the next poll resends it. At-least-once.
for (OutboxEvent e : batch) {
sink.send(e);
}
conn.setAutoCommit(false);
try (PreparedStatement ps = conn.prepareStatement(
"UPDATE outbox SET published = TRUE WHERE id = ? AND published = FALSE")) {
for (OutboxEvent e : batch) {
ps.setLong(1, e.outboxId());
ps.addBatch();
}
ps.executeBatch();
}
conn.commit();
conn.setAutoCommit(true);
return batch.size();
}
}
Running the relay twice over the one unpublished row from scenario (a) — genuine output:
=== (c) Relay: each unpublished row published exactly once ===
relay run #1 published = 1
sink received = 1
outbox unpublished = 0
relay run #2 published = 0 (nothing left to publish)
sink received = 1 (no duplicates from the relay)
=> the relay never re-publishes a row it marked
A second poll finds nothing because published = TRUE is the cursor. Two honest caveats, because interviews probe exactly these. First: this demo ran one relay instance. With two relay instances racing, both can select the same unpublished row before either marks it — the AND published = FALSE guard in the UPDATE only narrows the window. Production fixes this by claiming rows atomically (UPDATE ... WHERE published = FALSE ... RETURNING, or a lease column with an owner id and expiry). The duplicate that slips through is harmless anyway, because the consumer dedups — which is the point of the next section. Second: polling adds latency equal to your poll interval and a small steady SELECT load; that's the price of the simple version.
Decision rule: the relay owns exactly one job — move unpublished rows to the broker and mark them. It never interprets payloads, never reorders, never filters. All business logic stays in the producer (what to write) and the consumer (what to do).
Debezium and CDC: the production-grade alternative
Polling works, but at scale teams replace the polling loop with change data capture. Debezium tails the database's own transaction log — the write-ahead log in PostgreSQL, the binlog in MySQL — and streams every committed outbox row into Kafka through Kafka Connect. Nothing polls, and the relay code above disappears entirely; the application only writes the outbox row.
Why is log-tailing better than polling? Three reasons. Latency: events stream within milliseconds of commit instead of waiting for the next poll interval. Zero query load: polling runs a SELECT against your production table forever; CDC reads the log the database was already writing. Exactly the committed truth: the log only contains committed transactions, so a rolled-back outbox row can never leak — with polling you have to be careful never to read uncommitted rows (this demo reads only committed data because each statement runs in its own auto-commit transaction).
The trade is operational: Debezium means running a Kafka Connect cluster, managing connector configuration and offsets, and handling schema evolution. The pattern's contract doesn't change either way — the outbox table is the API, the relay mechanism is an implementation detail:
| Polling relay | Debezium / CDC | |
|---|---|---|
| Latency | One poll interval | Near real-time from the log |
| Load on the DB | A SELECT every interval, forever | None beyond normal logging |
| What you operate | A thread in your app | A Kafka Connect cluster + connector |
| Reads uncommitted rows? | Only if you misconfigure isolation | Impossible — the log has committed data only |
| Right for | Starting out, modest volume | High volume, strict latency, dedicated platform team |
I did not implement Debezium here — there is no Kafka broker or Connect cluster in this environment — but the outbox table this post builds is exactly the table a Debezium connector would tail in production. Write the table now; swap the relay later without touching placeOrder.
Ordering: same order, same partition
The relay reads ORDER BY id, so events leave in commit order. But ordering only survives the trip to Kafka if the producer cooperates: Kafka preserves order within a partition, not across them. So the producer keys every event by the aggregate id — the order id — and every event for one order lands in the same partition, in order. Different orders land in different partitions and interleave freely, which is exactly what you want for parallelism.
The production wiring the demo's in-memory sink stands in for (shown for completeness — not executed here; the demo has no broker):
// Production wiring for the relay's EventSink. The demo above uses an
// in-memory sink; this is what the sink looks like against a real broker.
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka:9092");
props.put(ProducerConfig.ACKS_CONFIG, "all"); // wait for all in-sync replicas
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true"); // broker dedups producer retries
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
EventSink kafkaSink = event -> {
// Key by aggregate id: one order's events -> one partition -> in order.
String key = String.valueOf(event.aggregateId());
try {
// Block for the ack BEFORE the relay marks the row published:
// slow but correct. Production code batches records and flush()es.
producer.send(new ProducerRecord<>(event.topic(), key, event.payload())).get();
} catch (InterruptedException | ExecutionException e) {
throw new RuntimeException("broker ack failed, row stays unpublished", e);
}
};
Two settings do real work there. acks=all means the broker acknowledges only after every in-sync replica has the record — a leader crash can't silently eat your event. enable.idempotence=true gives the producer a sequence number per partition so the broker discards retried sends instead of duplicating them; it protects the network hop, not the relay-crash window. The relay-crash window is the consumer's job.
Decision rule: guarantee ordering per aggregate (one order's events, in order), never globally. If your design needs a total order across all orders, you have a design problem, not a Kafka configuration problem.
The other half: consumer-side dedup
The relay is at-least-once, which means duplicates are not a bug — they're the design. So the consumer must be idempotent: processing the same event twice has the same effect as processing it once. The standard mechanism is a processed_events table keyed by the outbox's event_id, and the primary key is the arbiter. Whichever delivery wins the INSERT does the work; every later redelivery hits the duplicate key and skips:
package com.javamakeuse.orders;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.sql.SQLIntegrityConstraintViolationException;
public final class Consumer {
private Consumer() {}
/** @return true if this delivery did the work, false if it was a duplicate */
public static boolean handleOrderPlaced(Connection conn, OutboxEvent event) throws SQLException {
conn.setAutoCommit(false);
try {
try (PreparedStatement ps = conn.prepareStatement(
"INSERT INTO processed_events (event_id) VALUES (?)")) {
ps.setString(1, event.eventId());
ps.executeUpdate();
} catch (SQLIntegrityConstraintViolationException dup) {
// event_id already seen: this is a redelivery — do nothing.
conn.rollback();
return false;
}
// The real work goes here: reserve inventory, schedule shipment.
conn.commit();
return true;
} catch (SQLException e) {
conn.rollback();
throw e;
} finally {
conn.setAutoCommit(true);
}
}
}
To prove the pair, the demo simulates the relay's crash window: it flips the published row back to unpublished, as if the relay had sent the event but died before the marking UPDATE committed. The relay resends — the duplicate genuinely reaches the sink — and the consumer processes each delivery. Genuine output:
=== (d) Redelivery: relay crashes before marking, consumer dedups ===
simulated crash: row 1 reset to unpublished
relay run #3 published = 1 (at-least-once: the duplicate WAS re-sent)
sink received = 2 (sink now holds the same event twice)
consumer: fulfilled order from event c132906c...
consumer: duplicate skipped for event c132906c...
fulfilled = 1
duplicates skipped = 1
processed_events rows = 1
=> at-least-once relay + idempotent consumer = exactly-once EFFECT
Note what the dedup table must be: durable, in the consumer's own database, written in the same transaction as the consumer's side effects. An in-memory Set of seen ids loses everything on restart — exactly when redeliveries are most likely. And the check-then-act has to be one atomic step: the INSERT either wins or violates the key. A separate "check if exists, then insert" has a race between the check and the insert, and races are where duplicates breed.
Decision rule: never trust the transport to deliver exactly once — no broker or relay can promise that across crashes. Put the exactly-once effect in the consumer: a durable dedup key per event, claimed atomically, in the same transaction as the work.
Why interviews love this topic
This is the production piece that tutorial-level Kafka material skips. A candidate who says "we publish to Kafka after saving" has described both failure modes from the first diagram without knowing it. The interview-ready version of this post fits in four sentences: the dual-write problem means the database and the broker can never share a transaction, so you write the event into an outbox table in the same local transaction as the business row; a relay publishes unpublished rows and marks them, which is at-least-once because it can crash between publish and mark; the producer keys by aggregate id so one order's events stay ordered in one partition; and the consumer dedups on a durable event-id key, which is what makes the whole pipeline exactly-once in effect. Every sentence there maps to a decision in the code above — which is why "have you actually built one" is the follow-up question, and now you have.
What's next
The checkout is now reliable end to end: the order and its event commit atomically, the relay moves the event without losing it, and the consumer fulfills each order exactly once. But every one of those steps hits the database — the outbox poll, the dedup check, the order read itself. The next post in this track, "Caching in Java: Caffeine In-Process, Redis Distributed", takes the load off: what to cache, where the cache lives, and how invalidation stops stale data from undoing everything this post guaranteed.
Field check before you move on: take any service you own that writes to a database and then calls an external system — a webhook, a queue, an email sender. Add an outbox table and move the external call into a relay loop with a one-second poll, keeping the outbox insert in the existing transaction. Then simulate the crash window: kill the relay between publish and mark (a System.exit in the demo stands in for kill -9 in production), restart it, and assert the downstream side effect happened exactly once. If it happened twice, your dedup is the bug — fix that before you ship the pattern.
Continue: Java Learning Roadmap 2026
Comments
Post a Comment