이 문서는 단일 가이드 문서다.

관련 문서 : 이벤트 기반 설계란 · @TransactionalEventListener · @Async와 비동기 처리

왜 Kafka가 필요한가

이벤트 기반 설계란에서 Spring Event로 서비스 간 결합도를 낮추는 방법을 배웠다. ApplicationEventPublisher로 이벤트를 발행하고, @EventListener로 수신한다. 잘 동작한다. 서버가 한 대일 때는.

서비스가 커지면 문제가 생긴다.

  • 알림 서비스를 별도 서버로 분리하고 싶다. 알림 연산량이 급증해서 메인 서버 성능에 영향을 주니까.
  • 그런데 Spring Event는 같은 JVM 안에서만 동작한다. 서버 A에서 발행한 이벤트를 서버 B에서 수신할 수 없다.

Kafka는 이 한계를 돌파한다. 서버 외부에 존재하는 메시지 브로커로, 서로 다른 서버(프로세스) 간에 이벤트를 주고받을 수 있게 해준다.

graph LR subgraph 기존 - Spring Event A1[메인 서비스] -->|publishEvent| A2[이벤트 리스너] Note1["같은 JVM 안에서만 동작"] end subgraph Kafka 도입 후 B1[메인 서비스] -->|send| K[Kafka 브로커] K -->|consume| B2[알림 서비스] Note2["서로 다른 서버에서 동작"] end style K fill:#fff3e0,stroke:#FF9800

우체국 비유로 설명하면, Spring Event는 같은 건물 안에서 쪽지를 전달하는 것이고, Kafka는 우체국을 거쳐 다른 건물로 편지를 보내는 것이다. 보내는 쪽은 받는 쪽이 어디 있는지 몰라도 된다. 우체국(Kafka)이 알아서 전달한다.

Kafka 핵심 개념

Kafka를 사용하기 전에 네 가지 핵심 개념을 알아야 한다.

Topic

메시지가 저장되는 논리적 채널이다. 이름을 붙여서 구분한다.

  • discodeit.MessageCreatedEvent — 메시지 생성 이벤트
  • discodeit.RoleUpdatedEvent — 권한 변경 이벤트

이메일 주소와 비슷하다. 보내는 쪽이 특정 토픽에 메시지를 보내면, 그 토픽을 구독하는 쪽이 받아간다.

Producer

토픽에 메시지를 보내는 쪽이다. Spring에서는 KafkaTemplate을 사용한다.

kafkaTemplate.send("discodeit.MessageCreatedEvent", payload);

Consumer

토픽의 메시지를 받는 쪽이다. Spring에서는 @KafkaListener를 사용한다.

@KafkaListener(topics = "discodeit.MessageCreatedEvent")
public void onMessageCreated(String message) { ... }

Consumer Group

같은 토픽을 여러 Consumer가 구독할 때, 같은 그룹에 속한 Consumer끼리는 메시지를 나눠 받고, 다른 그룹은 각각 전체 메시지를 받는다.

graph TD P[Producer] --> T[Topic] subgraph "Group A (알림 서비스)" C1[Consumer 1] C2[Consumer 2] end subgraph "Group B (로그 서비스)" C3[Consumer 3] end T --> C1 T --> C2 T --> C3 Note1["Group A: 메시지를 반씩 나눠 처리
(병렬 처리)"] Note2["Group B: 모든 메시지를 받음"]

같은 그룹의 Consumer가 여러 개면 메시지를 분산 처리할 수 있다. 알림이 밀릴 때 Consumer를 늘리면 처리량이 올라간다.

전체 흐름

Spring Event에서 Kafka로 전환할 때의 전체 아키텍처를 보자.

sequenceDiagram participant Service as 메인 서비스 participant Spring as Spring Event participant Producer as KafkaProduceListener participant Kafka as Kafka 브로커 participant Consumer as KafkaTopicListener participant NS as 알림 서비스 Service->>Spring: publishEvent(MessageCreatedEvent) Note over Spring: Spring Event는
그대로 유지 Spring->>Producer: @TransactionalEventListener Producer->>Kafka: kafkaTemplate.send("topic", payload) Note over Kafka: 메시지 보관
(디스크에 저장) Kafka->>Consumer: @KafkaListener Consumer->>NS: 알림 생성

핵심 포인트는 Spring Event를 없애는 게 아니다. Spring Event는 그대로 유지하고, 리스너에서 Kafka로 중계하는 구조다. 메인 서비스의 코드(publishEvent())는 바뀌지 않는다.

Docker로 Kafka 실행

Kafka는 외부 프로세스이므로 별도로 실행해야 한다. Docker Compose가 가장 간편하다.

# docker-compose-kafka.yaml
services:
  broker:
    image: apache/kafka:latest
    hostname: broker
    container_name: broker
    ports:
      - 9092:9092
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT,CONTROLLER:PLAINTEXT
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://broker:29092,PLAINTEXT_HOST://localhost:9092
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_PROCESS_ROLES: broker,controller
      KAFKA_NODE_ID: 1
      KAFKA_CONTROLLER_QUORUM_VOTERS: 1@broker:29093
      KAFKA_LISTENERS: PLAINTEXT://broker:29092,CONTROLLER://broker:29093,PLAINTEXT_HOST://0.0.0.0:9092
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
      KAFKA_LOG_DIRS: /tmp/kraft-combined-logs
      CLUSTER_ID: MkU3OEVBNTcwNTJENDM2Qk
docker compose -f docker-compose-kafka.yaml up -d

이 설정은 KRaft 모드(Zookeeper 없이 동작)로 Kafka 단일 브로커를 실행한다. 개발 환경에서는 이 정도면 충분하다.

설정이 복잡해 보이는 이유

대부분 Kafka 내부 통신 설정이다. 중요한 건 ports: 9092:9092로, 애플리케이션에서 localhost:9092로 접속할 수 있게 포트를 매핑하는 것이다.

Spring Kafka 설정

의존성

implementation 'org.springframework.kafka:spring-kafka'

application.yml

spring:
  kafka:
    bootstrap-servers: localhost:9092
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
    consumer:
      group-id: discodeit-group
      auto-offset-reset: earliest
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
설정의미
bootstrap-serversKafka 브로커 주소. Docker에서 9092로 매핑했으므로 localhost:9092
key/value-serializer메시지를 보낼 때 직렬화 형식. 문자열(JSON)로 보낸다
group-idConsumer Group 이름. 같은 그룹의 Consumer끼리 메시지를 나눠 받는다
auto-offset-resetConsumer가 처음 구독할 때 어디부터 읽을지. earliest면 처음부터
key-serializer와 value-serializer

Producer는 Serializer, Consumer는 Deserializer를 사용한다. 짝이 맞아야 한다. 둘 다 String으로 맞추고, 이벤트 객체를 JSON 문자열로 변환해서 보내는 것이 가장 간단하다.

Producer 구현 — Spring Event → Kafka 중계

기존 @TransactionalEventListener를 직접 수정하는 대신, Kafka 전용 리스너를 새로 만들고 기존 리스너를 비활성화하는 구조를 쓴다.

@Slf4j
@RequiredArgsConstructor
@Component
public class KafkaProduceEventListener {

    private final KafkaTemplate<String, String> kafkaTemplate;
    private final ObjectMapper objectMapper;

    @Async("eventTaskExecutor")
    @TransactionalEventListener
    public void on(MessageCreatedEvent event) {
        try {
            String payload = objectMapper.writeValueAsString(event);
            kafkaTemplate.send("discodeit.MessageCreatedEvent", payload);
            log.info("Kafka 이벤트 발행 완료 - topic: discodeit.MessageCreatedEvent");
        } catch (JsonProcessingException e) {
            log.error("이벤트 직렬화 실패", e);
        }
    }

    @Async("eventTaskExecutor")
    @TransactionalEventListener
    public void on(RoleUpdatedEvent event) {
        try {
            String payload = objectMapper.writeValueAsString(event);
            kafkaTemplate.send("discodeit.RoleUpdatedEvent", payload);
            log.info("Kafka 이벤트 발행 완료 - topic: discodeit.RoleUpdatedEvent");
        } catch (JsonProcessingException e) {
            log.error("이벤트 직렬화 실패", e);
        }
    }
}

이 클래스가 하는 일은 단순하다. Spring Event를 받아서 → JSON으로 변환 → Kafka 토픽으로 전송한다.

  • KafkaTemplate<String, String> : Spring Kafka가 제공하는 메시지 전송 도구. send(topic, value)로 토픽에 메시지를 보낸다.
  • ObjectMapper : 이벤트 객체를 JSON 문자열로 변환한다. Consumer 쪽에서 같은 형식으로 역직렬화한다.
  • @Async + @TransactionalEventListener : 기존 패턴 그대로. 트랜잭션 커밋 후 비동기로 Kafka에 전송한다.
sequenceDiagram participant Service as MessageService participant Spring as Spring Event participant Listener as KafkaProduceListener participant Mapper as ObjectMapper participant Kafka as KafkaTemplate Service->>Spring: publishEvent(MessageCreatedEvent) Note over Service: publishEvent() 코드는
변경 없음! Spring->>Listener: on(MessageCreatedEvent) Listener->>Mapper: writeValueAsString(event) Mapper-->>Listener: JSON 문자열 Listener->>Kafka: send("discodeit.MessageCreatedEvent", json) Kafka-->>Listener: 전송 완료

Consumer 구현 — Kafka → 알림 생성

Kafka 토픽을 구독해서 알림을 생성하는 Consumer를 만든다. 실제로는 별도 서비스(별도 서버)로 분리하지만, 학습 단계에서는 같은 프로젝트에 둔다.

@Slf4j
@RequiredArgsConstructor
@Component
public class NotificationTopicListener {

    private final NotificationService notificationService;
    private final UserRepository userRepository;
    private final ObjectMapper objectMapper;

    @KafkaListener(topics = "discodeit.MessageCreatedEvent")
    public void onMessageCreated(String kafkaEvent) {
        try {
            MessageCreatedEvent event = objectMapper.readValue(
                kafkaEvent, MessageCreatedEvent.class);
            log.info("Kafka 이벤트 수신 - MessageCreatedEvent");
            // 알림 생성 로직
        } catch (JsonProcessingException e) {
            log.error("이벤트 역직렬화 실패", e);
        }
    }

    @KafkaListener(topics = "discodeit.RoleUpdatedEvent")
    public void onRoleUpdated(String kafkaEvent) {
        try {
            RoleUpdatedEvent event = objectMapper.readValue(
                kafkaEvent, RoleUpdatedEvent.class);
            log.info("Kafka 이벤트 수신 - RoleUpdatedEvent");
            // 알림 생성 로직
        } catch (JsonProcessingException e) {
            log.error("이벤트 역직렬화 실패", e);
        }
    }
}
  • @KafkaListener(topics = "...") : 지정한 토픽을 구독한다. 새 메시지가 도착하면 이 메서드가 호출된다.
  • objectMapper.readValue() : JSON 문자열을 다시 이벤트 객체로 변환한다.
  • 기존 NotificationRequiredEventListener의 알림 생성 로직을 여기로 옮기면 된다.

Kafka 콘솔로 이벤트 확인

메시지가 Kafka에 잘 전달되는지 콘솔로 확인할 수 있다.

# broker 컨테이너에 접속
docker exec -it -w /opt/kafka/bin broker sh

# 토픽 목록 확인
./kafka-topics.sh --list --bootstrap-server broker:29092

# 특정 토픽의 메시지 구독 (실시간 확인)
./kafka-console-consumer.sh --topic discodeit.MessageCreatedEvent \
  --from-beginning --bootstrap-server broker:29092

메시지를 생성하면 콘솔에 JSON이 출력된다.

{"channelId":"abc-123","authorName":"alice","content":"안녕하세요"}

Spring Event → Kafka 전환 정리

전환의 핵심을 정리하면 이렇다.

구분Spring Event (기존)Kafka (심화)
범위같은 JVM 내부서버 간 통신 가능
발행eventPublisher.publishEvent()kafkaTemplate.send()
수신@EventListener / @TransactionalEventListener@KafkaListener
메시지 형식Java 객체 그대로JSON 문자열 (직렬화 필요)
메시지 보존JVM 메모리 (휘발성)디스크 (영속성)
publishEvent() 코드그대로 유지변경 없음

메인 서비스의 publishEvent() 코드는 바꿀 필요가 없다. 리스너 쪽에서 Spring Event를 받아 Kafka로 중계하는 구조이기 때문이다. 이것이 이벤트 기반 설계의 장점이다. 발행자는 수신자가 누구인지, 어떻게 처리하는지 모른다.

실제 분산 환경에서는

Producer(메인 서비스)와 Consumer(알림 서비스)가 완전히 별도 프로젝트, 별도 서버다. Consumer 프로젝트에는 @KafkaListener만 있고, 이벤트 클래스는 공유 라이브러리로 분리하거나 JSON 스키마를 맞춘다. 이번 미션에서는 같은 프로젝트에 두지만, "나중에 분리할 수 있다"는 점이 중요하다.