Stop dual-writing to Postgres and Kafka. A small, boring pattern that keeps them in sync, with the real code.
If you've ever written something like this, you've got a bug. It just hasn't happened yet:
orderRepository.save(order);
kafkaTemplate.send("order.events", orderJson); // 💥 what if this fails? or the DB commit fails after it?
That's a dual write. Two systems, no shared transaction. If the DB commit fails after the send, downstream services react to an order that doesn't exist. If the send fails after the commit, the order exists and nobody ever hears about it. Retrying doesn't fix it, it just moves the window around.
The fix most teams end up with is the transactional outbox, plus idempotent consumers on the other side. Here's a version that runs on Spring Boot 4.1 / Java 21 with Spring Kafka, JPA and Flyway. The snippets are taken straight from a working project (an Order → Inventory saga), so they compile and they're covered by tests.
The idea in one paragraph
Don't publish to Kafka inside your business transaction. Insert a row into an outbox_events table in the same DB transaction as your business change. A separate publisher reads unpublished rows and sends them to Kafka. The DB is now the single source of truth: either both the order and its event exist, or neither does. You get at-least-once delivery, so consumers have to tolerate duplicates. That's the second half.
Step 1: one envelope for every event
Each message carries an eventId. That ID becomes the idempotency key later.
/**
* Envelope for every Kafka message. {@code eventId} is the idempotency key.
*/
public record DomainEvent(
UUID eventId,
String type,
Instant occurredAt,
String aggregateType,
UUID aggregateId,
int version,
String payloadJson
) {
public static DomainEvent of(String type, String aggregateType, UUID aggregateId, int version, String payloadJson) {
return new DomainEvent(UUID.randomUUID(), type, Instant.now(), aggregateType, aggregateId, version, payloadJson);
}
}
Step 2: business row + outbox row, one transaction
The outbox table (Flyway migration):
CREATE TABLE outbox_events (
id UUID PRIMARY KEY,
topic VARCHAR(120) NOT NULL,
partition_key VARCHAR(120) NOT NULL,
payload TEXT NOT NULL,
created_at TIMESTAMP WITH TIME ZONE NOT NULL,
published_at TIMESTAMP WITH TIME ZONE
);
CREATE INDEX idx_outbox_created ON outbox_events (created_at);
The writer just serializes the envelope and saves a row. Note that it has no @Transactional of its own. It joins whatever transaction the caller already opened:
public void enqueue(String topic, String partitionKey, DomainEvent event) {
try {
String json = objectMapper.writeValueAsString(event);
repository.save(new OutboxEvent(event.eventId(), topic, partitionKey, json, Instant.now(clock)));
} catch (JacksonException e) {
throw new IllegalStateException("Cannot serialize domain event", e);
}
}
(Small Boot 4 detail: that's Jackson 3, so the imports are tools.jackson.* and the exception is JacksonException.)
And the service that places an order:
@Transactional
public Order placeOrder(String productSku, int quantity) {
if (quantity < 1) {
throw new IllegalArgumentException("quantity must be >= 1");
}
Instant now = Instant.now(clock);
Order order = new Order(UUID.randomUUID(), productSku, quantity, now);
orderRepository.save(order);
String payload = toJson(new OrderPayloads.OrderCreated(order.getId(), productSku, quantity));
DomainEvent event = DomainEvent.of(
EventTypes.ORDER_CREATED, "Order", order.getId(), 1, payload);
outboxWriter.enqueue(props.topics().orders(), order.getId().toString(), event);
return order;
}
Two things worth noticing:
-
No Kafka call here at all. The HTTP request returns
201with the order inPENDING, and the event goes out asynchronously. - The partition key is the order ID. All events for the same order land on the same partition, so they're consumed in order.
Step 3: the publisher
A scheduled job polls unpublished rows, sends them, and marks them published:
@Scheduled(fixedDelayString = "${app.outbox.poll-interval-ms:500}")
@Transactional
public void publishPending() {
List<OutboxEvent> batch = repository.findUnpublished().stream()
.limit(props.outbox().batchSize())
.toList();
for (OutboxEvent event : batch) {
try {
kafkaTemplate.send(event.getTopic(), event.getPartitionKey(), event.getPayload()).get();
event.markPublished(Instant.now(clock));
log.info("Outbox published eventId={} topic={}", event.getId(), event.getTopic());
} catch (Exception ex) {
log.warn("Outbox publish failed eventId={}: {}", event.getId(), ex.getMessage());
break;
}
}
}
@Query("""
select e from OutboxEvent e
where e.publishedAt is null
order by e.createdAt asc
""")
List<OutboxEvent> findUnpublished();
Why it's written this way:
-
.get()makes the send synchronous. A row is only marked published after the broker acknowledged it. -
breakon the first failure keeps ordering. We don't skip a failed event and publish the next one for the same order. -
Where duplicates come from: if the send succeeds but the DB commit that sets
published_atfails (crash, connection drop), the row is still "unpublished" and gets sent again next tick. That's the at-least-once part, and it's expected.
The producer side is set up to avoid making things worse:
spring:
kafka:
producer:
acks: all
properties:
enable.idempotence: true
The idempotent producer prevents duplicates caused by the producer's own retries. It can't prevent the "sent, then crashed before commit" duplicate. Only the consumer can deal with that one.
Step 4: idempotent consumers
Every consumer records which events it already handled, keyed by (event_id, consumer_name):
CREATE TABLE processed_events (
event_id UUID NOT NULL,
consumer_name VARCHAR(80) NOT NULL,
processed_at TIMESTAMP WITH TIME ZONE NOT NULL,
PRIMARY KEY (event_id, consumer_name)
);
The consumer name is part of the key on purpose. The same event can legitimately be processed by inventory-service and by some other consumer later.
The core of it is small:
@Transactional
public void processOnce(String consumerName, DomainEvent event, Consumer<DomainEvent> handler) {
if (processedEventRepository.existsByEventIdAndConsumerName(event.eventId(), consumerName)) {
log.info("Skipping duplicate eventId={} consumer={}", event.eventId(), consumerName);
return;
}
handler.accept(event);
processedEventRepository.save(new ProcessedEvent(event.eventId(), consumerName, Instant.now(clock)));
}
And a listener uses it like this:
@KafkaListener(topics = "${app.topics.orders}", groupId = "inventory-service")
public void onOrderEvent(String payload) {
DomainEvent event = processor.parse(payload);
processor.processOnce(CONSUMER, event, this::handle);
}
The important part is what's inside that one transaction: the duplicate check, the handler's DB changes, and the processed_events insert. In the sample, the inventory handler reserves stock and then writes its own inventory.reserved (or inventory.failed) event to the outbox, all inside the same transaction. So "consume → update state → emit next event" is atomic from the database's point of view:
- Crash before commit? Nothing happened, and Kafka redelivers.
- Crash after commit but before the offset is committed? Kafka redelivers,
existsBy...returns true, and it's skipped. - Two copies processed concurrently? The composite primary key makes the second insert fail, its transaction rolls back (including the handler's changes), and on retry it's skipped as a duplicate.
Offsets are only committed after the listener returns (enable-auto-commit: false, so the Spring Kafka container handles commits), which is exactly the ordering you want here.
Step 5: retries and a dead letter topic
Not every failure is transient. The listener container gets a DefaultErrorHandler with exponential backoff and a dead-letter recoverer:
DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(kafkaTemplate,
(record, ex) -> {
// routes to order.events.DLT or inventory.events.DLT based on record.topic()
// ...
});
ExponentialBackOffWithMaxRetries backOff = new ExponentialBackOffWithMaxRetries(3);
backOff.setInitialInterval(200L);
backOff.setMultiplier(2.0);
backOff.setMaxInterval(2_000L);
DefaultErrorHandler errorHandler = new DefaultErrorHandler(recoverer, backOff);
errorHandler.addNotRetryableExceptions(IllegalArgumentException.class);
factory.setCommonErrorHandler(errorHandler);
So that's 3 retries at roughly 200 ms, 400 ms and 800 ms, then the record goes to the DLT, and the partition keeps moving instead of blocking on one poison message. IllegalArgumentException is marked non-retryable, and that's what the parser throws for malformed JSON. Retrying a broken payload three times just wastes time.
Two practical notes if you adapt this:
-
Choose your exception types on purpose. Whatever you mark non-retryable skips the backoff entirely. In the sample, the handlers wrap their own failures in
IllegalArgumentException. If your handler can hit transient problems (a DB timeout, a downstream call), let those escape as a different exception type so they get retried. - Watch DLT partition counts. If your resolver targets the same partition number as the source record, the DLT needs at least as many partitions as the source topic. Otherwise pick a partition that exists.
Testing it without Docker
The unit test for the dedupe logic is almost boring, which is the point:
@Test
void skipsDuplicateEvent() {
DomainEvent event = DomainEvent.of("order.created", "Order", UUID.randomUUID(), 1, "{}");
when(repository.existsByEventIdAndConsumerName(event.eventId(), "c1")).thenReturn(true);
AtomicInteger calls = new AtomicInteger();
processor.processOnce("c1", event, e -> calls.incrementAndGet());
assertThat(calls.get()).isZero();
verify(repository, never()).save(any());
}
For the full flow, an @EmbeddedKafka + H2 integration test posts an order over HTTP and uses Awaitility to wait until it becomes CONFIRMED (or REJECTED when stock is insufficient). It exercises the outbox, the publisher, both listeners and the saga, and ./mvnw verify runs it without Docker.
Trade-offs (read before copying into prod)
This is the simple, polling flavor of the outbox. It's a good default, but be honest with yourself about scale:
- Polling vs CDC. Polling every 500 ms is easy to reason about. At high volume, log-based CDC (e.g. Debezium reading the outbox table) cuts latency and DB load.
- Batch size. Here the batch limit is applied after the query runs. With a large backlog, push the limit into the query itself (paging).
-
Multiple instances. Two app instances can pick up the same unpublished rows. Consumers dedupe, so it's correct, but it's wasteful.
SELECT ... FOR UPDATE SKIP LOCKEDor a single active publisher fixes that. - The transaction stays open during the sends. That's fine for small batches. Keep batches small.
-
Housekeeping.
outbox_eventsandprocessed_eventsgrow forever unless you prune them. A periodic delete of old published rows is enough.
None of this is exotic. It's just the stuff that's easy to skip on day one and painful to retrofit on day ninety.
I packaged this whole setup (outbox, idempotent consumers, retries + DLT, the Order → Inventory saga sample, Docker Compose with Postgres 17 + Kafka KRaft, CI, and 16 tests) as a paid source-code starter, if you'd rather not wire it from scratch: Spring Boot 4 Kafka Event Starter (25 USD). Built by me, a Java backend engineer in Mexico City (GitHub: kodkodmx). The article stands on its own, and you don't need to buy anything to use the pattern.
Top comments (0)