Implementing Messaging with Kafka

1. Adding Kafka Dependencies

Example: Kafka dependency

<dependency>
  <groupId>org.springframework.kafka</groupId>
  <artifactId>spring-kafka</artifactId>
</dependency>

2. Configuring Kafka Producer Properties

Example: Kafka producer with idempotence

spring:
  kafka:
    bootstrap-servers: localhost:9092
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
      acks: all
      retries: 5
      properties:
        enable.idempotence: true
        compression.type: lz4

3. Configuring Kafka Consumer Properties

Example: Kafka consumer with manual commit

spring:
  kafka:
    consumer:
      group-id: orders-consumer
      auto-offset-reset: earliest
      enable-auto-commit: false
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
      properties:
        spring.json.trusted.packages: "*"
    listener:
      ack-mode: MANUAL_IMMEDIATE

4. Creating KafkaTemplate for Sending Messages

Example: Send Kafka message with result callback

@Service
public class OrderPublisher {
  private final KafkaTemplate<String, OrderPlaced> kafka;
  public OrderPublisher(KafkaTemplate<String, OrderPlaced> kafka) { this.kafka = kafka; }
  public void publish(OrderPlaced evt) {
    kafka.send("orders.placed", String.valueOf(evt.orderId()), evt)
         .whenComplete((res, ex) -> { if (ex != null) log.error("send failed", ex); });
  }
}

5. Sending Messages to Topics

Method Use
send(topic, value) Round-robin partition
send(topic, key, value) Hash key → partition
send(topic, partition, key, value) Explicit partition
send(ProducerRecord<K,V>) Headers + custom timestamp
executeInTransaction(t -> ...) Transactional producer

6. Creating Kafka Listeners (@KafkaListener)

Example: Consume Kafka record with manual ack

@Component
public class OrderConsumer {
  @KafkaListener(topics = "orders.placed", groupId = "orders-consumer", concurrency = "3")
  public void handle(ConsumerRecord<String, OrderPlaced> rec, Acknowledgment ack) {
    process(rec.value());
    ack.acknowledge();
  }
}

7. Configuring Consumer Groups

Aspect Effect
Same group-id Partitions split among consumers
Different group-id Each group gets all messages
Rebalance Triggered by membership/topic change

8. Handling Message Serialization

Serializer Use
StringSerializer Plain strings
JsonSerializer POJOs as JSON
ByteArraySerializer Raw bytes
Avro/Protobuf Schema-registry backed
ErrorHandlingDeserializer Wrap to bypass poison messages

9. Managing Offset Commits

ack-mode Behavior
RECORD After each record
BATCH (default) After poll batch
TIME After interval
COUNT After N records
MANUAL Acknowledgment.acknowledge()
MANUAL_IMMEDIATE Commit synchronously immediately

10. Implementing Error Handlers and Retry

Example: Kafka dead letter publisher on failure

@Bean
DefaultErrorHandler errorHandler(KafkaTemplate<String,Object> template) {
  DeadLetterPublishingRecoverer dlr = new DeadLetterPublishingRecoverer(template);
  return new DefaultErrorHandler(dlr, new FixedBackOff(1000L, 3));
}

11. Using Kafka Headers and Timestamps

Example: Access Kafka headers in listener

@KafkaListener(topics="orders.placed")
public void handle(@Payload OrderPlaced evt,
                   @Header(KafkaHeaders.RECEIVED_KEY) String key,
                   @Header(KafkaHeaders.RECEIVED_TIMESTAMP) long ts,
                   @Headers Map<String,Object> headers) { /* … */ }

12. Using Kafka Streams for Processing

Example: Kafka Streams filter to high-value topic

@Configuration
@EnableKafkaStreams
public class StreamsConfig {
  @Bean
  KStream<String, OrderPlaced> pipeline(StreamsBuilder b) {
    KStream<String,OrderPlaced> s = b.stream("orders.placed");
    s.filter((k,v) -> v.total().compareTo(BigDecimal.valueOf(1000)) > 0)
     .to("orders.high-value");
    return s;
  }
}