Command Palette

Search for a command to run...

Hectal
PHASE 7Intermediate ~8 min· topic 2 of 4

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.

0/4 · 0%

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

  1. 01

    Ack modes: RECORD, BATCH (default: commit after each poll's records are processed), MANUAL / MANUAL_IMMEDIATE (you call Acknowledgment.acknowledge()). With the default container settings, records are committed only after the listener returns without error.

  2. 02

    Error handling: DefaultErrorHandler(recoverer, new ExponentialBackOffWithMaxRetries(3)) retries a failing record in place (blocking) with backoff and then calls the recoverer, typically DeadLetterPublishingRecoverer which writes to <topic>-dlt with exception headers. Mark exceptions as not retryable with addNotRetryableExceptions.

  3. 03

    Deserialization: wrap your deserializer with ErrorHandlingDeserializer so bad bytes produce a DeserializationException routed to the error handler (and DLT) instead of an infinite poll-fail loop.

  4. 04

    Concurrency: concurrency: 3 creates three consumers in one application instance; total consumers across instances should not exceed partitions.

  5. 05

    Shutdown: Spring stops containers on context close, waiting shutdownTimeout for the listener to finish and committing offsets; in Kubernetes, give pods a terminationGracePeriodSeconds longer than that.

Code & diagrams

KafkaConsumerConfig.javajava
@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;
    }
}
PaymentListener.javajava
@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
}
consumer-application.ymlyaml
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.CooperativeStickyAssignor

When 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

01

Walk through what happens when the listener throws for a record under the configured error handler.

Practice

01

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.