목차
Spring boot에서 kafka 구축
본문의 JDK 8·Kafka 2.x·Spring Boot 예제는 2019년 실습 기록이다. 당시의 ZooKeeper 명령과 Gradle/라이브러리 코드를 현재 버전에서 그대로 실행하면 맞지 않을 수 있다. 특히
GET /send로 데이터를 쓰고GET /receiver가 마지막 메시지 하나를 공유 필드에서 읽는 구조는 운영 API로 쓰기 어렵다. 글 끝에 현재 설계 기준과 간결한 예제를 덧붙였다.
Spring Boot를 이용하여 kafka API를 사용해보고자 한다.
사전조건
-
kafka 설치 상태
-
개발도구 설치 상태
개발환경
-
jdk : 1.8
-
IDE : intelliJ
-
build : gradle
-
kafka 서버는 Ubuntu 18.04LTS에 설치되어 있다.
준비
- Producer를 통해 메시지를 보내고 Consumer를 통해 topic에 있는 데이터를 받아오는 프로그램을 간단하게 만들어보자.
| URL | 내용 |
|---|---|
| /send?msg= | Kafka source-topic에 메시지를 보낸다. |
| /receiver | Kafka source-topic에 들어간 메시지 한개 받는다. |
진행
1. zooKeeper, kafka 실행
- zooKeeper, kafka 실행
Zookeeper실행
$ ./bin/kafka-zookeeper-server-start.sh config/zookeeper.properties
Kafka 실행
$ ./bin/kafka-server-start.sh config/server.properties
2. Kafka topic 생성
- Kafka topic 생성
- source-topic을 생성한다.
$ ./bin/kafka-topics.sh --create --replication-factor 1 --partitions 1 --topic source-topic
3. 구현
- 구현
1. 프로젝트 생성
File -> module or project 선택

2. Project Metadata 설정

3. Dependencies 설정
-
Web > Spring Web Starter 선택
-
Developer Tools > Spring Boot DevTools, Lombok 선택
-
Messaging > Spring for Apache kafka, Spring for ApacheKafka Streams 선택

4. 프로젝트 명 설정

5. Gradle 설정

6. kafka Producer, Consumer Properties 설정
- 스프링 부트 프로젝트를 생성하면 application.properties 가 생성되는데 application.properties는 지우고 application.yaml으로 바꿔서 사용했다.
아래와 같이 설정
[application.yaml]
## kafka Producer
spring:
kafka:
producer:
bootstrap-servers: localhost:9092
acks: all
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.apache.kafka.common.serialization.StringSerializer
batch-size: 1000
consumer:
bootstrap-servers: localhost:9092
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
group-id: SpringKafka
---
server:
port: 8096
---
kafka:
topic:
source-topic: source-topic

7. Producer, Consumer Config 설정
- @EnableKafka : Kafka를 사용할 수 있도록 자동으로 와이어링 해준다.
@EnableKafka
@Configuration
public class KafkaConfig {
}
@EnableKafka@Configurationpublic class KafkaConfig {}
- Producer 설정
- producerProps()는 Producer를 실행하기위한 정보를 설정한다.
- ProducerFactory는 properties의 정보를 넣고 transaction을 설정할 수 있다.
- KafkaTemplate은 Kafka broker의 topic으로 데이터를 전송하도록 도와주는 역할을 한다.
@Bean
public Map<String, Object> producerProps() {
Map<String, Object> props = new HashMap<>();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, env.getProperty("spring.kafka.producer.bootstrap-servers"));
props.put(ProducerConfig.ACKS_CONFIG, env.getProperty("spring.kafka.producer.acks"));
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, env.getProperty("spring.kafka.producer.key-serializer"));
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, env.getProperty("spring.kafka.producer.value-serializer"));
props.put(ProducerConfig.RETRIES_CONFIG, 0);
return props;
}
@Bean
public ProducerFactory<String, String> producerFactory() {
return new DefaultKafkaProducerFactory<String, String>(producerProps());
}
@Bean
public KafkaTemplate<String, String> kafkaTemplate() {
return new KafkaTemplate<>(producerFactory());
}
- Consumer설정
- consumerProps()는 Consumer 환경설정 정보를 설정해준다.
- ConsumerFactory는 consumerProps()의 정보를 초기화를 해준다.
- KafkaListenerContainerFactory는 여러 개의 컨테이너를 만들 수 있고, 여러개의 factor를 구성할 수도 있다. ConcurrentMessageListenerContainer는 @KafkaListener 를 사용할 수 있게되고 해당 어노테이션을 통해 카프카 Consumer 데이터를 가져올 수 있다.
public Map<String, Object> consumerProps() {
Map<String, Object> props = new HashMap<>();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, env.getProperty("spring.kafka.consumer.bootstrap-servers"));
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, env.getProperty("spring.kafka.consumer.key-deserializer"));
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, env.getProperty("spring.kafka.consumer.value-deserializer"));
props.put(ConsumerConfig.GROUP_ID_CONFIG, env.getProperty("spring.kafka.consumer.group-id"));
return props;
}
@Bean
public ConsumerFactory<Integer, String> consumerFactory() {
return new DefaultKafkaConsumerFactory<>(consumerProps());
}
@Bean
KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<Integer, String>> kafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<Integer, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory());
factory.setConcurrency(3);
factory.getContainerProperties().setPollTimeout(3000);
return factory;
}
[KafkaConfig.java]
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.env.Environment;
import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.config.KafkaListenerContainerFactory;
import org.springframework.kafka.core.*;
import org.springframework.kafka.listener.ConcurrentMessageListenerContainer;
import java.util.HashMap;
import java.util.Map;
@EnableKafka
@Configuration
public class KafkaConfig {
private Environment env;
@Autowired
public KafkaConfig(Environment env) {
this.env = env;
}
/**
* Producer Properties
*
* @return
*/
@Bean
public Map<String, Object> producerProps() {
Map<String, Object> props = new HashMap<>();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, env.getProperty("spring.kafka.producer.bootstrap-servers"));
props.put(ProducerConfig.ACKS_CONFIG, env.getProperty("spring.kafka.producer.acks"));
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, env.getProperty("spring.kafka.producer.key-serializer"));
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, env.getProperty("spring.kafka.producer.value-serializer"));
props.put(ProducerConfig.RETRIES_CONFIG, 0);
return props;
}
/**
* Producer Factory
*
* @return
*/
@Bean
public ProducerFactory<String, String> producerFactory() {
return new DefaultKafkaProducerFactory<String, String>(producerProps());
}
/**
* KafkaTemplate는 kafka broker의 topic으로 데이터를 전송하도록 도와주는 역할을 한다.
*
*
* @return
*/
@Bean
public KafkaTemplate<String, String> kafkaTemplate() {
return new KafkaTemplate<>(producerFactory());
}
/**
* Consumer Properties
*
* @return
*/
@Bean
public Map<String, Object> consumerProps() {
Map<String, Object> props = new HashMap<>();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, env.getProperty("spring.kafka.consumer.bootstrap-servers"));
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, env.getProperty("spring.kafka.consumer.key-deserializer"));
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, env.getProperty("spring.kafka.consumer.value-deserializer"));
props.put(ConsumerConfig.GROUP_ID_CONFIG, env.getProperty("spring.kafka.consumer.group-id"));
return props;
}
/**
* Consumer Factory
*
* @return
*/
@Bean
public ConsumerFactory<Integer, String> consumerFactory() {
return new DefaultKafkaConsumerFactory<>(consumerProps());
}
/**
* Consumer Container factory
*ConcurrentKafkaListenerContainerFactory는
* @return
*/
@Bean
KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<Integer, String>> kafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<Integer, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory());
factory.setConcurrency(3);
factory.getContainerProperties().setPollTimeout(3000);
return factory;
}
}
8. Producer 구현
- Producer는 KafkaTemplate를 통해 메시지를 해당 Topic에게 전송할 수 있다.
kafkaTemplate.send(topic, message); 이 Kafka 서버로 해당 Topic으로 데이터를 전송하는 역할을 하고 ListenableFuture를 통해 Collback을 받을 수 있다.
ListenableFuture<SendResult<String, String>> future = kafkaTemplate.send(topic, message);
future.addCallback(new ListenableFutureCallback<SendResult<String, String>>() {
@Override
public void onSuccess(SendResult<String, String> result) {
log.info("Send Message : [{}] with offset=[{}]", message, result.getRecordMetadata().offset());
}
@Override
public void onFailure(Throwable ex) {
log.info("fail Message : {}, Exception : {}", message, ex.getMessage());
}
});
[ProducerService.java]
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.stereotype.Service;
import org.springframework.util.concurrent.ListenableFuture;
import org.springframework.util.concurrent.ListenableFutureCallback;
@Slf4j
@Service
public class ProducerService {
private final KafkaTemplate<String, String> kafkaTemplate;
@Autowired
public ProducerService(KafkaTemplate<String, String> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
@Value("${kafka.topic.source-topic}")
private String topic;
public void send(String message) {
log.info("Send Producer message : {}, topic={}", message, topic);
ListenableFuture<SendResult<String, String>> future = kafkaTemplate.send(topic, message);
future.addCallback(new ListenableFutureCallback<SendResult<String, String>>() {
@Override
public void onSuccess(SendResult<String, String> result) {
log.info("Send Message : [{}] with offset=[{}]", message, result.getRecordMetadata().offset());
}
@Override
public void onFailure(Throwable ex) {
log.info("fail Message : {}, Exception : {}", message, ex.getMessage());
}
});
}
}
9. Consumer 구현
-
@KafkaListener 어노테이션을 통해 Consumer 정보를 받을 수 있다. 원래는 KafkaConsumer 객체를 통해 topic을 구독하고 consumer.poll() 을통해 topic정보를 Polling 했지만, Spring Kafka에서는 @KafkaListener를 통해 이 역할을 수행한다.
-
@KafkaListener 의 상세정보를 보려면 Spring for Apache Kafka reference를 참조
기존 Kafka API호출
private final KafkaConsumer<String, String> consumer = new KafkaConsumer<String, String>(Consumer.createConfig(servers, groupId));
public ConsumerRecords<String, String> consume() {
this.consumer.subscribe(Collections.singleton(this.topic)); //2. topic publish
return consumer.poll(Duration.ofMillis(100)); //3. set timeout
}
Spring-Kafka API를 사용할 경우
[ConsumerService.java]
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.messaging.MessageHeaders;
import org.springframework.messaging.handler.annotation.Headers;
import org.springframework.messaging.handler.annotation.Payload;
import org.springframework.stereotype.Service;
@Slf4j
@Data
@Service
public class ConsumerService {
private String LastMsg;
@KafkaListener(topics = "${kafka.topic.source-topic}")
public void receiver(@Payload String message, @Headers MessageHeaders headers) {
headers.keySet().forEach(key -> log.info("key={}, value={}", key, headers.get(key)));
log.info("Received Message : {}", message);
setLastMsg(message);
}
}
10. Controller 구현
import com.springkafka.kafka.service.ConsumerService;
import com.springkafka.kafka.service.ProducerService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
@RestController
public class KafkaController {
private final ProducerService producerService;
private final ConsumerService consumerService;
@Autowired
public KafkaController(ProducerService producerService, ConsumerService consumerService) {
this.producerService = producerService;
this.consumerService = consumerService;
}
@GetMapping("/send")
public String sendProducer(@RequestParam(value = "msg") String message) {
producerService.send(message);
return "success";
}
@GetMapping("/receiver")
public String getMsg() {
return consumerService.getLastMsg();
}
}
11. build
- 오른쪽 창에서 gradle -> Tasks -> build 에서 clean을 한 후 build를 수행한다.

12. 실행 결과
- build 한 jar파일 실행(war를 만들었으면 톰캣에서 실행하면된다. Spring Boot Web을 선택했기에 내장 톰캣이 있어서 jar로 도 실행가능)
### .java -jar 파일명.jar
java -jar spring-kafka.jar
[spring-kafka.jar 실행]

- kafka-console-consumer.sh를 실행하여 topic 데이터가 잘들어오는지 확인해보자
$ ./bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic source-topic
- curl을 통해 테스트 수행
## /send?msg 요청
$ curl -X GET http://127.0.0.1:8096/send?msg=helloworld
## /receiver 요청
$ curl -X GET http://127.0.0.1:8096/receiver
[API 테스트 요청결과]

[전체 결과]

참조 사이트
https://docs.spring.io/spring-kafka/reference/html/
참조 사이트 https://docs.spring.io/spring-kafka/reference/html/
HTTP 요청과 Kafka 소비는 서로 다른 시점
원문의 /send?msg=는 데이터를 바꾸는 요청이므로 HTTP GET보다 POST가 적합하다. KafkaTemplate.send()는 비동기 전송이며, 호출 직후 문자열 success를 돌려주면 브로커의 저장 확인보다 먼저 성공을 알릴 수 있다. 반대로 /receiver가 ConsumerService.LastMsg 한 칸을 읽는 방식은 Kafka 토픽 조회가 아니다. 여러 요청과 여러 애플리케이션 인스턴스가 같은 필드를 덮어쓰고, 재시작하면 값도 사라진다.
flowchart LR H[POST 요청] --> P[KafkaTemplate 비동기 전송] P --> B[브로커 저장 확인] B --> R[HTTP 응답] B --> L["@KafkaListener 소비"] L --> D[업무 처리 또는 DB 저장]
Spring Boot가 spring-kafka 의존성을 발견하면 기본 ProducerFactory, KafkaTemplate, Listener 컨테이너 같은 구성을 자동화할 수 있다. 그래서 단순 문자열 전송·소비에 모든 Factory를 수동 작성할 필요는 없다. 다음은 핵심 흐름을 보여 주는 예이며 HTTP API에 인증·크기 제한·중복 요청 처리는 프로젝트에 맞게 더해야 한다.
spring:
kafka:
bootstrap-servers: localhost:9092
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.apache.kafka.common.serialization.StringSerializer
acks: all
consumer:
group-id: order-events-demo
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
@RestController
class EventController {
private final KafkaTemplate<String, String> kafka;
EventController(KafkaTemplate<String, String> kafka) {
this.kafka = kafka;
}
@PostMapping("/events")
CompletableFuture<ResponseEntity<String>> publish(@RequestBody String body) {
if (body == null || body.isBlank() || body.length() > 4096) {
return CompletableFuture.completedFuture(ResponseEntity.badRequest().body("invalid event"));
}
return kafka.send("source-topic", body)
.thenApply(result -> ResponseEntity.ok("stored at offset " + result.getRecordMetadata().offset()))
.exceptionally(error -> ResponseEntity.status(503).body("broker unavailable"));
}
}
@Component
class EventListener {
@KafkaListener(topics = "source-topic", groupId = "order-events-demo")
void handle(String value) {
// 실제 서비스라면 이벤트 ID로 중복 처리하고 결과를 영속 저장한다.
System.out.println(value);
}
}
위 코드는 해당 Java·Spring Boot 버전의 imports와 프로젝트 구성을 전제로 한 설명용 클래스다. HTTP 200은 브로커가 이 전송에 대해 성공을 알렸다는 뜻이지 소비자가 업무 처리를 끝냈다는 뜻이 아니다. 소비 완료를 사용자에게 알려야 한다면 이벤트 ID를 반환하고 별도 상태 조회 API나 알림 경로를 설계한다. 컨트롤러에서 받는 본문의 크기·형식·권한을 검증하고, 실제 운영에서는 임의의 사용자 입력을 원문 그대로 토픽에 게시하지 않는다.
| 상황 | 확인할 지점 |
|---|---|
| HTTP는 성공인데 listener 출력이 없음 | 토픽 이름, consumer 그룹과 오프셋, listener 로그 |
| 요청이 503으로 끝남 | 브로커 주소, 보안 설정, 전송 Future 오류 |
| 메시지가 중복 처리됨 | 재시도·오프셋 커밋과 이벤트 ID 중복 방지 |
| 서비스 인스턴스마다 마지막 메시지가 다름 | 원문의 메모리 필드 대신 DB/상태 저장소 사용 |
Kafka 클러스터의 주소가 컨테이너 밖/안에서 다르면 bootstrap-servers와 브로커의 advertised.listeners를 함께 확인한다. 소비자가 예외를 던질 때 무한 재시도만 하게 두면 같은 레코드가 파티션 처리를 막을 수 있다. 재시도 횟수, 지연, 실패 토픽과 알림 정책을 업무에 맞게 정한다. acks=all은 성공 확인의 조건을 강화하지만 단일 브로커 복제 계수 1에서는 물리적 장애에 대한 내구성을 보장하지 않는다.
참고: Spring Boot Kafka 지원, Spring Kafka 메시지 전송, @KafkaListener.
