Topic 7.2
A Production Spring Boot Consumer
In one line
A robust listener uses manual acknowledgement after processing, an ErrorHandlingDeserializer, a DefaultErrorHandler with backoff and a DeadLetterPublishingRecoverer, cooperative rebalancing, sensible concurrency, and graceful shutdown so in-flight records finish and offsets are committed.
Think of it like this
A well-run kitchen line. Each cook finishes a dish before ticking the ticket (manual ack), unreadable tickets go to the manager instead of stopping the line (DLT), a cook who burns a dish tries again after a pause (backoff), and at closing time cooks finish what's on the stove before leaving (graceful shutdown).
Key ideas
- 01
Ack modes:
RECORD,BATCH(default: commit after each poll's records are processed),MANUAL/MANUAL_IMMEDIATE(you callAcknowledgment.acknowledge()). With the default container settings, records are committed only after the listener returns without error. - 02
Error handling:
DefaultErrorHandler(recoverer, new ExponentialBackOffWithMaxRetries(3))retries a failing record in place (blocking) with backoff and then calls the recoverer, typicallyDeadLetterPublishingRecovererwhich writes to<topic>-dltwith exception headers. Mark exceptions as not retryable withaddNotRetryableExceptions. - 03
Deserialization: wrap your deserializer with
ErrorHandlingDeserializerso bad bytes produce aDeserializationExceptionrouted to the error handler (and DLT) instead of an infinite poll-fail loop. - 04
Concurrency:
concurrency: 3creates three consumers in one application instance; total consumers across instances should not exceed partitions. - 05
Shutdown: Spring stops containers on context close, waiting
shutdownTimeoutfor the listener to finish and committing offsets; in Kubernetes, give pods aterminationGracePeriodSecondslonger than that.
Code & diagrams
@Configuration
public class KafkaConsumerConfig {
@Bean
DefaultErrorHandler errorHandler(KafkaTemplate<Object, Object> template) {
var recoverer = new DeadLetterPublishingRecoverer(template,
(rec, ex) -> new TopicPartition(rec.topic() + "-dlt", rec.partition()));
var backoff = new ExponentialBackOffWithMaxRetries(3);
backoff.setInitialInterval(1_000);
backoff.setMultiplier(3.0);
var handler = new DefaultErrorHandler(recoverer, backoff);
handler.addNotRetryableExceptions(DeserializationException.class, ValidationException.class);
return handler;
}
@Bean
ConcurrentKafkaListenerContainerFactory<String, OrderEvent> kafkaListenerContainerFactory(
ConsumerFactory<String, OrderEvent> cf, DefaultErrorHandler eh) {
var f = new ConcurrentKafkaListenerContainerFactory<String, OrderEvent>();
f.setConsumerFactory(cf);
f.setCommonErrorHandler(eh);
f.setConcurrency(3);
f.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL);
f.getContainerProperties().setShutdownTimeout(30_000);
return f;
}
}@KafkaListener(topics = "commerce.orders", groupId = "payment-service")
public void onOrder(ConsumerRecord<String, OrderEvent> rec, Acknowledgment ack) {
MDC.put("orderId", rec.key());
log.info("topic={} partition={} offset={} key={} ts={}",
rec.topic(), rec.partition(), rec.offset(), rec.key(), rec.timestamp());
payments.handle(rec.value()); // idempotent (inbox or idempotency key)
ack.acknowledge(); // commit only after success
}spring:
kafka:
consumer:
group-id: payment-service
auto-offset-reset: earliest
enable-auto-commit: false
key-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
value-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
max-poll-records: 200
properties:
spring.deserializer.key.delegate.class: org.apache.kafka.common.serialization.StringDeserializer
spring.deserializer.value.delegate.class: org.springframework.kafka.support.serializer.JsonDeserializer
spring.json.value.default.type: com.acme.orders.OrderEvent
spring.json.trusted.packages: com.acme.orders
partition.assignment.strategy: org.apache.kafka.clients.consumer.CooperativeStickyAssignorWhen it breaks
spring.json.trusted.packages: "*" with type headers from untrusted producers
What you see
Consumers deserialize whatever class a header names, an unsafe deserialization surface.
Fix & prevent
Use a fixed default type, trust only your packages, and disable type headers on producers.
Explain it without notes
Walk through what happens when the listener throws for a record under the configured error handler.
Practice
Implement the consumer and force failures: a 503 that recovers, a validation error, and malformed JSON. Observe retries and DLT contents.
Trade-offs
- ↔
Blocking retries in the error handler keep order but pause the partition; use retry topics for long outages.
Done when you can
I can build a Spring consumer with manual acks, safe deserialization, retries, DLT and graceful shutdown.