SPRING
BOOT
KAFKA
Apache Kafka is a distributed event streaming platform designed for high-throughput, fault-tolerant, and real-time data processing. Originally developed by the engineering team at LinkedIn, Kafka was later contributed to the Apache Software Foundation under the Apache License, and is now widely known as Apache Kafka—an open-source solution for building scalable event-driven architectures.
What Kafka Does
Kafka allows applications to send, store, and read streams of data (events/messages) in real time. It’s widely used for:
- check_circle Messaging system – like RabbitMQ, but much more scalable.
- check_circle Event streaming – continuously process data streams.
- check_circle Data pipeline – move data between systems reliably.
- check_circle Log aggregation – collect logs from multiple servers.
Core Concepts

- check_circle Topic: A category or feed name where messages are published.
- check_circle Partition: A topic can be split into multiple partitions to allow parallelism and scalability.
- check_circle Producer: Application that sends messages to Kafka topics.
- check_circle Consumer: Application that reads messages from Kafka topics.
- check_circle Broker: A Kafka server that stores messages and serves clients.
- check_circle Cluster: Multiple brokers working together to handle more data and provide redundancy.
- check_circle ZooKeeper (old versions): Used to manage brokers in the cluster (modern Kafka can work without it).
Analogy
Imagine a mail system:
- check_circle Topics are like mailboxes.
- check_circle Producers are people sending letters.
- check_circle Consumers are people reading letters.
- check_circle Partitions are like sub-mailboxes for parallel delivery.
- check_circle Kafka ensures every letter is stored reliably and can be read by multiple people independently.
Lets Code It!
I will demonstrate spring boot service using kafka to produce message and listen message and how to handling success/failure case in simple way.
Preparation:
- check_circle Java 17.
- check_circle Spring Boot 3.5.6.
- check_circle Kafka 4.1.0.
- check_circle Girlfriend (if any)
Lets start with Producer Service!
Config:
@EnableKafka
@Configuration
public class KafkaProducerConfig {
@Value(value = "${spring.kafka.bootstrap-servers}")
private String bootstrapAddress;
@Value(value = "${spring.kafka.consumer.group-id}")
private String groupId;
@Bean
public ProducerFactory<String, String> producerFactory() {
Map<String, Object> configProps = new HashMap<>();
configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
configProps.put(ProducerConfig.RETRIES_CONFIG, 3);// Number of retries
configProps.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 60000);// 60 seconds total timeout
configProps.put(ProducerConfig.RETRY_BACKOFF_MS_CONFIG, 100);// Time between retries
configProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);// Enable idempotence
configProps.put(ProducerConfig.ACKS_CONFIG, "all");
configProps.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5); // set it to One if Ordering of messages are important
return new DefaultKafkaProducerFactory<>(configProps);
}
@Bean
public KafkaTemplate<String, String> kafkaTemplate() {
return new KafkaTemplate<>(producerFactory());
}
}Publisher Service:
@Slf4j
@Service
@RequiredArgsConstructor
public class KafkaPublisher {
private final KafkaTemplate<String, String> kafkaTemplate;
private final KafkaHelper kafkaHelper;
public void sendMessage(String message) {
CompletableFuture<SendResult<String, String>> future = kafkaTemplate.send(TOPIC_GENERAL, message);
future.whenComplete((result, ex) -> {
if (ex == null) {
kafkaHelper.onSuccess(result, message);
} else {
kafkaHelper.handleException(ex, message, TOPIC_GENERAL);
}
});
}
public void sendMessageToTopic(String topic, String message) {
CompletableFuture<SendResult<String, String>> future = kafkaTemplate.send(topic, message);
future.whenComplete((result, ex) -> {
if (ex == null) {
kafkaHelper.onSuccess(result, message);
} else {
kafkaHelper.handleException(ex, message, topic);
}
});
}
}Handler Success/Failure:
@Slf4j
@Service
@RequiredArgsConstructor
public class KafkaHelper {
private final KafkaTemplate<String, String> kafkaTemplate;
public <T> void onSuccess(final SendResult<String, String> result, final T t) {
log.info("Sent Message={} to topic-partition={}-{} with offset={}",
t.toString(),
result.getRecordMetadata().topic(),
result.getRecordMetadata().partition(),
result.getRecordMetadata().offset());
}
public <T> void onFailure(final Throwable ex, final T t) {
log.error("Unable to send message={} due to {}", t.toString(), ex.getMessage());
}
public void handleException(Throwable ex, String payload, final String targetTopic) {
log.error("Unable to send message={} due to {}", payload, ex.getMessage());
if (ex instanceof UnsupportedVersionException
|| ex instanceof RecordTooLargeException
|| ex instanceof CorruptRecordException) {
log.error("Non-Retriable exception occurred. Sending message to dead-letter queue: [{}] Error: {}", payload, ex.getMessage());
sendToDeadLetterTopic(payload, targetTopic);
} else if (ex instanceof SerializationException) {
log.error("Serialization error! message=[{}]", payload);
} else if (ex instanceof RetriableException) {
log.warn("Retriable exception occurred. Kafka will retry automatically: {}", ex.getMessage());
}
}
private void sendToDeadLetterTopic(String payload, final String topic) {
String deadLetterTopic = topic + ".DLT";
log.info("Sending failed message to Dead Letter Topic: {}", deadLetterTopic);
kafkaTemplate.send(deadLetterTopic, payload);
}
}Now we move to Consumer Service!
Config:
@EnableKafka
@Configuration
public class KafkaConsumerConfig {
@Value(value = "${spring.kafka.bootstrap-servers}")
private String bootstrapAddress;
@Value(value = "${spring.kafka.consumer.group-id}")
private String groupId;
@Bean
public ConsumerFactory<String, String> consumerFactory() {
Map<String, Object> props = new HashMap<>();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
return new DefaultKafkaConsumerFactory<>(props);
}
@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory());
return factory;
}
}Listener Service
@Slf4j
@Service
@RequiredArgsConstructor
public class KafkaListenerService {
@KafkaListener(topics = TOPIC_GENERAL, groupId = "${spring.kafka.consumer.group-id}")
public void listenMessage(String message) {
try {
log.info("Incoming Message: {}", message);
} catch (Exception e) {
log.error("error when listen message: {}", e.getMessage());
}
}
TypeReference<KafkaBaseMessage<CustomerInfo>> typeRef = new TypeReference<KafkaBaseMessage<CustomerInfo>>() {};
@KafkaListener(topics = TOPIC_CUSTOMER, groupId = "${spring.kafka.consumer.group-id}")
public void listenCustomerMessage(String message) {
try {
log.info("Incoming Message: {}", message);
KafkaBaseMessage<CustomerInfo> baseMessage = (new ObjectMapper()).readValue(message, typeRef);
log.info("Result message: " + baseMessage);
log.info("Customer: " + baseMessage.getData());
} catch (Exception e) {
log.error("error when listen message: {}", e.getMessage());
}
}
}Run Kafka
Run Kafka with config in config/server.properties.
Im running my kafka on port 9094.

Handbook kafka command
Jump to folder Kafka
cd kafka_2.13-4.1.0
Run kafka:
KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"
bin/kafka-storage.sh format --standalone -t $KAFKA_CLUSTER_ID -c config/server.properties
bin/kafka-server-start.sh config/server.properties
Stop Kafka:
bin/kafka-server-stop.sh
Delete log before re format:
rm -rf /tmp/kraft-combined-logs
Create topic:
bin/kafka-topics.sh --create --topic <topic-name> --bootstrap-server localhost:9092
Sample:
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--create --topic <topic-name> \
--partitions 3 \
--replication-factor 1 \
--config retention.ms=86400000
List display topic:
bin/kafka-topics.sh --describe --topic <topic-name> --bootstrap-server localhost:9092
Testing
Hit Producer Service with Postman

Sample log in Producer Service:
16-10-2025 10:45:07.886 [kafka-producer-network-thread | spring-kafka-producer-producer-1] [app_name=spring-kafka-producer, trace_id=, span_id=] INFO c.p.s.service.helper.KafkaHelper.onSuccess - Sent Message=HIDOEP JOEKAWAI 996 to topic-partition=message-string-topic-0 with offset=996
16-10-2025 10:45:07.886 [kafka-producer-network-thread | spring-kafka-producer-producer-1] [app_name=spring-kafka-producer, trace_id=, span_id=] INFO c.p.s.service.helper.KafkaHelper.onSuccess - Sent Message=HIDOEP JOEKAWAI 997 to topic-partition=message-string-topic-0 with offset=997
16-10-2025 10:45:07.886 [kafka-producer-network-thread | spring-kafka-producer-producer-1] [app_name=spring-kafka-producer, trace_id=, span_id=] INFO c.p.s.service.helper.KafkaHelper.onSuccess - Sent Message=HIDOEP JOEKAWAI 998 to topic-partition=message-string-topic-0 with offset=998
16-10-2025 10:45:07.887 [kafka-producer-network-thread | spring-kafka-producer-producer-1] [app_name=spring-kafka-producer, trace_id=, span_id=] INFO c.p.s.service.helper.KafkaHelper.onSuccess - Sent Message=HIDOEP JOEKAWAI 999 to topic-partition=message-string-topic-0 with offset=999Sample log in Consumer Service:
16-10-2025 10:45:08.052 [org.springframework.kafka.KafkaListenerEndpointContainer#2-0-C-1] [app_name=spring-kafka-consumer, trace_id=, span_id=] INFO c.p.s.s.s.kafka.KafkaListenerService.listenMessage - Incoming Message: HIDOEP JOEKAWAI 996
16-10-2025 10:45:08.052 [org.springframework.kafka.KafkaListenerEndpointContainer#2-0-C-1] [app_name=spring-kafka-consumer, trace_id=, span_id=] INFO c.p.s.s.s.kafka.KafkaListenerService.listenMessage - Incoming Message: HIDOEP JOEKAWAI 997
16-10-2025 10:45:08.053 [org.springframework.kafka.KafkaListenerEndpointContainer#2-0-C-1] [app_name=spring-kafka-consumer, trace_id=, span_id=] INFO c.p.s.s.s.kafka.KafkaListenerService.listenMessage - Incoming Message: HIDOEP JOEKAWAI 998
16-10-2025 10:45:08.053 [org.springframework.kafka.KafkaListenerEndpointContainer#2-0-C-1] [app_name=spring-kafka-consumer, trace_id=, span_id=] INFO c.p.s.s.s.kafka.KafkaListenerService.listenMessage - Incoming Message: HIDOEP JOEKAWAI 999
That's it. You can find source code here:
Producer Service
Consumer Service
Reference: