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;
}
}