kafka의 기본 개념을 기반으로 Spring에 적용하고 동작하는 방식을 확인하는 과정을 정리하고자 한다.
앞서 정리했던 개념을 다시한번 간단하게 정리하자면 아래와 같았다.
✅ Kafka 기본 개념
Producer, Consumer, Topic
- Producer
- Kafka로 메시지를 보내는 클라이언트
- 특정 Topic에 메시지를 기록
- Consumer
- Kafka에 저장된 메시지를 읽어오는 클라이언트
- 한 개 이상의 Topic을 구독(subscribe)
- Topic
- 메시지를 저장하는 논리적인 이름 공간
- 로그처럼 메시지가 순서대로 저장된다.
Kafka 프로젝트를 진행하기 전에 앞서 우리는 docker-compose 설정을 진행하고, KafkaConfig 설정을 통해 Spring과의 연동을 준비했다.
Kafka Cluster (docker-compose.yml)
version: '3.8'
services:
# --------------------------
# Broker 1 (Node ID: 1)
# --------------------------
kafka-1:
image: apache/kafka:3.7.0
container_name: kafka-1
ports:
- "9092:9092"
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-1:9093,2@kafka-2:9093,3@kafka-3:9093
# listeners
KAFKA_LISTENERS: CONTROLLER://:9093,INTERNAL://:29092,EXTERNAL://:9092
KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-1:29092,EXTERNAL://localhost:9092
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL
# replication configs
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2
KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
KAFKA_LOG_DIRS: /var/lib/kafka/data
volumes:
- kafka1_data:/var/lib/kafka/data
networks:
- kafka-net
# --------------------------
# Broker 2 (Node ID: 2)
# --------------------------
kafka-2:
image: apache/kafka:3.7.0
container_name: kafka-2
ports:
- "9093:9092"
environment:
KAFKA_NODE_ID: 2
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-1:9093,2@kafka-2:9093,3@kafka-3:9093
# listeners
KAFKA_LISTENERS: CONTROLLER://:9093,INTERNAL://:29092,EXTERNAL://:9092
KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-2:29092,EXTERNAL://localhost:9093
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL
# replication configs
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2
KAFKA_LOG_DIRS: /var/lib/kafka/data
volumes:
- kafka2_data:/var/lib/kafka/data
networks:
- kafka-net
# --------------------------
# Broker 3 (Node ID: 3)
# --------------------------
kafka-3:
image: apache/kafka:3.7.0
container_name: kafka-3
ports:
- "9094:9092"
environment:
KAFKA_NODE_ID: 3
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-1:9093,2@kafka-2:9093,3@kafka-3:9093
# listeners
KAFKA_LISTENERS: CONTROLLER://:9093,INTERNAL://:29092,EXTERNAL://:9092
KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-3:29092,EXTERNAL://localhost:9094
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL
# replication configs
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2
KAFKA_LOG_DIRS: /var/lib/kafka/data
volumes:
- kafka3_data:/var/lib/kafka/data
networks:
- kafka-net
# --------------------------
# Kafka UI
# --------------------------
kafka-ui:
image: provectuslabs/kafka-ui:latest
container_name: kafka-ui
ports:
- "8088:8080"
environment:
KAFKA_CLUSTERS_0_NAME: local-kraft-cluster
KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka-1:29092,kafka-2:29092,kafka-3:29092
KAFKA_CLUSTERS_0_METADATAPROVIDER: "KAFKA"
DYNAMIC_CONFIG_ENABLED: true
depends_on:
- kafka-1
- kafka-2
- kafka-3
networks:
- kafka-net
volumes:
kafka1_data:
kafka2_data:
kafka3_data:
networks:
kafka-net:
driver: bridge
✅ Kafka Config 설정
import java.util.HashMap;
import java.util.Map;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaProducerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.ProducerFactory;
@EnableKafka
@Configuration
public class KafkaConfig {
@Value("${spring.kafka.bootstrap-servers}")
private String bootstrapServers;
// 1) Producer 설정 (String, String)
@Bean
public ProducerFactory<String, String> stringProducerFactory() {
Map<String, Object> props = new HashMap<>();
// Kafka 서버 위치
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
// key, value를 문자열로 직렬화
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
return new DefaultKafkaProducerFactory<>(props);
}
// 2) Producer가 사용할 KafkaTemplate
@Bean
public KafkaTemplate<String, String> stringKafkaTemplate() {
return new KafkaTemplate<>(stringProducerFactory());
}
// 3) Consumer 설정 (String, String)
@Bean
public ConsumerFactory<String, String> stringConsumerFactory() {
Map<String, Object> props = new HashMap<>();
// Kafka 서버 위치
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
// 이 Consumer가 속할 그룹 아이디
props.put(ConsumerConfig.GROUP_ID_CONFIG, "simple-group");
// key, value를 문자열로 역직렬화
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
// 처음 시작할 때, 토픽 맨 앞(earliest)부터 읽도록 설정
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
return new DefaultKafkaConsumerFactory<>(props);
}
// 4) @KafkaListener가 사용할 ListenerContainerFactory
@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> stringKafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
// 3) 에서 만든 consumer 등록 !
factory.setConsumerFactory(stringConsumerFactory());
return factory;
}
}
위와 같이 설정한 내용을 글로 다시 풀어보자면
Producer 설정
- ProducerFactory<String, String>
- Producer가 사용할 설정 정보를 담는다
- Key, Value 모두 문자로 직렬화하기 위해 StringSerializer를 사용한다.
- KafkaTemplate<String, String>
- Producer 대신 실제로 send()를 호출하는 역할을 한다.
- Kafka와 Spring boot가 통신하는 객체
Consumer 설정
- ConsumerFactory<String, String>
- Consumer가 어떤 방식으로 메시지를 읽을지 설정
- Key, Value를 문자열로 역질렬화 하기 위해 StringDesrializer를 사용
- Group_ID_CONFIG = "simple-group"
→ 이 컨슈머가 속한 그룹 이름 - AUTO_OFFSET_RESET_CONFIG = "earliest"
→ 처음 실행 시, 토픽에 쌓여 있던 예전 메시지까지 모두 읽음
- ConCurrnetKafkaListenerContainerFactory<String, String>
- @KafkaListener가 내부적으로 사용할 컨테이너를 생성해주는 팩토리
- containerFactory 이름으로 @KafkaListener에 연결해서 사용
✅ 객체를 이벤트로 쓰기 : SimpleEvent 설계
단순한 문자열을 Kafka에 저장하고 소비하는 방식이 아닌 객체(DTO)를 Kafka를 통해 주고받을 수 있는 과정 또한 경험해보았다.
SimpleEvent 클래스
import java.time.LocalDateTime;
import lombok.AllArgsConstructor;
import lombok.Getter;
import lombok.NoArgsConstructor;
import lombok.ToString;
@Getter
@NoArgsConstructor
@AllArgsConstructor
@ToString
public class SimpleEvent {
// 실제 메시지 내용
private String message;
// 작업자
private String worker;
// 이벤트 생성 시각
private LocalDateTime createdAt;
}
이러한 객체를 받아주기 위한 Producer 설정도 진행해야 하는데, 방식은 크게 다르지 않았다.
// KafkaConfig.java 내부
import com.example.kafka.event.SimpleEvent;
import org.springframework.kafka.support.serializer.JsonSerializer;
// ...
// Producer 설정 (String, SimpleEvent) - JSON 직렬화
@Bean
public ProducerFactory<String, SimpleEvent> eventProducerFactory() {
Map<String, Object> props = new HashMap<>();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
return new DefaultKafkaProducerFactory<>(props);
}
@Bean
public KafkaTemplate<String, SimpleEvent> eventKafkaTemplate() {
return new KafkaTemplate<>(eventProducerFactory());
}
- ProducerFactory<String, SimpleEvent>
- key는 String, value는 SimpleEvent 타입으로 설정하고
- JsonSerializer
- SimpleEvent 객체를 JSON 문자열로 변환하여 Kafka에 전송
SimpleEvent용 Consumer 설정
// KafkaConfig.java 내부
import org.springframework.kafka.support.serializer.JsonDeserializer;
// ...
// Consumer 설정 (String, SimpleEvent) - JSON 역직렬화
@Bean
public ConsumerFactory<String, SimpleEvent> eventConsumerFactory() {
Map<String, Object> props = new HashMap<>();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
props.put(ConsumerConfig.GROUP_ID_CONFIG, "simple-event-group");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
// SimpleEvent 타입을 역직렬화하기 위한 JsonDeserializer 생성
JsonDeserializer<SimpleEvent> deserializer = new JsonDeserializer<>(SimpleEvent.class);
// 역직렬화를 허용할 패키지 지정 (보안 및 타입 안전성)
// !!! 개인마다 패키지 경로가 다를수있습니다. !!!
// SimpleEvent.class가 있는 위치를 넣어주시면 됩니다.
deserializer.addTrustedPackages("com.example.kafka.event");
return new DefaultKafkaConsumerFactory<>(
props,
new StringDeserializer(),
deserializer
);
}
@Bean
public ConcurrentKafkaListenerContainerFactory<String, SimpleEvent> eventKafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, SimpleEvent> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(eventConsumerFactory());
return factory;
}
- JsonDeserializer<SimpleEvent>
- Kafka에서 읽어온 JSON 데이터를 SimpleEvent 객체로 변환
- addTrustPackages("com.example.kafka.event")
- 패키지에 있는 클래스는 신뢰할 수 있다고 지정해주고
- 역직렬화 시 임의의 타입으로 변환하는 것을 방지하기 위한 보안 장치이다.
- eventKafkaListenerContainerFactory
- @KafkaListener( containerFactory = "eventKafkaListenerContainerFactory")에서 사용
Config 구현 이후 Producer , Listener는 아래와 같이 진행하였다.
SimpleEvent Producer
import lombok.RequiredArgsConstructor;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Service;
@Service
@RequiredArgsConstructor
public class SimpleEventProducer {
private static final String TOPIC = "simple-events";
private final KafkaTemplate<String, SimpleEvent> eventKafkaTemplate;
public void send(SimpleEvent event) {
eventKafkaTemplate.send(TOPIC, event);
}
}
- KafkaTemplate<String, SimpleEvent> 타입을 주입받는다
- send() 메서드에 SimpleEvent 객체를 그대로 넘긴다.
- JSON 변환은 JsonSerializer가 자동으로 처리한다
- 토픽 이름은 simple-events로 설정
SimpleEventListener
import lombok.extern.slf4j.Slf4j;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
@Slf4j
@Component
public class SimpleEventListener {
@KafkaListener(
topics = "simple-events",
groupId = "simple-event-group",
containerFactory = "eventKafkaListenerContainerFactory"
)
public void consume(SimpleEvent event) {
log.info("받은 이벤트: {}", event);
}
}
- topics = "simple-evenets"
- SimpleEventProducer 에서 사용한 토픽과 동일한 토픽을 구독
- groupId = "simple-event-group"
- KafkaConfig의 eventConsumerFactory에서 사용한 그룹 아이디와 동일해야 한다.
- containerFactory = "eventKafkaListenerContainerFactory"
- KafkaConfig에 정의한 이벤트용 ListenerContainerFactory를 사용
- 메서드 파라미터 타입이 SimpleEvent이다
- Consumer는 JSON 문자열을 직접 다루지 않는다
- JsonDeserializer<SimpleEvent>가 JSON → SimpleEvent 변환을 자동으로 수행한다.
'Spring > 백엔드 기초' 카테고리의 다른 글
| [개념정리] Kafka 클러스터 구조 이해하기 (0) | 2026.01.05 |
|---|---|
| [개념정리] Kafka - Broker와 Consumer Group (0) | 2026.01.05 |
| [개념정리] Kafka란 무엇인가? (0) | 2026.01.02 |
| [백엔드 기초] Redis 기초 개념 정리 (0) | 2025.12.08 |
| [개념정리] WebSocket 과 STOMP (0) | 2025.12.02 |