Apache Kafka와 이벤트 기반 아키텍처
Kafka 토픽과 소비자 그룹의 역할을 이해하고 사용자 생성 이벤트를 발행해 알림 서비스에서 비동기로 처리합니다.
지난 절에서는 NestJS와 gRPC를 활용하여 고성능 서비스 간 동기 통신을 구현하는 방법을 살펴보았습니다.
이제 7장의 마지막 절로, 마이크로서비스 아키텍처에서 비동기 통신과 이벤트 기반 아키텍처(EDA)의 핵심 요소인 Apache Kafka를 NestJS에 적용하는 방법에 대해 알아보겠습니다.
동기 통신(Synchronous Communication)은 요청-응답 패턴에 적합하지만, 서비스 간 강한 결합을 유발하고 한 서비스의 장애가 다른 서비스로 전파될 위험이 있습니다.
또한 대규모 이벤트를 처리하거나 실시간 데이터 파이프라인을 구축할 때는 비효율적입니다.
Apache Kafka는 이러한 문제를 해결하고 확장성, 내결함성, 높은 처리량을 제공하는 분산 스트리밍 플랫폼입니다.
동기 호출은 상대 서비스의 응답을 기다리지만, Kafka 기반 흐름은 이벤트를 토픽에 남기고 소비자가 자신의 속도로 처리합니다.
- 1요청한 서비스가 응답을 기다린다
Synchronous Notifications가 느리거나 장애가 나면 Users 흐름까지 영향을 받을 수 있습니다.
- 2이벤트를 남기고 흐름을 분리한다
Event-driven Users는 `user_created`를 발행하고, Notifications는 나중에 읽어 처리합니다.
- 3Users Service
Producer 사용자 생성 후 이벤트를 발행합니다.
- 4user_created
Topic log 이벤트가 Kafka 로그에 저장되어 재처리 가능한 기록이 됩니다.
- 5Notifications Service
Consumer 자신의 consumer group에서 이벤트를 읽어 환영 메일을 보냅니다.
이벤트 기반 아키텍처 (EDA)와 Kafka
이벤트 기반 아키텍처(EDA, Event-Driven Architecture)는 시스템의 구성 요소들이 직접적으로 호출하는 대신, 이벤트(Event)를 발행하고 구독하는 방식으로 상호작용하는 아키텍처 스타일입니다.
각 서비스는 자신이 관심 있는 이벤트를 구독하고, 해당 이벤트가 발생하면 필요한 작업을 수행합니다.
EDA의 주요 특징- 느슨한 결합(Loose Coupling): 서비스들이 서로를 직접 호출하지 않으므로, 한 서비스의 변경이 다른 서비스에 미치는 영향이 최소화됩니다.
- 확장성: 이벤트 생산자와 소비자 모두 독립적으로 확장할 수 있습니다.
- 복원력(Resilience): 한 서비스가 다운되더라도 이벤트는 메시지 큐에 저장되어 나중에 처리될 수 있으므로, 시스템 전체의 가용성이 높아집니다.
- 비동기 통신: 응답을 즉시 기다리지 않으므로, 장시간 실행되는 작업을 처리하거나 백그라운드 작업을 실행하는 데 적합합니다.
- 실시간 처리: 데이터 변경 이벤트를 즉시 처리하여 실시간 시스템 구축에 용이합니다.
Apache Kafka는 이러한 EDA를 구현하는 데 가장 널리 사용되는 플랫폼 중 하나입니다.
Kafka는 분산된 커밋 로그(Commit Log) 또는 분산 메시지 브로커로 볼 수 있으며, 다음과 같은 특징을 가집니다.
- 높은 처리량: 초당 수백만 건의 이벤트를 처리할 수 있습니다.
- 내결함성: 여러 서버에 데이터를 복제하여 서버 장애 시에도 데이터 손실 없이 안정적인 운영이 가능합니다.
- 확장성: 수평적 확장이 용이하여 데이터 양이 증가해도 유연하게 대응할 수 있습니다.
- 영속성: 발행된 이벤트는 디스크에 저장되어 설정된 기간 동안 유지되므로, 소비자가 나중에 이벤트를 다시 읽을 수 있습니다.
- 다중 소비자: 여러 소비자가 동일한 이벤트를 독립적으로 소비할 수 있습니다.
- Producer(생산자): Kafka 토픽으로 메시지(이벤트)를 발행하는 애플리케이션.
- Consumer(소비자): Kafka 토픽에서 메시지를 구독하고 처리하는 애플리케이션.
- Topic(토픽): 메시지(이벤트)가 발행되고 소비되는 논리적인 카테고리.
- Partition(파티션): 토픽을 물리적으로 나누는 단위. 파티션 덕분에 Kafka는 높은 처리량을 가집니다. 각 파티션의 메시지는 순서가 보장됩니다.
- Broker(브로커): Kafka 서버. 여러 브로커가 클러스터를 구성하여 내결함성과 확장성을 제공합니다.
- Zookeeper(주키퍼): Kafka 클러스터의 메타데이터(토픽, 파티션 정보 등)를 관리하는 분산 코디네이션 서비스 (최신 Kafka 버전에서는 주키퍼 없이도 동작 가능).
Kafka 용어는 따로 외우기보다 이벤트가 어디에 저장되고 누가 어느 조각을 읽는지 위치로 보면 빠르게 정리됩니다.
- 이벤트 발행자
Producer Users Service가 `user_created` 이벤트를 topic에 씁니다.
- Topic: user_created
partition 0 0 1 2 partition 1 0 1 2 각 partition 안에서는 순서가 유지되고 offset으로 읽은 위치를 추적합니다.
- 작업을 나눠 읽는 소비자
Consumer group 같은 group의 consumer들은 partition을 나눠 맡아 중복 처리를 줄입니다.
- Broker
토픽 데이터를 저장하는 Kafka 서버입니다.
- Replication
파티션 복제본을 두어 broker 장애에도 데이터를 지킵니다.
- Retention
읽었는지와 별개로 이벤트를 정해진 기간 보관합니다.
Kafka 이벤트는 단순한 메시지가 아니라 재처리, 순서, 소비자 확장까지 포함한 운영 계약입니다.
토픽을 만들기 전에 아래 기준으로 이벤트 경계를 먼저 잡아두면 구현과 검증이 쉬워집니다.
이벤트는 한 번 흘려보내는 알림이 아니라, 나중에 다시 읽고 여러 소비자가 나눠 처리할 운영 기록입니다.
- 업무 사건 단위
Topic 업무 사건 단위 `user_created`처럼 발생한 사실을 이름으로 둡니다.
- 순서를 묶는 기준
Key 순서를 묶는 기준 같은 사용자 이벤트가 같은 partition에 가도록 key를 정합니다.
- 소비자가 필요한 최소 정보
Payload 소비자가 필요한 최소 정보 변경 사실과 식별자를 담고, 큰 조회 데이터는 줄입니다.
- 처리 책임 단위
Consumer group 처리 책임 단위 메일, 분석, 감사 로그가 서로 독립적으로 읽을 수 있습니다.
| 설계 항목 | 정할 것 | 놓치면 생기는 문제 |
|---|---|---|
| 파티션 키 | 같은 순서가 필요한 이벤트 묶음 | 사용자별 이벤트 순서가 뒤섞임 |
| 보관 기간 | 재처리 가능한 retention 기간 | 장애 후 복구할 이벤트가 사라짐 |
| 중복 처리 | 같은 이벤트가 두 번 와도 안전한 처리 | 메일 중복 발송이나 정산 중복 반영 |
NestJS에서 Kafka 설정 및 구현
NestJS는 @nestjs/microservices 패키지를 통해 Kafka 전송 계층을 지원합니다.
시나리오: Users 서비스에서 새로운 사용자가 생성되면 이벤트를 Kafka에 발행하고, Notifications 서비스에서 이 이벤트를 구독하여 사용자에게 환영 이메일을 보내는 시나리오를 구축해 보겠습니다.
사전 준비: 로컬 환경에서 Kafka와 Zookeeper를 실행해야 합니다.
Docker를 사용하는 것이 가장 편리합니다.
# Docker Compose 파일 (docker-compose.yml) 예시
version: '3.8'
services:
zookeeper:
image: confluentinc/cp-zookeeper:7.5.0
hostname: zookeeper
container_name: zookeeper
ports:
- "2181:2181"
environment:
ZOOKEEPER_CLIENT_PORT: 2181
ZOOKEEPER_TICK_TIME: 2000
kafka:
image: confluentinc/cp-kafka:7.5.0
hostname: kafka
container_name: kafka
ports:
- "9092:9092"
- "9093:9093" # 내부 통신용 포트
environment:
KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9093,PLAINTEXT_HOST://localhost:9092
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
depends_on:
- zookeeper
# 실행: docker-compose up -dDocker Compose 파일로 Kafka를 실행한 후 다음 단계를 진행합니다.
Kafka 사용자 이벤트 서비스 구축
새로운 사용자가 생성될 때 Kafka에 이벤트를 발행하는 서비스입니다.
단계 1: 필요한 패키지 설치users-service 프로젝트에서 설치합니다.
npm install @nestjs/microservices @nestjs/platform-express
npm install --save-dev @types/express # 웹 서버 역할을 겸할 경우kafka-node(레거시): 이전에는kafka-node를 사용했지만, NestJS 최신 버전에서는@nestjs/microservices가 내부적으로kafkajs를 기반으로 하므로 별도로 설치할 필요가 없습니다.
main.ts 파일 수정 (HTTP 서버 유지)
사용자 서비스는 클라이언트의 HTTP 요청을 받아 사용자 생성 후 Kafka 이벤트를 발행합니다.
import { NestFactory } from '@nestjs/core';
import { AppModule } from './app.module';
async function bootstrap() {
const app = await NestFactory.create(AppModule);
// NestJS 마이크로서비스 챕터의 컨벤션을 위해 포트 3000 사용
await app.listen(3000);
console.log('Users Service (HTTP Server) is listening on port 3000');
}
bootstrap();UsersModule에 Kafka 클라이언트 등록
ClientsModule을 사용하여 Kafka 클라이언트를 등록하고, 이를 UsersController에 주입하여 이벤트를 발행합니다.
import { Module } from '@nestjs/common';
import { ClientsModule, Transport } from '@nestjs/microservices';
import { UsersController } from './users.controller';
import { UsersService } from './users.service';
@Module({
imports: [
ClientsModule.register([
{
name: 'KAFKA_SERVICE', // Kafka 클라이언트 토큰
transport: Transport.KAFKA, // Kafka 전송 방식 사용
options: {
client: {
clientId: 'user-producer', // 클라이언트 ID
brokers: ['localhost:9092'], // Kafka 브로커 주소
},
producer: {
allowAutoTopicCreation: true, // 토픽이 없으면 자동으로 생성
},
},
},
]),
],
controllers: [UsersController],
providers: [UsersService],
})
export class UsersModule {}name: 'KAFKA_SERVICE': 이 클라이언트 프록시를 의존성 주입을 통해 참조할 토큰입니다.transport: Transport.KAFKA: Kafka 전송 방식을 사용합니다.options.client: Kafka 클라이언트 설정 (고유한clientId, 브로커 목록brokers).options.producer: Kafka 생산자 설정 (예:allowAutoTopicCreation으로 토픽 자동 생성 활성화).
사용자 생성 시 Kafka emit 메서드를 통해 이벤트를 발행합니다.
import { Controller, Post, Body, Get, Param, Inject } from '@nestjs/common';
import { ClientKafka } from '@nestjs/microservices'; // ClientKafka 임포트 (ClientProxy의 하위 타입)
import { UsersService } from './users.service';
import { OnModuleInit } from '@nestjs/common'; // OnModuleInit 임포트
import { lastValueFrom } from 'rxjs'; // lastValueFrom 임포트
interface User {
id: number;
name: string;
email: string;
}
@Controller('users') // HTTP 엔드포인트
export class UsersController implements OnModuleInit {
constructor(
private readonly usersService: UsersService,
@Inject('KAFKA_SERVICE') private readonly client: ClientKafka, // Kafka 클라이언트 주입
) {}
async onModuleInit() {
// 이 클라이언트가 메시지를 발행할 모든 토픽을 구독해야 합니다.
// NestJS는 이 토픽을 미리 연결하여 메시지 전송 실패를 방지합니다.
this.client.subscribeToResponseOf('user_created');
await this.client.connect(); // 클라이언트 연결
}
@Post()
async createUser(@Body() userDto: { name: string; email: string }): Promise<User> {
const newUser = this.usersService.create(userDto);
console.log(`Users Service: New user created: ${JSON.stringify(newUser)}`);
// 'user_created' 토픽으로 이벤트 발행
// emit()은 응답을 기다리지 않는 이벤트 기반 통신에 사용합니다.
this.client.emit('user_created', newUser); // 메시지 키는 생략하고 값만 전송
return newUser;
}
// 기존 사용자 조회 엔드포인트 (HTTP)
@Get(':id')
getUserById(@Param('id') id: string): User {
const userId = parseInt(id, 10);
return this.usersService.findOne(userId);
}
}ClientKafka: Kafka 전송 방식을 위한 특화된 클라이언트입니다.onModuleInit():client.connect()를 호출하여 Kafka 브로커와 연결을 설정합니다.subscribeToResponseOf()는 클라이언트가 서버로부터 응답을 받을 때 사용하지만,emit을 사용할 때는 필수적이지 않습니다. 하지만 안전하게 모든 토픽을 미리 연결하고 싶을 때 유용합니다.this.client.emit('topic_name', message_payload):topic_name으로message_payload를 발행합니다.emit은 비동기적이며 응답을 기다리지 않습니다.
UsersService 구현 (동일)
users-service/src/users/users.service.ts 파일은 이전과 동일하게 유지됩니다.
AppModule에 UsersModule 임포트 (동일)
users-service/src/app.module.ts는 변경 없이 UsersModule을 임포트합니다.
알림 서비스 구축
Users Service는 사용자를 저장한 뒤 즉시 HTTP 응답을 돌려주고, 알림 처리는 Kafka 이벤트를 구독한 서비스가 따로 수행합니다.
- UsersController
Producer client.emit('user_created', user) 사용자 생성 사실을 토픽으로 발행합니다.
- topic: user_created
Kafka brokers: ['localhost:9092'] 브로커가 이벤트를 저장하고 consumer group에 전달합니다.
- NotificationsController
Consumer @EventPattern('user_created') 이벤트를 받아 환영 이메일을 보냅니다.
- HTTP 책임
사용자 생성 성공 여부는 Users Service가 바로 응답합니다.
- 이벤트 책임
알림 전송은 Kafka 이벤트 처리 성공 여부로 따로 관측합니다.
- 연결 책임
`clientId`, `brokers`, `groupId`가 서비스별로 구분되어야 합니다.
user_created 이벤트를 구독하여 처리하는 서비스입니다.
nest new notifications-service --skip-install
cd notifications-service
npm install @nestjs/microservices
npm installmain.ts 파일 수정 (Kafka 마이크로서비스 서버 설정)
import { NestFactory } from '@nestjs/core';
import { AppModule } from './app.module';
import { MicroserviceOptions, Transport } from '@nestjs/microservices';
import { join } from 'path';
async function bootstrap() {
const app = await NestFactory.createMicroservice<MicroserviceOptions>(AppModule, {
transport: Transport.KAFKA, // Kafka 전송 방식 사용
options: {
client: {
clientId: 'notification-consumer', // 클라이언트 ID
brokers: ['localhost:9092'], // Kafka 브로커 주소
},
consumer: {
groupId: 'notification-group', // 컨슈머 그룹 ID (필수)
},
},
});
await app.listen();
console.log('Notifications Microservice (Kafka Consumer) is listening.');
}
bootstrap();transport: Transport.KAFKA: Kafka 전송 방식을 사용합니다.options.client: Kafka 클라이언트 설정.options.consumer.groupId: 필수 항목입니다. Kafka 컨슈머는 반드시 컨슈머 그룹에 속해야 합니다. 동일한 그룹 ID를 가진 컨슈머들은 토픽의 파티션을 나누어 소비합니다.
@EventPattern() 데코레이터를 사용하여 특정 Kafka 토픽에서 발생하는 이벤트를 처리합니다.
import { Controller } from '@nestjs/common';
import { EventPattern, Payload } from '@nestjs/microservices'; // EventPattern, Payload 임포트
import { NotificationsService } from './notifications.service';
interface UserCreatedEvent { // user_created 이벤트의 페이로드 구조 정의
id: number;
name: string;
email: string;
}
@Controller()
export class NotificationsController {
constructor(private readonly notificationsService: NotificationsService) {}
// 'user_created' 토픽의 이벤트를 처리하는 핸들러
@EventPattern('user_created')
async handleUserCreated(@Payload() data: UserCreatedEvent) {
console.log(`Notifications Service: Received 'user_created' event for user: ${JSON.stringify(data)}`);
// 여기에서 실제 알림 로직 (예: 이메일 전송)을 구현합니다.
await this.notificationsService.sendWelcomeEmail(data.email, data.name);
}
}@EventPattern('topic_name'): 특정 Kafka 토픽에서 발행된 이벤트를 구독하여 처리합니다.send()와 달리 응답을 반환할 필요가 없습니다.@Payload(): 이벤트 메시지의 페이로드(값)를 추출하여 메서드 인자로 주입합니다.
import { Injectable } from '@nestjs/common';
@Injectable()
export class NotificationsService {
async sendWelcomeEmail(email: string, name: string): Promise<void> {
console.log(`Sending welcome email to ${name} (${email})...`);
// 실제 이메일 전송 로직 (예: Nodemailer, SendGrid 등 사용)
await new Promise(resolve => setTimeout(resolve, 1000)); // 1초 지연 시뮬레이션
console.log(`Welcome email sent to ${name}!`);
}
}NotificationsModule 및 AppModule 구성
import { Module } from '@nestjs/common';
import { NotificationsController } from './notifications.controller';
import { NotificationsService } from './notifications.service';
@Module({
controllers: [NotificationsController],
providers: [NotificationsService],
})
export class NotificationsModule {}import { Module } from '@nestjs/common';
import { NotificationsModule } from './notifications/notifications.module';
@Module({
imports: [NotificationsModule],
controllers: [],
providers: [],
})
export class AppModule {}컨테이너 환경 실행 검증
- Docker Compose 파일이 있는 디렉토리에서:
docker-compose up -d - 모든 컨테이너가 정상적으로 실행되었는지 확인:
docker-compose ps
cd notifications-servicenpm run start:dev- 콘솔에
Notifications Microservice (Kafka Consumer) is listening.메시지 확인.
cd users-servicenpm run start:dev- 콘솔에
Users Service (HTTP Server) is listening on port 3000메시지 확인.
POST http://localhost:3000/users- Headers:
Content-Type: application/json - Body
{ "name": "David", "email": "david@example.com" } - Users Service 콘솔:
Users Service: New user created: {"id":3,"name":"David","email":"david@example.com"}메시지가 즉시 출력됩니다. - Notifications Service 콘솔: 잠시 후 (Kafka 브로커를 거쳐 메시지가 전달된 뒤),
Notifications Service: Received 'user_created' event for user: {"id":3,"name":"David","email":"david@example.com"}와Sending welcome email to David (david@example.com)...,Welcome email sent to David!메시지가 출력되는 것을 확인할 수 있습니다.
이 과정을 통해 Users 서비스에서 Notifications 서비스로 직접 호출 없이 Kafka를 통해 비동기적으로 이벤트가 전달되고 처리되는 것을 확인했습니다.
Kafka 검증은 API 응답보다 이벤트가 브로커를 지나 소비자까지 도착했는지가 핵심입니다.
아래 다이어그램은 컨테이너, 프로듀서, 토픽, 컨슈머 로그를 순서대로 확인하는 체크 포인트를 보여줍니다.
사용자 생성 API가 성공해도 알림 이벤트가 소비자까지 도착했는지는 별도로 확인해야 합니다.
- Broker 실행
Container docker-compose ps 정상 이면 zookeeper와 kafka가 떠 있습니다.
- Users log
Producer New user created HTTP 요청이 이벤트 발행 지점까지 도달했습니다.
- user_created
Topic offset 증가 이벤트가 브로커 로그에 기록됩니다.
- Notifications log
Consumer Welcome email sent consumer group이 이벤트를 읽고 처리했습니다.
| 멈춘 위치 | 확인할 것 | 다음 조치 |
|---|---|---|
| Broker | 9092 포트, advertised listeners | Docker Compose 설정과 컨테이너 로그 확인 |
| Producer | `KAFKA_SERVICE`, broker 주소, topic 이름 | Users Service 연결과 `emit()` 호출 확인 |
| Consumer | `groupId`, `@EventPattern`, payload 구조 | Notifications Service 로그와 handler 실행 확인 |
Kafka 이벤트 기반 설계는 연결 자체보다 실패 후에도 같은 상태로 수렴하는지가 더 중요합니다.
아래 다이어그램은 토픽 계약, 컨슈머 처리 경계, 재시도와 DLQ를 한 흐름으로 묶어 점검하는 기준을 정리합니다.
브로커에 이벤트가 남아 있기 때문에 소비자는 다시 읽을 수 있습니다. 그래서 재시도, 중복 처리, DLQ 규칙을 함께 정해야 합니다.
- 이벤트 소비
Read 이벤트 소비 consumer가 `user_created` offset을 읽습니다.
- 업무 처리
Handle 업무 처리 환영 메일 발송처럼 외부 작업을 수행합니다.
- 일시 실패 재시도
Retry 일시 실패 재시도 네트워크 오류는 지연 후 다시 처리합니다.
- 계속 실패한 이벤트 격리
DLQ 계속 실패한 이벤트 격리 형식 오류나 영구 실패는 별도 토픽으로 보냅니다.
- Idempotency
같은 이벤트가 두 번 와도 결과가 한 번 처리된 것과 같아야 합니다.
- Offset commit
업무 처리가 끝난 뒤 읽은 위치를 확정해야 손실을 줄입니다.
- DLQ review
죽은 편지 토픽은 알림과 재처리 도구가 함께 있어야 합니다.
Kafka 기반 아키텍처는 서비스 간 동기 호출을 줄이고 이벤트 스트림을 중심으로 비동기 흐름을 구성합니다.
구현할 때는 토픽 설계, 파티션 키, 컨슈머 그룹, 재시도, DLQ, 중복 처리 기준을 함께 정해야 합니다.
이 장에서는 외부 요청을 처리하는 REST API와 서비스 간 통신에 쓰는 TCP, gRPC, Kafka를 비교하고, 서비스 경계와 통신 패턴을 선택할 때 확인할 기준을 정리했습니다.
마지막으로 Kafka를 도입할 때 반드시 분리해서 봐야 하는 토픽, 소비 처리, 재처리 책임을 한 번 더 정리합니다.
NestJS는 Kafka transport를 쉽게 연결해주지만, 이벤트가 어떤 계약으로 오래 살아남을지는 애플리케이션이 정해야 합니다.
- 업무 사건을 기록한다
Topic 업무 사건을 기록한다 토픽 이름, payload, key가 변경 이력을 설명해야 합니다.
- 처리 책임을 분리한다
Consumer 처리 책임을 분리한다 groupId로 작업 단위를 나누고, 장애가 다른 소비자에 번지지 않게 합니다.
- 실패 후 수렴을 설계한다
Recovery 실패 후 수렴을 설계한다 retry, DLQ, idempotency가 있어야 재처리가 안전합니다.
| 책임 | 정할 기준 | 검증 신호 |
|---|---|---|
| 발행 | 언제 이벤트를 쓰고 어떤 payload를 담는가 | producer 로그와 topic offset 증가 |
| 소비 | 어떤 group이 어떤 부작업을 처리하는가 | consumer 로그와 처리 성공 기록 |
| 복구 | 실패 이벤트를 언제 재시도하고 언제 DLQ로 보내는가 | 재시도 횟수, DLQ 건수, 중복 처리 여부 |