본문으로 건너뛰기

안동민 개발노트

본문 시작

Apache Kafka와 이벤트 기반 아키텍처

Kafka 토픽과 소비자 그룹의 역할을 이해하고 사용자 생성 이벤트를 발행해 알림 서비스에서 비동기로 처리합니다.

지난 절에서는 NestJS와 gRPC를 활용하여 고성능 서비스 간 동기 통신을 구현하는 방법을 살펴보았습니다.

이제 7장의 마지막 절로, 마이크로서비스 아키텍처에서 비동기 통신과 이벤트 기반 아키텍처(EDA)의 핵심 요소인 Apache Kafka를 NestJS에 적용하는 방법에 대해 알아보겠습니다.

동기 통신(Synchronous Communication)은 요청-응답 패턴에 적합하지만, 서비스 간 강한 결합을 유발하고 한 서비스의 장애가 다른 서비스로 전파될 위험이 있습니다.

또한 대규모 이벤트를 처리하거나 실시간 데이터 파이프라인을 구축할 때는 비효율적입니다.

Apache Kafka는 이벤트를 파티션 로그에 보관하고 생산자와 소비자의 처리 시점을 분리하는 분산 스트리밍 플랫폼입니다. 실제 처리량과 장애 허용 범위는 파티션 수, 복제, 승인, 보관 정책과 실행 환경에 따라 달라집니다.


이벤트 기반 아키텍처 (EDA)와 Kafka

이벤트 기반 아키텍처(EDA, Event-Driven Architecture)는 시스템의 구성 요소들이 직접적으로 호출하는 대신, 이벤트(Event)를 발행하고 구독하는 방식으로 상호작용하는 아키텍처 스타일입니다.

각 서비스는 자신이 관심 있는 이벤트를 구독하고, 해당 이벤트가 발생하면 필요한 작업을 수행합니다.

EDA의 주요 특징
  • 느슨한 결합(Loose Coupling): 서비스들이 서로를 직접 호출하지 않으므로, 한 서비스의 변경이 다른 서비스에 미치는 영향이 최소화됩니다.
  • 확장성: 이벤트 생산자와 소비자 모두 독립적으로 확장할 수 있습니다.
  • 복원력(Resilience): 이벤트가 브로커에 성공적으로 기록되고 보관 기간 안에 남아 있다면, 소비자가 잠시 중단되어도 다시 읽어 처리할 수 있습니다.
  • 비동기 통신: 응답을 즉시 기다리지 않으므로, 장시간 실행되는 작업을 처리하거나 백그라운드 작업을 실행하는 데 적합합니다.
  • 낮은 지연 처리: 데이터 변경 이벤트를 지속적으로 전달할 수 있지만, 실제 지연은 생산자·브로커·소비자 설정과 부하에 따라 달라집니다.

Apache Kafka는 이러한 EDA를 구현하는 데 가장 널리 사용되는 플랫폼 중 하나입니다.

Kafka는 분산된 커밋 로그(Commit Log) 또는 분산 메시지 브로커로 볼 수 있으며, 다음과 같은 특징을 가집니다.

  • 처리량 확장: 토픽을 여러 파티션으로 나누고 생산자·소비자를 확장할 수 있습니다. 달성 가능한 처리량은 워크로드와 클러스터 구성에 따라 달라집니다.
  • 장애 허용: 파티션 복제와 ISR, 생산자 acks 설정을 조합해 장애 허용 범위를 정합니다. 복제만으로 모든 장애에서 무손실이 보장되는 것은 아닙니다.
  • 확장성: 수평적 확장이 용이하여 데이터 양이 증가해도 유연하게 대응할 수 있습니다.
  • 영속성: broker가 승인한 이벤트는 디스크 로그에 기록되며, 소비자가 다시 읽을 수 있는 범위는 topic의 cleanup·retention 정책에 따라 달라집니다.
  • 다중 소비자: 서로 다른 consumer group은 같은 이벤트 스트림을 독립적으로 읽고, 같은 group의 인스턴스들은 파티션을 나눠 처리합니다.
Kafka의 핵심 개념
  • Producer(생산자): Kafka 토픽으로 메시지(이벤트)를 발행하는 애플리케이션.
  • Consumer(소비자): Kafka 토픽에서 메시지를 구독하고 처리하는 애플리케이션.
  • Topic(토픽): 메시지(이벤트)가 발행되고 소비되는 논리적인 카테고리.
  • Partition(파티션): 토픽 로그를 나누는 단위. 레코드 순서는 파티션 안에서만 정의되며, 같은 업무 키의 순서가 필요하면 그 키가 일관되게 같은 파티션으로 가도록 설계해야 합니다.
  • Broker(브로커): 파티션 로그를 저장하고 생산·소비 요청을 처리하는 Kafka 서버. 여러 브로커의 복제와 배치는 설정한 장애 허용 범위를 구성합니다.
  • KRaft controller: 메타데이터 요청과 로그를 관리하는 서버 역할입니다. 여러 controller가 Raft metadata quorum을 구성하며, Kafka 4.x는 KRaft 모드만 지원합니다.

Kafka 이벤트는 단순한 메시지가 아니라 재처리, 순서, 소비자 확장까지 포함한 운영 계약입니다.

토픽을 만들기 전에 아래 기준으로 이벤트 경계를 먼저 잡아두면 구현과 검증이 쉬워집니다.

HTTP 사용자 생성 요청이 NestJS emit 발행, Kafka 토픽 파티션과 consumer group 배분, EventPattern 처리로 이어지고 각 경계를 별도로 관측하는 흐름

Nest · Kafka Event Path

직접 서비스 호출 대신 Users가 사건을 토픽에 남기고 Notifications가 자기 속도로 읽습니다. HTTP 결과, 브로커 기록, 소비자 업무 결과는 서로 다른 성공 신호입니다.

publish → consume

요청 진입부터 소비 완료까지

  1. HTTP 경계

    POST /users를 받은 UsersController가 사용자를 생성합니다. 사용자 저장 성공과 이벤트 전달 성공은 별도 결과이므로 둘을 하나의 성공으로 간주하지 않습니다.

  2. NestJS 이벤트 발행

    emit에 토픽 'user_created', 업무 키, payload를 전달합니다. 이는 reply topic이 없는 이벤트 방식이므로 subscribeToResponseOf가 필요하지 않습니다.

  3. Broker 기록과 partition 선택

    생산이 성공하면 broker가 선택된 partition 로그에 record와 offset을 남깁니다. 같은 키의 순서는 그 키가 같은 partition으로 갈 때 해당 partition 안에서만 정의됩니다.

  4. Consumer group 배분

    notification-group의 인스턴스들은 partition을 나눠 맡습니다. 다른 group은 같은 스트림을 독립적으로 읽지만, group 자체가 재전달이나 업무 중복을 제거하지는 않습니다.

  5. NestJS 이벤트 처리

    @EventPattern과 토픽 'user_created'가 연결되고, @Payload로 value를 받아 환영 메일을 보냅니다. 성공 후 offset을 확정하는 경계에서도 장애 시 재처리가 생길 수 있습니다.

HTTP ≠ event

응답 시점을 계약으로 정한다

emit은 hot Observable로 즉시 발행을 시도합니다. 완료를 기다리면 생산자 오류를 요청 경계에서 관측할 수 있지만 소비자 처리 완료까지 기다리는 것은 아닙니다. 먼저 HTTP를 반환한다면 미발행 복구 책임을 별도로 둡니다.

wiring

로컬 이름과 Kafka 계약을 구분한다

'KAFKA_SERVICE'는 Nest DI token입니다. clientId는 클라이언트를 식별하고 brokers는 bootstrap 주소를 제공하며, groupId는 소비 진행 위치와 partition 배분의 단위가 됩니다.

broker

저장은 승인·복제·보관 설정에 달려 있다

Broker는 topic의 partition 로그를 보관합니다. 생산자 acks, ISR과 복제 계수는 승인 가능한 장애 범위를 정하고, retention은 소비 여부와 별개로 재읽기 가능한 기간을 정합니다.

delivery

Consumer group은 중복 제거 장치가 아니다

같은 group 안에서는 한 partition을 한 시점에 한 member가 담당하지만, 처리 후 commit 전에 중단되면 record가 다시 전달될 수 있습니다. 메일·정산 같은 부작용은 event ID로 멱등하게 처리합니다.

observability

이벤트가 멈춘 경계를 순서대로 좁힌다

1 · Runtime / broker

컨테이너 상태와 broker 로그, bootstrap 포트 9092, advertised listener를 확인합니다.

2 · Producer

HTTP trace와 event ID, emit 완료 또는 오류, 기록된 topic·partition·offset을 남깁니다.

3 · Topic / group

토픽의 log end offset, group의 current offset과 lag, partition assignment와 rebalance를 확인합니다.

4 · Handler

@EventPattern 진입, payload 검증, 업무 성공, retry와 DLQ 전환을 같은 event ID로 연결합니다.

검증 순서: broker 준비 → producer 발행 결과 → topic/partition 기록 → consumer group 진행 위치 → handler의 실제 부작용. API가 성공했다는 사실만으로 이 전체 경로가 성공했다고 결론내리지 않습니다.


NestJS에서 Kafka 설정 및 구현

NestJS는 @nestjs/microservices 패키지를 통해 Kafka 전송 계층을 지원합니다.

시나리오: Users 서비스에서 새로운 사용자가 생성되면 이벤트를 Kafka에 발행하고, Notifications 서비스에서 이 이벤트를 구독하여 사용자에게 환영 이메일을 보내는 시나리오를 구축해 보겠습니다.

사전 준비: 로컬 환경에서 Kafka broker를 실행해야 합니다.

Docker를 사용하는 것이 가장 편리합니다.

# 로컬 학습용 단일 KRaft broker (docker-compose.yml)
services:
  kafka:
    image: apache/kafka:4.3.1
    ports:
      - "9092:9092"

# 실행: docker compose up -d

공식 이미지의 기본 단일 노드 KRaft 구성은 로컬 학습용입니다. 운영 환경에서는 controller/broker 역할, 복제 계수, ISR, 인증과 보관 정책을 별도로 설계해야 합니다.

Kafka 사용자 이벤트 서비스 구축

새로운 사용자가 생성될 때 Kafka에 이벤트를 발행하는 서비스입니다.

단계 1: 필요한 패키지 설치

users-service 프로젝트에서 설치합니다.

npm install @nestjs/microservices @nestjs/platform-express kafkajs
npm install --save-dev @types/express # 웹 서버 역할을 겸할 경우
  • NestJS Kafka transporter는 kafkajs 드라이버를 런타임에 사용하므로 애플리케이션 의존성으로 함께 설치합니다. kafka-node는 이 예제에서 사용하지 않습니다.
단계 2: main.ts 파일 수정 (HTTP 서버 유지)

사용자 서비스는 클라이언트의 HTTP 요청을 받아 사용자 생성 후 Kafka 이벤트를 발행합니다.

users-service/src/main.ts
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();
단계 3: UsersModule에 Kafka 클라이언트 등록

ClientsModule을 사용하여 Kafka 클라이언트를 등록하고, 이를 UsersController에 주입하여 이벤트를 발행합니다.

users-service/src/users/users.module.ts (수정)
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, // 토픽이 없으면 자동으로 생성
          },
          producerOnlyMode: true, // 이벤트 발행 전용 클라이언트
        },
      },
    ]),
  ],
  controllers: [UsersController],
  providers: [UsersService],
})
export class UsersModule {}
  • name: 'KAFKA_SERVICE': 이 클라이언트 프록시를 의존성 주입을 통해 참조할 토큰입니다.
  • transport: Transport.KAFKA: Kafka 전송 방식을 사용합니다.
  • options.client: Kafka 클라이언트 설정 (고유한 clientId, 브로커 목록 brokers).
  • options.producer: Kafka 생산자 설정 (예: allowAutoTopicCreation으로 토픽 자동 생성 활성화).
단계 4: 사용자 컨트롤러 및 서비스 구현 (Kafka 이벤트 발행)

사용자 생성 시 Kafka emit 메서드를 통해 이벤트를 발행합니다.

users-service/src/users/users.controller.ts (수정)
import { Body, Controller, Get, Inject, OnModuleInit, Param, Post } from '@nestjs/common';
import { ClientKafkaProxy } from '@nestjs/microservices';
import { UsersService } from './users.service';
import { lastValueFrom } from 'rxjs';

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: ClientKafkaProxy,
  ) {}

  async onModuleInit() {
    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)}`);

    // emit()은 reply topic을 만들지 않는 이벤트 발행입니다.
    // 같은 사용자의 이벤트가 같은 파티션으로 가도록 key를 함께 보냅니다.
    await lastValueFrom(
      this.client.emit('user_created', {
        key: String(newUser.id),
        value: newUser,
      }),
    );

    return newUser;
  }

  // 기존 사용자 조회 엔드포인트 (HTTP)
  @Get(':id')
  getUserById(@Param('id') id: string): User {
    const userId = parseInt(id, 10);
    return this.usersService.findOne(userId);
  }
}
  • ClientKafkaProxy: Kafka 전송 방식을 위한 클라이언트 계약입니다.
  • onModuleInit(): client.connect()를 호출하여 애플리케이션 시작 시 브로커 연결 실패를 드러냅니다.
  • emit: 이벤트 발행용 hot Observable을 반환합니다. lastValueFrom()으로 완료를 기다리면 이 요청에서 생산자 오류를 관측할 수 있지만, 소비자의 업무 처리 결과를 기다리는 것은 아닙니다.
  • subscribeToResponseOf(): send()@MessagePattern()을 쓰는 요청-응답 패턴의 reply topic에만 필요합니다. emit()@EventPattern()에는 사용하지 않습니다.

위 예제는 생산자 승인을 관측하지만 사용자 저장과 Kafka 기록을 하나의 원자적 트랜잭션으로 만들지는 않습니다. 사용자 저장 성공 뒤에도 이벤트 발행을 반드시 이어가야 한다면 같은 데이터베이스 트랜잭션에 outbox 레코드를 남기고 별도 relay가 Kafka에 발행하도록 설계합니다.

단계 5: UsersService 구현 (동일)

users-service/src/users/users.service.ts 파일은 이전과 동일하게 유지됩니다.

단계 6: AppModuleUsersModule 임포트 (동일)

users-service/src/app.module.ts는 변경 없이 UsersModule을 임포트합니다.

알림 서비스 구축

user_created 이벤트를 구독하여 처리하는 서비스입니다.

단계 1: 새 NestJS 프로젝트 생성
nest new notifications-service --skip-install
cd notifications-service
npm install @nestjs/microservices kafkajs
단계 2: main.ts 파일 수정 (Kafka 마이크로서비스 서버 설정)
notifications-service/src/main.ts
import { NestFactory } from '@nestjs/core';
import { AppModule } from './app.module';
import { MicroserviceOptions, Transport } from '@nestjs/microservices';
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: Nest가 생략 시 기본값을 제공하므로 문법적으로 필수는 아닙니다. 하지만 @EventPattern()으로 topic을 구독하는 운영 서비스는 목적별로 안정적인 ID를 명시해야 배포 후에도 같은 offset과 partition 배분 단위를 유지할 수 있습니다. 같은 유효 group ID의 인스턴스들은 파티션을 나누어 소비합니다.
단계 3: 알림 마이크로서비스 컨트롤러 (Kafka 핸들러) 구현

@EventPattern() 데코레이터를 사용하여 특정 Kafka 토픽에서 발생하는 이벤트를 처리합니다.

notifications-service/src/notifications/notifications.controller.ts (새로 생성)
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(): 이벤트 메시지의 페이로드(값)를 추출하여 메서드 인자로 주입합니다.
단계 4: 알림 서비스 구현 (이메일 전송 시뮬레이션)
notifications-service/src/notifications/notifications.service.ts (새로 생성)
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}!`);
  }
}
단계 5: NotificationsModuleAppModule 구성
notifications-service/src/notifications/notifications.module.ts
import { Module } from '@nestjs/common';
import { NotificationsController } from './notifications.controller';
import { NotificationsService } from './notifications.service';

@Module({
  controllers: [NotificationsController],
  providers: [NotificationsService],
})
export class NotificationsModule {}
notifications-service/src/app.module.ts
import { Module } from '@nestjs/common';
import { NotificationsModule } from './notifications/notifications.module';

@Module({
  imports: [NotificationsModule],
  controllers: [],
  providers: [],
})
export class AppModule {}

컨테이너 환경 실행 검증

Kafka broker 실행
  • Docker Compose 파일이 있는 디렉토리에서: docker compose up -d
  • 컨테이너가 정상적으로 실행되었는지 확인: docker compose ps
Notifications Service (Kafka Consumer) 시작
  • cd notifications-service
  • npm run start:dev
  • 콘솔에 Notifications Microservice (Kafka Consumer) is listening. 메시지 확인.
Users Service (Kafka Producer) 시작
  • cd users-service
  • npm run start:dev
  • 콘솔에 Users Service (HTTP Server) is listening on port 3000 메시지 확인.
사용자 생성 API 호출 (Postman 또는 cURL)
  • 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 응답보다 이벤트가 브로커를 지나 소비자까지 도착했는지가 핵심입니다.

아래 다이어그램은 컨테이너, 프로듀서, 토픽, 컨슈머 로그를 순서대로 확인하는 체크 포인트를 보여줍니다.

Kafka 이벤트 기반 설계는 연결 자체보다 실패 후에도 같은 상태로 수렴하는지가 더 중요합니다.

아래 다이어그램은 토픽 계약, 컨슈머 처리 경계, 재시도와 DLQ를 한 흐름으로 묶어 점검하는 기준을 정리합니다.

Kafka 이벤트의 schema, key, partition 순서, retention과 consumer group 계약을 정하고 outbox, producer ack, 멱등 처리, offset commit, retry와 DLQ로 재처리에 안전하게 수렴하는 설계

Kafka · Delivery Contract

Kafka가 record를 보관해도 데이터베이스 변경과 발행의 원자성, 소비 부작용의 중복 안전성, 실패 격리와 재생 절차는 애플리케이션 계약입니다.

schema

업무 사건을 설명한다

Topic은 user_created처럼 이미 일어난 사실로 이름 짓고, payload에는 event ID, 발생 시각, 업무 식별자와 소비자가 필요한 최소 데이터를 둡니다. schema 변경에는 호환성·version 규칙이 필요합니다.

key · order

순서 범위를 먼저 고른다

일관된 key와 partitioner는 같은 사용자의 record를 한 partition으로 모읍니다. 보장 범위는 그 partition log의 offset 순서이며, retry 때 생산자 전송 순서까지 보존하려면 idempotence와 동시 in-flight 같은 ordering 조건도 맞춰야 합니다.

retention

복구 창을 보관 정책에 맞춘다

cleanup.policy=delete에서는 retention.ms/bytes를 넘은 오래된 segment가 소비 여부와 무관하게 삭제 대상이 됩니다. 그 뒤의 재생은 보장할 수 없으며, compact는 key별 과거 값을 별도 규칙으로 제거합니다.

at-least-once convergence

발행과 소비의 두 단절점을 복구 가능하게 만든다

  1. 업무 상태와 outbox를 함께 확정한다

    사용자 저장과 outbox record를 같은 데이터베이스 트랜잭션에 넣습니다. 그러면 저장만 성공하고 발행 의도가 사라지는 간격을 relay가 나중에 복구할 수 있습니다.

  2. Relay가 key와 event ID를 유지해 발행한다

    Relay는 broker의 생산자 승인을 관측한 뒤 outbox를 발행 완료로 표시합니다. 실패하거나 확인이 모호하면 다시 발행할 수 있으므로 event ID가 소비 멱등성의 기준이 됩니다.

  3. Consumer가 부작용을 멱등하게 처리한다

    DB 변경은 inbox의 event ID와 업무 상태를 같은 트랜잭션에 기록합니다. 이메일·결제 같은 외부 부작용은 별도 단절점이므로 downstream idempotency key나 notification outbox·relay로 보호합니다.

  4. 업무 성공 뒤 offset을 commit한다

    처리 뒤 commit 전에 중단되면 같은 record가 재처리될 수 있습니다. 반대로 업무 성공 전에 진행 위치를 확정하면 실패한 작업을 건너뛸 수 있으므로 경계를 명시적으로 선택합니다.

  5. 일시 실패는 제한적으로 retry하고 영구 실패는 격리한다

    Backoff와 최대 횟수를 둔 뒤 schema 오류나 반복 실패를 DLQ topic으로 보냅니다. DLQ는 Kafka의 자동 복구 보장이 아니라 애플리케이션이 구현하고 운영해야 할 정책입니다.

producer durability

승인은 강도를 높일 뿐 절대 무손실은 아니다

acks=all은 사용 가능한 가장 강한 broker 승인입니다. 복제 계수, ISR과 min.insync.replicas, 디스크·운영 장애와 오류 처리를 함께 설계해야 하며, producer idempotence도 DB와 Kafka를 하나의 트랜잭션으로 만들지는 않습니다.

duplicate safety

Consumer group과 idempotency는 다른 책임이다

같은 group의 member들은 partition을 분담하지만 업무 성공 뒤 offset commit 완료 전에 중단되거나 outbox가 재발행되거나 운영 replay를 하면 같은 사건을 다시 볼 수 있습니다. DB 중복 억제는 event ID와 업무 상태의 원자적 기록으로 구현합니다.

retry

재시도가 partition 진행을 막을 수 있다

NestJS Kafka에서 처리되지 않은 event handler 예외는 KafkaJS로 전파되고 offset이 commit되지 않아 record가 재전달될 수 있습니다. 오류 분류, backoff, 최대 횟수와 timeout을 정하고 같은 partition의 뒤 record가 얼마나 기다릴지도 판단합니다.

DLQ operations

격리 뒤의 소유자와 replay를 준비한다

DLQ에는 원래 topic·key·payload·schema version·실패 이유와 시도를 남깁니다. 알림, 담당자, 수정·재생 도구, 재생 시 멱등성 검증이 없으면 DLQ는 복구가 아니라 적체가 됩니다.

group contract

업무 목적마다 독립 group을 둔다

메일, 분석, 감사 로그가 모두 같은 event를 필요로 하면 서로 다른 groupId를 사용해 각자의 offset과 lag를 가집니다. 한 group 안에서 인스턴스를 늘리면 partition을 분담하지만, consumer 수가 partition 수를 넘으면 유휴 인스턴스가 생길 수 있습니다.

운영 신호: publish 실패와 ack 지연, outbox 적체, topic log end offset, group lag·rebalance, handler 성공률, retry 횟수, DLQ 건수, 중복 억제 건수와 commit 진행 위치를 같은 event ID로 연결합니다.


Kafka 기반 아키텍처는 서비스 간 동기 호출을 줄이고 이벤트 스트림을 중심으로 비동기 흐름을 구성합니다.

구현할 때는 토픽 설계, 파티션 키, 컨슈머 그룹, 재시도, DLQ, 중복 처리 기준을 함께 정해야 합니다.

이 장에서는 외부 요청을 처리하는 REST API와 서비스 간 통신에 쓰는 TCP, gRPC, Kafka를 비교하고, 서비스 경계와 통신 패턴을 선택할 때 확인할 기준을 정리했습니다.

마지막으로 Kafka를 도입할 때는 토픽, 소비 처리, 재처리 책임을 분리해 관측하고 같은 이벤트의 재전달에도 안전하게 수렴하는지 검증해야 합니다.