A message payload arrives with an unparseable schema variation, throwing a silent serialization exception. Within seconds, your consumer loop stalls, offsets halt, partition lag cascades across the cluster, and rebalance storms trigger cascading socket timeouts across downstream services. Native Apache Kafka client primitives provide blistering wire speed, but handling distributed synchronization, thread lifecycles, and resilient recovery manually requires hundreds of lines of fragile plumbing.
Spring Kafka bridges the gap between raw Apache Kafka client throughput and enterprise application architecture. By orchestrating consumer threads, abstracting connection lifecycles via managed templates, and decoupling acknowledgment logic from business processing, it delivers rock-solid distributed message pipelines.
This guide deconstructs modern Spring Kafka implementations for 2026. You will master fine-grained consumer concurrency, non-blocking dead letter queues, virtual thread scaling under Spring Boot 3, and bulletproof integration testing using real containerized brokers.
Core Architecture: Bridging Spring Apache Kafka and Native Clients
The native Java client for Apache Kafka relies on a low-level polling loop executed on a single thread. The consumer continuously executes KafkaConsumer.poll(Duration), requiring developers to manually manage batch dispatching, offset commit timers, consumer rebalance listeners, and network disconnect recovery. Spring Kafka (formally spring apache kafka) wraps this procedural loop inside high-level, declarative abstractions while preserving direct access to underlying client configurations.
+--------------------------------------------------------------------------+
| Spring Application |
| +--------------------------------------------------------------------+ |
| | Business Logic Layer (@KafkaListener / KafkaTemplate) | |
| +-----------------------------------+--------------------------------+ |
| | |
| +-----------------------------------v--------------------------------+ |
| | MessageListenerContainer (Thread Lifecycle Management) | |
| | - Partition Assignment Listener - Dynamic Record/Batch Dispatch| |
| | - Transaction Coordinator - ErrorHandling / Recovery | |
| +-----------------------------------+--------------------------------+ |
| | |
| +-----------------------------------v--------------------------------+ |
| | Native Kafka Client Driver (org.apache.kafka.clients) |
| | - NetworkClient Selector - RecordAccumulator Buffers |
| +-----------------------------------+--------------------------------+ |
+--------------------------------------|-----------------------------------+
v
+----------------------------------+
| Apache Kafka Broker Cluster |
+----------------------------------+
At the architectural core sits the MessageListenerContainer. Rather than blocking your application code, Spring Kafka spawns dedicated polling threads that handle broker heartbeats and record buffering in the background. Received records are normalized into Spring Message<T> envelopes, routed through interceptor chains, and dispatched directly to your business logic methods.
Architecture Rule: Never execute long-running or blocking I/O calls directly inside a standard listener thread without adjusting
max.poll.interval.ms. If message processing exceeds this window, the coordinator marks the consumer dead and triggers an immediate partition rebalance.
| Architectural Dimension | Native Apache Kafka Client | Spring for Apache Kafka |
|---|---|---|
| Thread Model | Single-threaded loop; manual worker dispatch | Managed thread pool via listener containers |
| Offset Management | Manual commit calls or crude auto-commit intervals | Granular AckModes (RECORD, BATCH, MANUAL) |
| Error Handling | Manual try-catch around poll iteration | Built-in backoff, poison-pill trapping, and DLT routing |
| Object Serialization | Static raw byte serializer/deserializer contracts | Dynamic JSON/Avro/Protobuf mapping with type headers |
| Broker Testing | Requires custom external test infrastructure | Integrated Testcontainers and dynamic property registry |
Setting Up Build Dependencies via Spring Kafka Maven Artifacts
Establishing a stable baseline begins with clean dependency management. The spring kafka maven ecosystem relies on the Spring Boot dependencies Bill of Materials (BOM) to harmonize the spring-kafka module, Spring framework core libraries, and the official kafka-clients driver. Defining unmanaged client versions often introduces binary incompatibilities across protocol serializers.
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>3.3.4</version>
<relativePath/>
</parent>
<groupId>com.example.messaging</groupId>
<artifactId>kafka-resilient-pipeline</artifactId>
<version>1.0.0-SNAPSHOT</version>
<dependencies>
<-- Spring Kafka Starter -->
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
<-- JSON Processing -->
<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson-databind</artifactId>
</dependency>
<-- Observability & Metrics -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId>
</dependency>
<dependency>
<groupId>io.micrometer</groupId>
<artifactId>micrometer-registry-prometheus</artifactId>
</dependency>
<-- Modern Test Harness -->
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.testcontainers</groupId>
<artifactId>kafka</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
</project>
Verify your build pipeline against these criteria to prevent classpath collisions and serialization breakages:
- BOM Delegation: Let
spring-boot-starter-parentgovern thespring-kafkaandkafka-clientsartifacts to prevent API signature mismatches. - SLF4J Routing: Ensure native Kafka client loggers route cleanly through Logback without pulling in outdated log4j-over-slf4j bridges.
- Snappy/Zstd Compression: If enabling wire compression at the producer level, confirm native dynamic C-libraries compile or load cleanly on container runtime platforms like Alpine Linux.
End-to-End Spring Boot Kafka Example: Producer and Consumer Flow
A production-ready pipeline requires deterministic message serialization, dynamic error handling, and end-to-end tracing headers. The following spring boot kafka example presents a complete payload publishing and ingestion cycle using standard Spring Boot 3 configurations.
package com.example.messaging.model;
import java.math.BigDecimal;
import java.time.Instant;
public record OrderPlacedEvent(
String orderId,
String customerId,
BigDecimal amount,
Instant timestamp
) {}
Configure your producer factory with delivery guarantees, idempotent retries, and high-throughput batching:
package com.example.messaging.producer;
import com.example.messaging.model.OrderPlacedEvent;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.StringSerializer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
import org.springframework.kafka.support.serializer.JsonSerializer;
import java.util.HashMap;
import java.util.Map;
@Configuration
public class KafkaProducerConfig {
@Bean
public ProducerFactory<String, OrderPlacedEvent> orderProducerFactory() {
Map<String, Object> configProps = new HashMap<>();
configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
configProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
configProps.put(ProducerConfig.ACKS_CONFIG, "all");
configProps.put(ProducerConfig.RETRIES_CONFIG, 10);
configProps.put(ProducerConfig.LINGER_MS_CONFIG, 20);
configProps.put(ProducerConfig.BATCH_SIZE_CONFIG, 32768);
return new DefaultKafkaProducerFactory<>(configProps);
}
@Bean
public KafkaTemplate<String, OrderPlacedEvent> kafkaTemplate() {
return new KafkaTemplate<>(orderProducerFactory());
}
}
Execute the asynchronous publishing flow by attaching non-blocking completion callbacks:
package com.example.messaging.service;
import com.example.messaging.model.OrderPlacedEvent;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.stereotype.Service;
import java.util.concurrent.CompletableFuture;
@Service
public class OrderEventPublisher {
private static final Logger log = LoggerFactory.getLogger(OrderEventPublisher);
private static final String TOPIC = "orders.v1.events";
private final KafkaTemplate<String, OrderPlacedEvent> kafkaTemplate;
public OrderEventPublisher(KafkaTemplate<String, OrderPlacedEvent> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
public CompletableFuture<SendResult<String, OrderPlacedEvent>> publish(OrderPlacedEvent event) {
return kafkaTemplate.send(TOPIC, event.orderId(), event).whenComplete((result, ex) -> {
if (ex == null) {
log.info("Published event: orderId={}, offset={}",
event.orderId(),
result.getRecordMetadata().offset());
} else {
log.error("Failed publishing event: orderId={}", event.orderId(), ex);
}
});
}
}
To execute and verify this flow from ingestion to storage, follow these operational steps:
- Initialize Broker Topic: Create the topic
orders.v1.eventswith at least 3 partitions and a replication factor of 3. - Trigger Publisher: Call
OrderEventPublisher.publish()via an HTTP boundary or incoming event stream. - Observe Ack Lifecycle: Ensure the
CompletableFuturereceives acknowledgment confirmation directly from the Kafka cluster without blocking the servlet thread.
Engineering a Robust Spring Boot Kafka Consumer with Dynamic Concurrency
A high-performance spring boot kafka consumer balances resource utilization against partition availability. By default, a @KafkaListener runs single-threaded per container instance. If your topic has 12 partitions, an unconfigured consumer leaves 11 partitions idle within the active process.
Using ConcurrentKafkaListenerContainerFactory, you can allocate concurrent listener threads across partition assignments dynamically. In modern Spring Boot setups running on Java 21, you can back these consumers with virtual threads to minimize kernel thread context-switching overhead during downstream I/O calls.
package com.example.messaging.consumer;
import com.example.messaging.model.OrderPlacedEvent;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.task.SimpleAsyncTaskExecutor;
import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.listener.ContainerProperties;
import org.springframework.kafka.support.serializer.ErrorHandlingDeserializer;
import org.springframework.kafka.support.serializer.JsonDeserializer;
import java.util.HashMap;
import java.util.Map;
@EnableKafka
@Configuration
public class KafkaConsumerConfig {
@Bean
public ConsumerFactory<String, OrderPlacedEvent> consumerFactory() {
Map<String, Object> props = new HashMap<>();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "order-processing-group");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
props.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, JsonDeserializer.class.getName());
props.put(JsonDeserializer.TRUSTED_PACKAGES, "com.example.messaging.model");
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 100);
return new DefaultKafkaConsumerFactory<>(props);
}
@Bean
public ConcurrentKafkaListenerContainerFactory<String, OrderPlacedEvent> kafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, OrderPlacedEvent> factory =
new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory());
factory.setConcurrency(3);
factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
// Back listener containers with Java 21 Virtual Threads
SimpleAsyncTaskExecutor executor = new SimpleAsyncTaskExecutor("k-virtual-");
executor.setVirtualThreads(true);
factory.getContainerProperties().setListenerTaskExecutor(executor);
return factory;
}
}
package com.example.messaging.consumer;
import com.example.messaging.model.OrderPlacedEvent;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.messaging.handler.annotation.Payload;
import org.springframework.stereotype.Component;
@Component
public class OrderEventListener {
private static final Logger log = LoggerFactory.getLogger(OrderEventListener.class);
@KafkaListener(
topics = "orders.v1.events",
groupId = "order-processing-group",
containerFactory = "kafkaListenerContainerFactory"
)
public void handleOrder(
@Payload OrderPlacedEvent event,
@Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
@Header(KafkaHeaders.OFFSET) long offset,
Acknowledgment ack) {
log.info("Processing order: id={}, partition={}, offset={}", event.orderId(), partition, offset);
// Business processing logic executed here
// Explicit manual commit on processing success
ack.acknowledge();
}
}
| Scaling Lever | Default Setting | Optimized Production Value | Engineering Impact |
|---|---|---|---|
| Container Concurrency | 1 | Match partition count per instance | Eliminates idle partitions across parallel consumer workers |
| Virtual Threads | false (Platform Threads) | true (SimpleAsyncTaskExecutor) | Eliminates thread starvation during downstream HTTP/DB I/O |
| max.poll.records | 500 | 50 to 100 | Prevents partition timeouts caused by slow batch processing |
| Partition Assignment | RangeAssignor | CooperativeStickyAssignor | Prevents full stop-the-world rebalance storms across nodes |
Offset Commit Strategies and Manual Acknowledgment Models
Relying on Kafka’s automatic background offset commits (enable.auto.commit=true) is a primary cause of silent data loss. Background commits occur strictly on wall-clock intervals (governed by auto.commit.interval.ms), completely decoupled from whether your downstream database transaction succeeded or crashed.
Spring Kafka introduces the ContainerProperties.AckMode enumeration, handing offset management authority directly to your architectural requirements.
| AckMode | Commit Mechanics | Latency Impact | Data Loss Risk | Duplicate Risk |
|---|---|---|---|---|
| RECORD | Commits after each individual record processes | High (Frequent broker sync) | Zero | Minimal |
| BATCH | Commits the entire poll batch once processed | Very Low (Batched write) | High on unhandled crash | High (Full batch replay) |
| MANUAL | Queues offset when ack.acknowledge() runs; flushes on next poll |
Low | Low | Moderate |
| MANUAL_IMMEDIATE | Commits synchronously to broker the instant ack.acknowledge() executes |
Moderate | Zero | Low |
| TIME / COUNT | Flushes commits when time elapses or count threshold is met | Configurable | Moderate | Moderate |
Production Best Practice: For mission-critical transactional pipelines like financial ledgers or order processing, select
MANUAL_IMMEDIATEcombined with transactional database outbox updates or idempotent writes. For telemetry or metrics ingestion pipelines, useBATCHacknowledgment to maximize overall ingest throughput.
Poison Pill Defense: Dead Letter Publishing and Backoff Recovery
A single malformed JSON payload, missing required schema attributes, or incompatible bytecode structure acts as a poison pill. If a consumer fails during deserialization, native Kafka throws an unhandled exception before reaching your @KafkaListener method. In a naive setup, the offset never moves forward, leading to an infinite retry loop that permanently stalls partition consumption.
Resilient architectures solve this through a layered strategy: configure an ErrorHandlingDeserializer at the perimeter to catch parsing failures, apply a DefaultErrorHandler with exponential backoff for transient issues, and forward persistent failures to a Dead Letter Topic (DLT) via DeadLetterPublishingRecoverer.
package com.example.messaging.error;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.common.TopicPartition;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.KafkaOperations;
import org.springframework.kafka.listener.CommonErrorHandler;
import org.springframework.kafka.listener.DeadLetterPublishingRecoverer;
import org.springframework.kafka.listener.DefaultErrorHandler;
import org.springframework.util.backoff.ExponentialBackOff;
@Configuration
public class KafkaErrorRecoveryConfig {
private static final Logger log = LoggerFactory.getLogger(KafkaErrorRecoveryConfig.class);
@Bean
public CommonErrorHandler errorHandler(KafkaOperations<Object, Object> kafkaOperations) {
// Route failed records to a topic with ".DLT" appended to the original topic name
DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(
kafkaOperations,
(ConsumerRecord<,> cr, Exception ex) -> {
log.error("Routing record to DLT: key={}, topic={}, partition={}, offset={}, cause={}",
cr.key(), cr.topic(), cr.partition(), cr.offset(), ex.getMessage());
return new TopicPartition(cr.topic() + ".DLT", cr.partition());
}
);
// 3 retry attempts: 1s initial delay, 2.0 multiplier, 10s max backoff interval
ExponentialBackOff backOff = new ExponentialBackOff(1000L, 2.0);
backOff.setMaxElapsedTime(10000L);
DefaultErrorHandler errorHandler = new DefaultErrorHandler(recoverer, backOff);
// Do not retry poison pills or parsing errors; forward immediately to DLT
errorHandler.addNotRetryableExceptions(
org.springframework.kafka.support.serializer.DeserializationException.class,
com.fasterxml.jackson.core.JsonParseException.class,
java.lang.IllegalArgumentException.class
);
return errorHandler;
}
}
Confirm your resilience design against this production readiness checklist:
- ErrorHandlingDeserializer Wrapper: Configured on both key and value deserializers to capture byte corruption without crashing the container thread.
- Non-Retryable Classification: Deserialization and payload validation exceptions are explicitly flagged as non-retryable to bypass pointless backoff cycles.
- DLT Partition Parity: Target dead letter topics must match or exceed the partition count of the source topic to prevent routing failures when preserving partition mapping.
- Audit Headers: Verify the recoverer appends diagnostic headers (e.g.
KafkaHeaders.DLT_EXCEPTION_MESSAGE,KafkaHeaders.DLT_ORIGINAL_OFFSET) for rapid incident triaging.
Integration Testing Pipelines Using Testcontainers and Awaitility
Historically, teams relied on @EmbeddedKafka for local integration tests. However, embedded in-memory brokers run within JVM constraints, masking genuine container network timeouts, client protocol negotiations, and real Kraft quorum behaviors. Modern pipelines use real, containerized Apache Kafka instances driven by Testcontainers and verified asynchronously with Awaitility.
package com.example.messaging;
import com.example.messaging.model.OrderPlacedEvent;
import com.example.messaging.producer.OrderEventPublisher;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.support.serializer.JsonDeserializer;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.kafka.KafkaContainer;
import java.math.BigDecimal;
import java.time.Duration;
import java.time.Instant;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.await;
@SpringBootTest
@Testcontainers
class OrderPipelineIntegrationTest {
@Container
static KafkaContainer kafka = new KafkaContainer("apache/kafka-native:3.8.0");
@DynamicPropertySource
static void overrideProperties(DynamicPropertyRegistry registry) {
registry.add("spring.kafka.bootstrap-servers", kafka:getBootstrapServers);
}
@Autowired
private OrderEventPublisher publisher;
private Consumer<String, OrderPlacedEvent> testConsumer;
@BeforeEach
void setUp() {
Map<String, Object> props = new HashMap<>();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafka.getBootstrapServers());
props.put(ConsumerConfig.GROUP_ID_CONFIG, "test-verify-group-" + UUID.randomUUID());
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
props.put(JsonDeserializer.TRUSTED_PACKAGES, "*");
DefaultKafkaConsumerFactory<String, OrderPlacedEvent> factory = new DefaultKafkaConsumerFactory<>(props);
testConsumer = factory.createConsumer();
testConsumer.subscribe(Collections.singletonList("orders.v1.events"));
}
@AfterEach
void tearDown() {
if (testConsumer!= null) {
testConsumer.close();
}
}
@Test
void shouldPublishAndConsumeOrderEventSuccessfully() {
String orderId = "ORD-" + UUID.randomUUID();
OrderPlacedEvent event = new OrderPlacedEvent(orderId, "CUST-99", new BigDecimal("149.99"), Instant.now());
publisher.publish(event);
await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> {
ConsumerRecords<String, OrderPlacedEvent> records = testConsumer.poll(Duration.ofMillis(200));
assertThat(records).isNotEmpty();
ConsumerRecord<String, OrderPlacedEvent> matched = null;
for (ConsumerRecord<String, OrderPlacedEvent> record: records) {
if (orderId.equals(record.value().orderId())) {
matched = record;
break;
}
}
assertThat(matched).isNotNull();
assertThat(matched.value().amount()).isEqualByComparingTo("149.99");
});
}
}
Structure your automated test execution around these validation phases:
- Spin up Ephemeral Cluster: Start the Kafka container via Testcontainers, binding its assigned dynamic port into Spring’s
DynamicPropertyRegistry. - Dispatch Domain Events: Use the application’s actual Spring
KafkaTemplateor domain publisher service to dispatch records. - Assert Non-Blocking State: Avoid hardcoded
Thread.sleep()calls. Instead, use Awaitility to continuously poll and assert state, verifying partition allocation and deserialization resilience cleanly.
Frequently Asked Questions
What is the difference between pure Apache Kafka client and Spring Kafka?
Spring Kafka wraps the native Java Kafka client, adding declarative programming through annotations like @KafkaListener, managed thread containers, automated template abstractions, and integrated error handling, reducing boilerplate infrastructure code significantly in Spring-managed applications.
How do I add Spring Kafka using Maven?
Add the spring-kafka dependency directly to your pom.xml file under dependencies. When using Spring Boot, inherit versions from spring-boot-starter-parent to automatically manage compatible versions of the native Kafka client without manual version overrides.
How do I scale a Spring Boot Kafka consumer for high throughput?
Scale consumers by increasing the concurrency property in ConcurrentKafkaListenerContainerFactory to match partition counts, configuring batch listener mode, enabling asynchronous acknowledgments, or utilizing Java 21 virtual threads via Spring Boot properties.
How does Spring Kafka prevent consumer poison pills?
Configure an ErrorHandlingDeserializer wrapping your payload deserializer alongside a DefaultErrorHandler. Unparseable records are routed directly to a Dead Letter Topic via DeadLetterPublishingRecoverer, preventing continuous offset processing failure loops.
Building resilient, high-throughput streaming systems requires mastering the boundary between raw Kafka broker protocols and Spring’s concurrency abstractions. By combining explicit manual acknowledgment modes, virtual-thread container factories, and automated Dead Letter Topic recoverers, you safeguard your pipelines against data loss, cascade rebalances, and poison-pill crashes.
As you scale these patterns across production clusters, pair your error-handling and concurrency setups with Micrometer metrics and OpenTelemetry tracing. Continuous observability across partition lag, rebalance durations, and deserialization faults ensures your distributed event architecture remains predictable, stable, and resilient under any workload.