PrepZone Logo
PrepZone

Kafka Producers and Consumers

KafkaTemplate, @KafkaListener, serialization and error handling with DLT.

Read these first

Why this matters

  • When the inventory service, notification service, and analytics pipeline all need to know about a new order, point-to-point HTTP calls create a fragile web of dependencies.
  • Kafka persists messages to a log — consumers can restart, replay from an offset, and scale independently by adding members to a consumer group.
  • Spring for Apache Kafka wraps the raw client with KafkaTemplate and @KafkaListener, handling serialization and error hooks with minimal boilerplate.

Spring Kafka essentials

  • KafkaTemplate<K, V> — the producer API; send synchronously, asynchronously, or with callbacks.
  • @KafkaListener — declares a method that consumes from one or more topics.
  • Consumer groups — partition assignment and offset tracking; each message is delivered to one member per group.
  • Serializers — convert Java objects to bytes; JSON with Jackson is common for BookStore DTOs.
  • Dead-letter topic (DLT) — routes poison messages after retries exhaust, so one bad payload does not block the partition.
OrderServiceVaultCommerce
vaultcommerce.orders.placed.v1
InventoryConsumer
KafkaTemplate publishes to a topic. @KafkaListener consumes from a consumer group.

Dependencies and configuration

Java
<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>
Java
spring:
  kafka:
    bootstrap-servers: 127.0.0.1:9092
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
    consumer:
      group-id: bookstore-order-service
      auto-offset-reset: earliest
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
      properties:
        spring.json.trusted.packages: com.bookstore.events

Publishing order events

Java
public record OrderPlacedMessage(
        Long orderId,
        String customerId,
        List<OrderLineMessage> lines,
        Instant timestamp
) {}
Java
@Service
public class OrderEventPublisher {

    private final KafkaTemplate<String, OrderPlacedMessage> kafkaTemplate;

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

    public void publish(Order order) {
        OrderPlacedMessage message = OrderPlacedMessage.from(order);
        kafkaTemplate.send("bookstore.orders.placed", order.getId().toString(), message)
                .whenComplete((result, ex) -> {
                    if (ex != null) {
                        log.error("Failed to publish order {}", order.getId(), ex);
                    }
                });
    }
}

Consuming with error handling

Java
@Service
public class InventoryOrderConsumer {

    private final InventoryService inventoryService;

    public InventoryOrderConsumer(InventoryService inventoryService) {
        this.inventoryService = inventoryService;
    }

    @KafkaListener(
            topics = "bookstore.orders.placed",
            groupId = "bookstore-inventory"
    )
    public void onOrderPlaced(OrderPlacedMessage message) {
        inventoryService.reserveStock(message.orderId(), message.lines());
    }
}
Java
@KafkaListener(topics = "bookstore.orders.placed.DLT", groupId = "bookstore-dlt-handler")
public void handlePoisonMessage(OrderPlacedMessage message,
                                @Header(KafkaHeaders.RECEIVED_TOPIC) String topic) {
    log.error("Poison message from {} for order {}", topic, message.orderId());
    alertOpsTeam(message);
}
Java
@Bean
public ConcurrentKafkaListenerContainerFactory<String, OrderPlacedMessage>
        kafkaListenerContainerFactory(ConsumerFactory<String, OrderPlacedMessage> consumerFactory) {
    ConcurrentKafkaListenerContainerFactory<String, OrderPlacedMessage> factory =
            new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory);
    factory.setCommonErrorHandler(new DefaultErrorHandler(
            new DeadLetterPublishingRecoverer(kafkaTemplate),
            new FixedBackOff(1000L, 3)
    ));
    return factory;
}

Quick recall

Everything you need if you only revisit this box.

  • KafkaTemplate.send(topic, key, value) publishes; the key controls partition assignment.
  • @KafkaListener(topics = "...", groupId = "...") consumes; each group gets its own copy of the stream.
  • Configure JSON serializers in application.yml and whitelist your event package for deserialization.
  • Retry with DefaultErrorHandler + DeadLetterPublishingRecoverer to avoid stuck consumers on bad messages.
  • Application events are in-process; Kafka is for cross-service, durable, replayable messaging.

Test yourself

Answer these before moving on — recall is what makes it stick.