Integrating Spring Boot with Apache Kafka for Exactly-Once Message Processing

Integrating Spring Boot with Apache Kafka for Exactly-Once Message Processing

Intended Readers and Outcomes

This comprehensive guide targets Java developers and architects who build event-driven microservices using Spring Boot and Apache Kafka. If you require a fault-tolerant system that guarantees messages are processed exactly once—even in the event of retries, restarts, or failures—this guide will clarify how to achieve exactly-once semantics (EOS) using Kafka transactions integrated with Spring Boot.

After completing this tutorial, you will know how to configure Kafka producers and consumers with transactional guarantees, implement transactional message handling within Spring Boot services, and validate that your solution reliably processes messages without duplicates or data loss.

Prerequisites and Version Assumptions

  • Java 11 or later
  • Spring Boot 3.x
  • Apache Kafka 2.5 or newer broker and client
  • Maven or Gradle for dependency management
  • Basic understanding of Kafka's producer-consumer model, topics, partitions, and offsets
  • Familiarity with Spring annotations and transaction management

Why Exactly-Once Processing Matters

When dealing with distributed systems and asynchronous messaging, achieving correct state consistency is notoriously difficult. Messages may be delivered multiple times due to retries, or lost due to failures. Exactly-once semantics ensure each message changes your system state exactly one time, which eliminates duplication and inconsistency.

Use Cases Ideal for Exactly-Once Processing

  • Payment or financial transaction systems where duplicate processing leads to reconciliation errors
  • Inventory or stock updates where consistency affects availability
  • Order management pipelines to prevent double shipments
  • Event-driven data synchronization between services or databases

When to Avoid Exactly-Once Semantics

  • Systems with relaxed consistency requirements where occasional duplication is acceptable
  • High-throughput analytics or logging systems where at-least-once semantics suffice
  • Cases where the additional latency and complexity of transactions are prohibitive

Alternatives and Trade-offs

Semantic LevelDescriptionGuaranteesTrade-offs
At-Most-OnceDeliver zero or one times; may lose messagesNo duplicates, possible lossSimplest, fastest, but unsafe
At-Least-OnceDeliver one or more times; possible duplicatesNo loss, possible duplicatesRequires idempotent consumers
Exactly-OnceDeliver one and only one timeNo loss, no duplicatesIncreased complexity and resource use

Kafka achieves EOS primarily through producer idempotence and transactions, which tie together sending data and committing offsets atomically.

Overview of Kafka Exactly-Once Semantics (EOS)

Key Kafka Features Enabling EOS

  1. Idempotent Producers: Automatically manage sequence numbers and producer IDs to avoid duplicate messages on retries.
  2. Transactions: Group multiple produce and offset commit operations into a single atomic unit.
  3. Consumer Isolation: Consumers configured with isolation.level=read_committed will only read data from completed transactions.

This combination lets producers write messages and commit offsets atomically. Consumers avoid seeing uncommitted or aborted messages, preventing duplicates on recovery.

SemanticsBehaviorGuarantees
At-Most-OnceMessages delivered zero or one timePossible loss, no duplicates
At-Least-OnceMessages delivered one or more timesNo loss, possible duplicates
Exactly-Once (EOS)Messages delivered exactly onceNo loss, no duplicates

Setting Up Your Development Environment

Kafka Broker Setup

  1. Download and unzip Apache Kafka 2.5 or later.
  2. Start ZooKeeper (required by Kafka 2.5):
bin/zookeeper-server-start.sh config/zookeeper.properties
  1. Start Kafka broker:
bin/kafka-server-start.sh config/server.properties
  1. Create a transactional topic (highly recommended 3 partitions for parallelism):
bin/kafka-topics.sh --create --topic transactional-topic --bootstrap-server localhost:9092 --partitions 3 --replication-factor 1

Ensure the topic supports transactions by default with your Kafka version.

Spring Boot Project Dependencies

Add the following to your Maven pom.xml:

<dependencies>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.kafka</groupId>
        <artifactId>spring-kafka</artifactId>
    </dependency>
</dependencies>

Or in Gradle build.gradle:

dependencies {
    implementation 'org.springframework.boot:spring-boot-starter'
    implementation 'org.springframework.kafka:spring-kafka'
}

Kafka Transactional Configuration in Spring Boot

Core Configuration Options

PropertyPurpose
enable.idempotence=trueEnable idempotent producer requests
transactional.idUniquely identifies the transactional producer
isolation.level=read_committedConsumers read only committed transactions

You must assign unique transactional IDs per producer instance for fault tolerance and recovery.

Complete Kafka Configuration Class

@Configuration
public class KafkaConfig {

    @Bean
    public ProducerFactory<String, String> producerFactory() {
        Map<String, Object> props = new HashMap<>();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); // Idempotence
        props.put(ProducerConfig.ACKS_CONFIG, "all"); // Strongest acknowledgement
        props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "txn-id-"); // Must be unique per instance
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        return new DefaultKafkaProducerFactory<>(props);
    }

    @Bean
    public KafkaTransactionManager<String, String> kafkaTransactionManager(ProducerFactory<String, String> producerFactory) {
        return new KafkaTransactionManager<>(producerFactory);
    }

    @Bean
    public KafkaTemplate<String, String> kafkaTemplate(ProducerFactory<String, String> producerFactory) {
        KafkaTemplate<String, String> kafkaTemplate = new KafkaTemplate<>(producerFactory);
        kafkaTemplate.setTransactionIdPrefix("txn-id-"); // Matches transactional.id prefix
        return kafkaTemplate;
    }

    @Bean
    public ConsumerFactory<String, String> consumerFactory() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "transactional-consumer-group");
        props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed"); // Only read committed transactions
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        // Disable auto commit to allow transactional offset management
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
        return new DefaultKafkaConsumerFactory<>(props);
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(
            ConsumerFactory<String, String> consumerFactory,
            KafkaTransactionManager<String, String> transactionManager) {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory);
        factory.setTransactionManager(transactionManager);
        // Acknowledge each record to commit offsets individually
        factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.RECORD);

        // Handle retries for transient errors
        factory.setErrorHandler(new SeekToCurrentErrorHandler());

        return factory;
    }
}

Explanation

  • The producer factory is configured with idempotence and a transaction ID prefix.
  • KafkaTransactionManager integrates Kafka transactions with Spring's transaction abstraction.
  • Kafka template uses transactions transparently.
  • Consumer factory reads only committed transactions and disables auto offset commit.
  • Listener container factory handles transactions and errors consistently.

Implementing Transactional Producer and Consumer Services

Transactional Producer Service

@Service
public class TransactionalKafkaProducer {

    private final KafkaTemplate<String, String> kafkaTemplate;

    public TransactionalKafkaProducer(KafkaTemplate<String, String> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    /**
     * Sends a message within a Kafka transaction.
     * @param topic topic to send to
     * @param key message key
     * @param message payload
     */
    public void sendMessageTransactional(String topic, String key, String message) {
        kafkaTemplate.executeInTransaction(kt -> {
            kt.send(topic, key, message);
            // Additional sends or operations can be included here
            return true;
        });
    }
}

How it works: The executeInTransaction ensures messages are sent atomically. If an exception occurs, the transaction aborts and no partial writes or offset commits occur.

Transactional Consumer Service

@Service
public class TransactionalKafkaConsumer {

    @KafkaListener(topics = "transactional-topic", containerFactory = "kafkaListenerContainerFactory")
    @Transactional("kafkaTransactionManager")
    public void listen(@Payload String message,
                       @Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) String key) {
        System.out.printf("Processing message with key=%s: %s\n", key, message);
        // Business logic goes here
        // If exceptions happen here, transaction is rolled back and message reprocessed
    }
}

Why it works: The listener runs within a Kafka-managed transaction, which commits offsets only on successful processing. When exceptions occur, the container rolls back, and messages are retried, guaranteeing exactly-once.

Validation and Verification Steps

Manual Testing

  1. Start Kafka and your Spring Boot application.
  2. Use your TransactionalKafkaProducer to send messages:
producer.sendMessageTransactional("transactional-topic", "key1", "Hello Kafka EOS");
  1. Verify the consumer logs each message exactly once.
  2. Restart the consumer application mid-processing; ensure no duplicate processing.

Using Embedded Kafka for Integration Test

@SpringBootTest
@EmbeddedKafka(partitions = 1, topics = {"transactional-topic"})
public class ExactlyOnceIntegrationTest {

    @Autowired
    private TransactionalKafkaProducer producer;

    @Autowired
    private EmbeddedKafkaBroker embeddedKafka;

    private KafkaConsumer<String, String> testConsumer;

    @BeforeEach
    public void setup() {
        Map<String, Object> configs = new HashMap<>(KafkaTestUtils.consumerProps("testGroup", "false", embeddedKafka));
        configs.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
        testConsumer = new KafkaConsumer<>(configs, new StringDeserializer(), new StringDeserializer());
        testConsumer.subscribe(Collections.singleton("transactional-topic"));
    }

    @Test
    public void testExactlyOnceProcessing() {
        String topic = "transactional-topic";
        String key = "key1";
        String message = "test message";

        producer.sendMessageTransactional(topic, key, message);

        ConsumerRecords<String, String> records = KafkaTestUtils.getRecords(testConsumer);
        assertThat(records.count()).isEqualTo(1);

        ConsumerRecord<String, String> record = records.iterator().next();
        assertThat(record.key()).isEqualTo(key);
        assertThat(record.value()).isEqualTo(message);
    }

    @AfterEach
    public void teardown() {
        testConsumer.close();
    }
}

Expected Results:

  • Only one message is consumed.
  • The key and value match exactly.
  • Offsets are committed accordingly.

Common Failure Modes and Troubleshooting

Failure ModeCauseMitigation
Transaction timeoutsTransactions blocking beyond configured timeoutKeep transactions short. Adjust transaction.timeout.ms in broker.
Missing transactional.idProducer is not configured for transactionsSet unique transactional.id per producer instance
Consumer reads uncommitted dataisolation.level set incorrectlyUse read_committed for consumers
Producer retries not idempotentenable.idempotence is false or missingEnable enable.idempotence=true on producers

Troubleshooting Tips

  • Enable debug logging for org.apache.kafka to trace Kafka internals.
  • Inspect Kafka broker logs for aborted or timed out transactions.
  • Use Spring-Kafka’s error handling and retry mechanisms.
  • Monitor Kafka metrics for transactional coordinator load.

Security and Operational Considerations

  • Use SSL encryption and SASL authentication for message security and integrity.
  • Ensure transactional producers use unique and meaningful transaction IDs to avoid clashes.
  • Monitor transaction coordinator performance to prevent bottlenecks.
  • Limit transaction size and duration to prevent resource starvation.
  • Implement dead-letter queues to handle poison messages that consistently fail.

Performance Implications

Transactions come with overhead due to extra coordination:

  • Increased latency due to atomic commit protocols.
  • Extra broker resource usage and network roundtrips.

Tune producer batch size, linger time, and broker transaction timeouts based on your throughput and latency needs.

Limitations

  • EOS applies within Kafka brokers but does not guarantee distributed transactional consistency across multiple microservices or external databases.
  • Long running transactional operations risk timing out.
  • Each producer instance must have a unique transactional ID to avoid conflicts.

Summary

Implementing exactly-once processing with Kafka and Spring Boot involves enabling key Kafka properties (idempotence and transactions), configuring consumers for committed isolation, and integrating Spring’s transaction management mechanisms. This architecture reliably avoids duplicate processing and data loss even under failure conditions, simplifying the development of reliable event-driven microservices.

Careful configuration, combined with monitoring and testing, is required to implement EOS effectively while balancing resource usage and latency.

FAQ

Can exactly-once semantics be achieved across multiple Kafka topics?

Yes. Kafka transactions can atomically write to multiple topics and partitions, committing both data and offsets in a single transaction, enabling exactly-once semantics across multiple topics.

Does enabling idempotence alone guarantee exactly-once processing?

No. Idempotence prevents duplicate records on producer retries, but exactly-once processing also requires transactional offset commits and consumers set to read only committed data.

How does read_committed isolation affect Kafka consumers?

Consumers set to read_committed will only receive messages from transactions that have been committed, skipping those in aborted or ongoing transactions, preventing consumers from reading inconsistent data.

Sources and further reading

Related reading