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 Level | Description | Guarantees | Trade-offs |
|---|---|---|---|
| At-Most-Once | Deliver zero or one times; may lose messages | No duplicates, possible loss | Simplest, fastest, but unsafe |
| At-Least-Once | Deliver one or more times; possible duplicates | No loss, possible duplicates | Requires idempotent consumers |
| Exactly-Once | Deliver one and only one time | No loss, no duplicates | Increased 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
- Idempotent Producers: Automatically manage sequence numbers and producer IDs to avoid duplicate messages on retries.
- Transactions: Group multiple produce and offset commit operations into a single atomic unit.
- Consumer Isolation: Consumers configured with
isolation.level=read_committedwill 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.
| Semantics | Behavior | Guarantees |
|---|---|---|
| At-Most-Once | Messages delivered zero or one time | Possible loss, no duplicates |
| At-Least-Once | Messages delivered one or more times | No loss, possible duplicates |
| Exactly-Once (EOS) | Messages delivered exactly once | No loss, no duplicates |
Setting Up Your Development Environment
Kafka Broker Setup
- Download and unzip Apache Kafka 2.5 or later.
- Start ZooKeeper (required by Kafka 2.5):
bin/zookeeper-server-start.sh config/zookeeper.properties
- Start Kafka broker:
bin/kafka-server-start.sh config/server.properties
- 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
| Property | Purpose |
|---|---|
enable.idempotence=true | Enable idempotent producer requests |
transactional.id | Uniquely identifies the transactional producer |
isolation.level=read_committed | Consumers 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.
KafkaTransactionManagerintegrates 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
- Start Kafka and your Spring Boot application.
- Use your
TransactionalKafkaProducerto send messages:
producer.sendMessageTransactional("transactional-topic", "key1", "Hello Kafka EOS");
- Verify the consumer logs each message exactly once.
- 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 Mode | Cause | Mitigation |
|---|---|---|
| Transaction timeouts | Transactions blocking beyond configured timeout | Keep transactions short. Adjust transaction.timeout.ms in broker. |
| Missing transactional.id | Producer is not configured for transactions | Set unique transactional.id per producer instance |
| Consumer reads uncommitted data | isolation.level set incorrectly | Use read_committed for consumers |
| Producer retries not idempotent | enable.idempotence is false or missing | Enable 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
- Apache Kafka Transactions Documentation
- Spring Kafka Reference Guide – Transactions
- Kafka Producer Configuration Documentation
- Spring Boot Kafka Support Documentation
