Saga Orchestration with the Transactional Outbox Pattern in Spring Boot 3

Introduction

In a microservices architecture, maintaining data consistency across multiple services is one of the hardest engineering problems. A monolith can wrap an entire business operation in a single ACID database transaction. Microservices cannot — each service owns its own data store, and there is no cross-service transaction manager that scales.

Two patterns address this problem together:

  • The Saga pattern decomposes a distributed operation into a sequence of local transactions, each within one service. When a step fails, compensating transactions undo previously completed steps.
  • The Transactional Outbox pattern eliminates the "dual-write" problem: rather than updating your database AND publishing to a message broker in two separate steps that can fail independently, you write the event into an outbox_events table as part of the same local database transaction. A publisher process then relays those rows to Kafka.

Together they form a battle-tested approach to distributed consistency that does not require a distributed lock or two-phase commit.

In this tutorial we build an Order Processing System using Spring Boot 3.5. An Orchestrator in the Order service coordinates the saga across independent Payment and Inventory services. The outbox guarantees that every command reaches Kafka even when the application crashes between a database write and a Kafka send.

Prerequisites

  • Java 21+
  • Docker and Docker Compose (for local PostgreSQL and Kafka)
  • Familiarity with Spring Data JPA and Spring Kafka basics

Understanding the Dual-Write Problem

Every event-driven service eventually hits this scenario:

// Step 1 — update the database
order.setStatus(CONFIRMED);
orderRepository.save(order);

// Step 2 — tell the world
kafkaTemplate.send("order.events", event);  // what if this crashes?

If the application restarts between steps 1 and 2, the order is confirmed in the database but the downstream services never received the event. The system is now inconsistent and there is no automatic recovery.

The Transactional Outbox fixes this by merging both writes into one atomic database transaction:

@Transactional
public void confirmOrder(Order order) {
    order.setStatus(CONFIRMED);
    orderRepository.save(order);

    // Written atomically with the order update — never lost
    outboxRepository.save(new OutboxEvent("OrderConfirmed", order));
}

A separate publisher process polls the outbox table and forwards events to Kafka. Because the outbox row and the domain change are in the same transaction, either both are durable or neither is.

Project Setup

pom.xml

<parent>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-parent</artifactId>
    <version>3.5.4</version>
</parent>

<dependencies>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-web</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-data-jpa</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.kafka</groupId>
        <artifactId>spring-kafka</artifactId>
    </dependency>
    <dependency>
        <groupId>org.postgresql</groupId>
        <artifactId>postgresql</artifactId>
        <scope>runtime</scope>
    </dependency>
    <dependency>
        <groupId>org.liquibase</groupId>
        <artifactId>liquibase-core</artifactId>
    </dependency>
    <!-- Test -->
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-test</artifactId>
        <scope>test</scope>
    </dependency>
    <dependency>
        <groupId>org.springframework.kafka</groupId>
        <artifactId>spring-kafka-test</artifactId>
        <scope>test</scope>
    </dependency>
    <dependency>
        <groupId>org.testcontainers</groupId>
        <artifactId>postgresql</artifactId>
        <version>1.20.4</version>
        <scope>test</scope>
    </dependency>
    <dependency>
        <groupId>org.testcontainers</groupId>
        <artifactId>kafka</artifactId>
        <version>1.20.4</version>
        <scope>test</scope>
    </dependency>
</dependencies>

application.yml

spring:
  datasource:
    url: jdbc:postgresql://localhost:5432/orderdb
    username: user
    password: password
  jpa:
    hibernate:
      ddl-auto: validate
    show-sql: false
  kafka:
    bootstrap-servers: localhost:9092
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
      acks: all
      retries: 3
    consumer:
      group-id: order-service
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      auto-offset-reset: earliest
  liquibase:
    change-log: classpath:db/changelog/master.yaml

Domain Model

The Order saga has three forward steps and two compensating steps:

STARTED
  └─► ReserveInventoryCommand ──► INVENTORY_RESERVED
        └─► ChargePaymentCommand ──► PAYMENT_CHARGED ──► COMPLETED
              (payment fails)
              └─► COMPENSATING ──► ReleaseInventoryCommand ──► FAILED

Order Entity

package com.vimleshpandey.demo.order.domain;

import jakarta.persistence.*;
import java.math.BigDecimal;
import java.util.UUID;

@Entity
@Table(name = "orders")
public class Order {

    @Id
    @GeneratedValue(strategy = GenerationType.UUID)
    private UUID id;

    @Column(nullable = false)
    private String customerId;

    @Column(nullable = false)
    private String productId;

    private int quantity;
    private BigDecimal totalAmount;

    @Enumerated(EnumType.STRING)
    @Column(nullable = false)
    private OrderStatus status = OrderStatus.PENDING;

    @Enumerated(EnumType.STRING)
    @Column(nullable = false)
    private SagaStatus sagaStatus = SagaStatus.STARTED;

    // getters and setters omitted for brevity
}
public enum OrderStatus  { PENDING, CONFIRMED, CANCELLED }
public enum SagaStatus   { STARTED, INVENTORY_RESERVED, PAYMENT_CHARGED,
                           COMPENSATING, INVENTORY_RELEASED, COMPLETED, FAILED }

Outbox Event Entity

package com.vimleshpandey.demo.outbox.domain;

import jakarta.persistence.*;
import java.time.Instant;
import java.util.UUID;

@Entity
@Table(name = "outbox_events",
       indexes = @Index(columnList = "status, created_at"))
public class OutboxEvent {

    @Id
    @GeneratedValue(strategy = GenerationType.UUID)
    private UUID id;

    private String aggregateType;   // "Order"
    private String aggregateId;     // order UUID

    @Column(nullable = false)
    private String eventType;       // "ReserveInventoryCommand", etc.

    @Column(nullable = false)
    private String topic;

    @Column(columnDefinition = "TEXT", nullable = false)
    private String payload;

    @Enumerated(EnumType.STRING)
    private OutboxStatus status = OutboxStatus.PENDING;

    private Instant createdAt = Instant.now();
    private Instant sentAt;
    private int retryCount;

    // getters and setters
}

public enum OutboxStatus { PENDING, SENT, FAILED }

Implementing the Saga Orchestrator

The orchestrator lives inside the Order service. Each method handles one incoming event, advances the saga state, and writes a new command to the outbox — all in one transaction.

package com.vimleshpandey.demo.saga;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.vimleshpandey.demo.order.domain.*;
import com.vimleshpandey.demo.order.repository.OrderRepository;
import com.vimleshpandey.demo.outbox.domain.*;
import com.vimleshpandey.demo.outbox.repository.OutboxEventRepository;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;

import java.util.UUID;

@Service
@Transactional
@RequiredArgsConstructor
@Slf4j
public class OrderSagaOrchestrator {

    private final OrderRepository orderRepo;
    private final OutboxEventRepository outboxRepo;
    private final ObjectMapper objectMapper;

    public Order startSaga(CreateOrderRequest req) {
        Order order = new Order();
        order.setCustomerId(req.customerId());
        order.setProductId(req.productId());
        order.setQuantity(req.quantity());
        order.setTotalAmount(req.totalAmount());
        order = orderRepo.save(order);

        publishCommand(order, "ReserveInventoryCommand", "inventory.commands");
        log.info("Saga started for order {}", order.getId());
        return order;
    }

    public void onInventoryReserved(UUID orderId) {
        Order order = load(orderId);
        order.setSagaStatus(SagaStatus.INVENTORY_RESERVED);
        publishCommand(order, "ChargePaymentCommand", "payment.commands");
        log.info("Inventory reserved for order {}; charging payment", orderId);
    }

    public void onPaymentCharged(UUID orderId) {
        Order order = load(orderId);
        order.setSagaStatus(SagaStatus.COMPLETED);
        order.setStatus(OrderStatus.CONFIRMED);
        publishCommand(order, "OrderConfirmedEvent", "order.events");
        log.info("Order {} completed", orderId);
    }

    public void onPaymentFailed(UUID orderId) {
        Order order = load(orderId);
        order.setSagaStatus(SagaStatus.COMPENSATING);
        publishCommand(order, "ReleaseInventoryCommand", "inventory.commands");
        log.warn("Payment failed for order {}; compensating", orderId);
    }

    public void onInventoryReleased(UUID orderId) {
        Order order = load(orderId);
        order.setSagaStatus(SagaStatus.FAILED);
        order.setStatus(OrderStatus.CANCELLED);
        log.warn("Compensation complete for order {}; saga FAILED", orderId);
    }

    private Order load(UUID id) {
        return orderRepo.findById(id)
            .orElseThrow(() -> new IllegalStateException("Order not found: " + id));
    }

    private void publishCommand(Order order, String eventType, String topic) {
        try {
            OutboxEvent event = new OutboxEvent();
            event.setAggregateType("Order");
            event.setAggregateId(order.getId().toString());
            event.setEventType(eventType);
            event.setTopic(topic);
            event.setPayload(objectMapper.writeValueAsString(order));
            outboxRepo.save(event);
        } catch (JsonProcessingException e) {
            throw new IllegalStateException("Cannot serialize order payload", e);
        }
    }
}

Implementing the Outbox Publisher

The publisher runs on a fixed schedule. The FOR UPDATE SKIP LOCKED clause in the query means multiple publisher instances can run concurrently — each grabs a different set of rows and there is no double-publishing.

package com.vimleshpandey.demo.outbox;

import com.vimleshpandey.demo.outbox.domain.*;
import com.vimleshpandey.demo.outbox.repository.OutboxEventRepository;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;

import java.time.Instant;
import java.util.List;
import java.util.concurrent.TimeUnit;

@Component
@RequiredArgsConstructor
@Slf4j
public class OutboxPublisher {

    private static final int BATCH_SIZE = 20;
    private static final int MAX_RETRIES = 5;

    private final OutboxEventRepository outboxRepo;
    private final KafkaTemplate<String, String> kafkaTemplate;

    @Scheduled(fixedDelay = 1000)
    @Transactional
    public void publish() {
        List<OutboxEvent> pending = outboxRepo.findPendingWithLock(BATCH_SIZE);
        for (OutboxEvent event : pending) {
            try {
                kafkaTemplate.send(event.getTopic(), event.getAggregateId(), event.getPayload())
                    .get(5, TimeUnit.SECONDS);   // wait for broker ack before marking SENT
                event.setStatus(OutboxStatus.SENT);
                event.setSentAt(Instant.now());
                log.debug("Published {} for aggregate {}", event.getEventType(), event.getAggregateId());
            } catch (Exception ex) {
                int retries = event.getRetryCount() + 1;
                event.setRetryCount(retries);
                if (retries >= MAX_RETRIES) {
                    event.setStatus(OutboxStatus.FAILED);
                    log.error("Outbox event {} exceeded retry limit — marking FAILED", event.getId());
                } else {
                    log.warn("Outbox publish attempt {} failed for event {}: {}",
                             retries, event.getId(), ex.getMessage());
                }
            }
            outboxRepo.save(event);
        }
    }
}

Outbox Repository

package com.vimleshpandey.demo.outbox.repository;

import com.vimleshpandey.demo.outbox.domain.OutboxEvent;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.data.jpa.repository.Query;
import org.springframework.data.repository.query.Param;

import java.util.List;
import java.util.UUID;

public interface OutboxEventRepository extends JpaRepository<OutboxEvent, UUID> {

    @Query(value = "SELECT * FROM outbox_events " +
                   "WHERE status = 'PENDING' " +
                   "ORDER BY created_at " +
                   "LIMIT :limit " +
                   "FOR UPDATE SKIP LOCKED",
           nativeQuery = true)
    List<OutboxEvent> findPendingWithLock(@Param("limit") int limit);

    List<OutboxEvent> findByAggregateId(String aggregateId);
}

Saga Event Listener

The Order service consumes response events from Inventory and Payment services and drives the saga forward:

package com.vimleshpandey.demo.saga;

import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;

import java.util.UUID;

@Component
@RequiredArgsConstructor
@Slf4j
public class SagaEventListener {

    private final OrderSagaOrchestrator orchestrator;
    private final ObjectMapper objectMapper;

    @KafkaListener(topics = "inventory.responses", groupId = "order-service")
    public void onInventoryResponse(String message) {
        SagaResponse response = parse(message);
        UUID orderId = response.orderId();
        if ("InventoryReserved".equals(response.eventType())) {
            orchestrator.onInventoryReserved(orderId);
        } else if ("InventoryReleased".equals(response.eventType())) {
            orchestrator.onInventoryReleased(orderId);
        }
    }

    @KafkaListener(topics = "payment.responses", groupId = "order-service")
    public void onPaymentResponse(String message) {
        SagaResponse response = parse(message);
        if (response.success()) {
            orchestrator.onPaymentCharged(response.orderId());
        } else {
            orchestrator.onPaymentFailed(response.orderId());
        }
    }

    private SagaResponse parse(String message) {
        try {
            return objectMapper.readValue(message, SagaResponse.class);
        } catch (Exception e) {
            throw new IllegalStateException("Cannot parse saga response: " + message, e);
        }
    }
}
// A simple record to carry the response payload
public record SagaResponse(UUID orderId, String eventType, boolean success) {}

Database Schema (Liquibase)

# src/main/resources/db/changelog/001-create-tables.yaml
databaseChangeLog:
  - changeSet:
      id: 001
      author: vimlesh
      changes:
        - createTable:
            tableName: orders
            columns:
              - column:
                  name: id
                  type: UUID
                  constraints: { primaryKey: true, nullable: false }
              - column: { name: customer_id,   type: VARCHAR(255), constraints: { nullable: false } }
              - column: { name: product_id,    type: VARCHAR(255), constraints: { nullable: false } }
              - column: { name: quantity,      type: INT,          constraints: { nullable: false } }
              - column: { name: total_amount,  type: DECIMAL(19,2) }
              - column: { name: status,        type: VARCHAR(50),  defaultValue: PENDING }
              - column: { name: saga_status,   type: VARCHAR(50),  defaultValue: STARTED }

        - createTable:
            tableName: outbox_events
            columns:
              - column:
                  name: id
                  type: UUID
                  constraints: { primaryKey: true, nullable: false }
              - column: { name: aggregate_type, type: VARCHAR(255) }
              - column: { name: aggregate_id,   type: VARCHAR(255) }
              - column: { name: event_type,     type: VARCHAR(255), constraints: { nullable: false } }
              - column: { name: topic,          type: VARCHAR(255), constraints: { nullable: false } }
              - column: { name: payload,        type: TEXT }
              - column: { name: status,         type: VARCHAR(50),  defaultValue: PENDING }
              - column: { name: created_at,     type: TIMESTAMP }
              - column: { name: sent_at,        type: TIMESTAMP }
              - column: { name: retry_count,    type: INT,          defaultValue: 0 }

        - createIndex:
            indexName: idx_outbox_status_created
            tableName: outbox_events
            columns:
              - column: { name: status }
              - column: { name: created_at }

Testing with Testcontainers

Testcontainers spins up real PostgreSQL and Kafka containers for integration tests, ensuring the FOR UPDATE SKIP LOCKED query runs against an actual database engine and the outbox publisher exercises real Kafka brokers.

package com.vimleshpandey.demo;

import com.vimleshpandey.demo.order.domain.*;
import com.vimleshpandey.demo.outbox.domain.*;
import com.vimleshpandey.demo.outbox.repository.OutboxEventRepository;
import com.vimleshpandey.demo.order.repository.OrderRepository;
import com.vimleshpandey.demo.saga.OrderSagaOrchestrator;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import org.testcontainers.containers.KafkaContainer;
import org.testcontainers.containers.PostgreSQLContainer;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.utility.DockerImageName;

import java.math.BigDecimal;
import java.util.List;
import java.util.UUID;

import static org.assertj.core.api.Assertions.assertThat;

@SpringBootTest
@Testcontainers
class OrderSagaIntegrationTest {

    @Container
    static PostgreSQLContainer<?> postgres =
        new PostgreSQLContainer<>("postgres:16-alpine");

    @Container
    static KafkaContainer kafka =
        new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:7.7.1"));

    @DynamicPropertySource
    static void overrideProps(DynamicPropertyRegistry registry) {
        registry.add("spring.datasource.url",      postgres::getJdbcUrl);
        registry.add("spring.datasource.username", postgres::getUsername);
        registry.add("spring.datasource.password", postgres::getPassword);
        registry.add("spring.kafka.bootstrap-servers", kafka::getBootstrapServers);
    }

    @Autowired OrderSagaOrchestrator orchestrator;
    @Autowired OutboxEventRepository outboxRepo;
    @Autowired OrderRepository orderRepo;

    @Test
    void startSaga_writesReserveInventoryCommandToOutbox() {
        var req = new CreateOrderRequest("cust-1", "prod-42", 2, new BigDecimal("49.99"));
        Order order = orchestrator.startSaga(req);

        List<OutboxEvent> events = outboxRepo.findByAggregateId(order.getId().toString());
        assertThat(events).hasSize(1);
        assertThat(events.get(0).getEventType()).isEqualTo("ReserveInventoryCommand");
        assertThat(events.get(0).getStatus()).isEqualTo(OutboxStatus.PENDING);
    }

    @Test
    void onPaymentFailed_writesReleaseInventoryCommandToOutbox() {
        var req = new CreateOrderRequest("cust-2", "prod-7", 1, new BigDecimal("19.99"));
        Order order = orchestrator.startSaga(req);
        orchestrator.onInventoryReserved(order.getId());
        orchestrator.onPaymentFailed(order.getId());

        Order updated = orderRepo.findById(order.getId()).orElseThrow();
        assertThat(updated.getSagaStatus()).isEqualTo(SagaStatus.COMPENSATING);

        boolean hasCompensation = outboxRepo
            .findByAggregateId(order.getId().toString()).stream()
            .anyMatch(e -> "ReleaseInventoryCommand".equals(e.getEventType()));
        assertThat(hasCompensation).isTrue();
    }

    @Test
    void fullHappyPath_orderConfirmed() {
        var req = new CreateOrderRequest("cust-3", "prod-99", 5, new BigDecimal("199.99"));
        Order order = orchestrator.startSaga(req);
        orchestrator.onInventoryReserved(order.getId());
        orchestrator.onPaymentCharged(order.getId());

        Order completed = orderRepo.findById(order.getId()).orElseThrow();
        assertThat(completed.getStatus()).isEqualTo(OrderStatus.CONFIRMED);
        assertThat(completed.getSagaStatus()).isEqualTo(SagaStatus.COMPLETED);
    }
}

Running the Application

Start PostgreSQL and Kafka with Docker Compose:

# docker-compose.yml
version: "3.9"
services:
  postgres:
    image: postgres:16-alpine
    environment:
      POSTGRES_DB: orderdb
      POSTGRES_USER: user
      POSTGRES_PASSWORD: password
    ports: ["5432:5432"]

  zookeeper:
    image: confluentinc/cp-zookeeper:7.7.1
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181

  kafka:
    image: confluentinc/cp-kafka:7.7.1
    depends_on: [zookeeper]
    environment:
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
    ports: ["9092:9092"]
docker-compose up -d
./mvnw spring-boot:run

Place an order:

curl -X POST http://localhost:8080/api/orders \
  -H "Content-Type: application/json" \
  -d '{"customerId":"c1","productId":"p42","quantity":2,"totalAmount":99.99}'

Watch the outbox table fill with commands and drain as the publisher forwards them:

SELECT event_type, status, created_at, sent_at
FROM outbox_events
ORDER BY created_at;

Common Pitfalls and Tips

1. Idempotent consumers are not optional. Kafka delivers at-least-once, so the same ReserveInventoryCommand may arrive twice. Your Inventory service must detect duplicate commands using the outbox event ID or a combination of aggregateId + eventType. Store processed IDs in an inbox_events table.

2. Don't swallow the Kafka ack. The .get(5, TimeUnit.SECONDS) call blocks until the broker acknowledges the record. If you fire-and-forget with just .send(...), a broker crash between send and ack means the event is lost while the row is marked SENT. Always wait for the ack before committing the status change.

3. Polling overhead at high throughput. A 1-second poll works fine up to thousands of events per minute. Beyond that, replace polling with Debezium CDC, which tails the PostgreSQL write-ahead log and forwards rows to Kafka with sub-100ms latency and zero application-level polling.

4. Clock skew breaks ORDER BY created_at. In a multi-instance deployment, two nodes can write outbox rows with identical or inverted timestamps. Switch your ordering key to a monotonic database sequence (serial or bigserial column) to guarantee causal ordering.

5. Saga state machine sprawl. As business requirements grow, the orchestrator accumulates special cases. Consider encoding transitions in Spring Statemachine or a simple transition table to prevent invalid state jumps (e.g., STARTED → COMPLETED bypassing INVENTORY_RESERVED).

6. Compensating transactions must also use the outbox. A compensation that publishes directly to Kafka without going through the outbox has the same dual-write problem as the original operation. Every command — forward and compensating — goes through the outbox table.

7. Monitor FAILED outbox rows. Rows that exceed the retry limit are marked FAILED rather than dropped. Alert on this metric and route them to a dead-letter topic for manual inspection. Silently discarding events makes saga failures invisible.

Conclusion

The Saga + Transactional Outbox combination addresses two distinct problems cleanly:

  • The Saga orchestrator models the business transaction as an explicit state machine, making the happy path and every failure mode visible in code and traceable in the database.
  • The Transactional Outbox decouples the reliability of event delivery from the application's runtime health, ensuring every saga step eventually reaches Kafka regardless of crashes between the database write and the publish.

The key architectural principle: treat the outbox table as the source of truth for outbound events, never Kafka directly. Kafka is a delivery mechanism; the database is the contract.

For production, the next step is replacing the polling publisher with Debezium CDC — but the orchestrator, outbox schema, and consumer design remain identical. That migration is purely operational, not architectural.