DEV Community

Alberto Haboba
Alberto Haboba

Posted on

Transactional outbox + idempotent Kafka consumers in Spring Boot 4

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?
Enter fullscreen mode Exit fullscreen mode

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);
    }
}
Enter fullscreen mode Exit fullscreen mode

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);
Enter fullscreen mode Exit fullscreen mode

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);
    }
}
Enter fullscreen mode Exit fullscreen mode

(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;
}
Enter fullscreen mode Exit fullscreen mode

Two things worth noticing:

  • No Kafka call here at all. The HTTP request returns 201 with the order in PENDING, 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;
        }
    }
}
Enter fullscreen mode Exit fullscreen mode
@Query("""
        select e from OutboxEvent e
        where e.publishedAt is null
        order by e.createdAt asc
        """)
List<OutboxEvent> findUnpublished();
Enter fullscreen mode Exit fullscreen mode

Why it's written this way:

  • .get() makes the send synchronous. A row is only marked published after the broker acknowledged it.
  • break on 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_at fails (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
Enter fullscreen mode Exit fullscreen mode

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)
);
Enter fullscreen mode Exit fullscreen mode

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)));
}
Enter fullscreen mode Exit fullscreen mode

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);
}
Enter fullscreen mode Exit fullscreen mode

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);
Enter fullscreen mode Exit fullscreen mode

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:

  1. 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.
  2. 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());
}
Enter fullscreen mode Exit fullscreen mode

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 LOCKED or a single active publisher fixes that.
  • The transaction stays open during the sends. That's fine for small batches. Keep batches small.
  • Housekeeping. outbox_events and processed_events grow 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)