실시간 채팅, 주문 시스템, 알림 서비스, 로그 수집과 같은 기능들은 서비스 간에 데이터를 주고 받는 것이 거의 필수이다. 이러한 작업들을 하나의 서버에서 직접 처리하면 응답 속도가 느려지고, 결합도가 높아진다. 이 문제를 해결하기 위해 있는 것이 kafka(Apache Kafka)이다.
여러 언어에서 사용할 수 있지만 필자는 java로 예시를 사용하겠다.
Kafka는 대용량 실시간 데이터를 안정적으로 전달하기 위한 분산 메세지 브로커(Message Broker)이다.
데이터를 보내는 시스템과 데이터를 처리하는 시스템 사이에서 데이터를 안전하게 전달해주는 플랫폼
구조
Producer ⇒ Kafka(Broker) ⇒ Consumer
1. 높은 처리량
Kafka는 초당 수십만~수백만 건 이상의 message를 처리한다.
2. 낮은 지연 시간
message가 매우 빠르게 전달된다.
3. 데이터 보존
일반적인 message queue는 소비하면 message가 사라지는 경우가 많다.
Kafka는 message를 일정 기간 저장한다.
4. 확장성
데이터를 분산 저장하기 때문에 서버를 추가하면서 성능을 확장할 수 있다.

사진 출처 https://dbrang.tistory.com/1626
Producer
message를 Kafka에 전송하는 역할을 한다.
kafkaTemplate.send("orders", "Order Created")
// topic = orders, message = Order Created
Consumer
Kafka에서 message를 읽어 처리하는 역할을 한다. Topic에 저장된 메세지를 읽는다.
@KafkaListener(topics = ["orders"])
fun consume(message: String) {
println(message)
}
Broker
kafka 서버 하나를 의미한다.
Kafka는 하나 이상의 Broker로 구성된다.
Broker는 메세지 저장, Topic 관리, Consumer 요청 처리를 담당한다.
여러 개의 Broker가 모여 Kafka Cluster를 구성한다.
Kafka Cluster
Broker들의 모임이다. Kafka는 확장성과 고가용성을 위하여 broker를 cluster로 구성한다.
cluster : '집단'을 의미하며 여러 대의 컴퓨터나 서버를 네트워크로 연결하여 하나의 고성능 시스템처럼 작동하게 만드는 기술
데이터 분산 저장, 장애 대응, 성능 확장 등을 위해 이처럼 구성한다.
Topic
message를 저장하는 논리적인 공간이다.
Producer는 Topic에 message를 저장하고, Consumer는 Topic에서 message를 읽는다.
Partition
Topic은 내부적으로 여러 개의 Partition으로 나뉜다.
Partition을 사용하는 이유는 병렬 처리를 하기 위해서이다.
Consumer도 Partition 단위로 읽는다.
Offset
Partition 내부에서 message의 위치를 나타내는 번호이다.
Consumer는 Offset을 기준으로 읽는다.
Pub/Sub
Publish/Subscribe 모델이다.
Publish ⇒ Producer, Subscribe ⇒ Consumer
예를 들어, Producer가 Order Created를 Publish하면, 여러 Consumer가 Subscribe할 수 있다.

각 Consumer는 동일한 message를 받을 수 있다.
Consumer Group
여러 Consumer를 하나의 그룹으로 묶는 것이다. 각각의 Group에서는 하나의 Partition을 하나의 Consumer만 읽는다. 따라서 중복 처리가 발생하지 않는다.
zooKeeper
ZooKeeper는 예전 Kafka에서 Cluster를 관리하던 시스템이다.
Broker 관리, Leader 선출, 메타데이터 관리, Cluster 상태 관리
Kafka 3.x 이후부터는 KRaft가 도입되었고, 4.x 버전부터는 ZooKeeper없이 Kafka 자체적으로 Cluster를 관리하는 KRaft 모드가 기본으로 되었다.
Producer가 메세지 전송 ⇒ Kafka 저장 ⇒ Consumer가 읽음 ⇒ Offset 저장
1. Producer
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Service;
@Service
public class OrderProducer {
private final KafkaTemplate<String, String> kafkaTemplate;
public OrderProducer(KafkaTemplate<String, String> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
public void sendOrder() {
kafkaTemplate.send("orders", "Order Created");
}
}
// Topic : orders
// Message : Order Created
2. 실행 코드
import org.springframework.boot.CommandLineRunner;
import org.springframework.stereotype.Component;
@Component
public class Runner implements CommandLineRunner {
private final OrderProducer producer;
public Runner(OrderProducer producer) {
this.producer = producer;
}
@Override
public void run(String... args) {
producer.sendOrder();
}
}
// 애플리케이션 실행 ⇒ producer.sendOrder() ⇒ Kafka Broker로 message 전송
3. Broker 내부
Topic: orders
--------------------------------------
Offset 0 Order Created
4. Consumer
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
@Component
public class OrderConsumer {
@KafkaListener(topics = "orders", groupId = "order-group")
public void consume(String message) {
System.out.println("받은 메시지: " + message);
}
}
Consumer는 orders Topic을 구독한다.
message가 들어오면 자동으로 consume() 메서드가 실행된다.
실행 결과
실행하면
받은 메세지: Order Created
가 출력된다.
1. Producer
Producer는 Broker에 요청을 보낸다.
kafkaTemplate.send("orders", "Order Created");
2. Broker
Broker는 Topic에 저장한다.
Offset0 : Order Created
3. Consumer
Consumer는 Offset 0을 읽는다.
Offset0 ⇒ consume()
4. Offsset 커밋
Consumer는 Offset0 처리 완료를 Kafka에 기록한다.
다음 메세지는 Offset1부터 읽는다.
메세지 저장 방법
Kafka의 Topic은 내부적으로 Partition으로 나뉜다.
메세지는 순서대로 저장된다.
여기서 Offset은 메세지의 번호이다.
Consumer는 Offset을 기준으로 어디까지 읽었는지 관리한다.