Kafka 입문: Topic과 Partition, Node.js Producer·Consumer, Consumer Group, 주문 처리 예제

이 글의 핵심

Kafka 설치와 Node.js 클라이언트 사용법, Partition과 Consumer Group이 병렬 처리와 순서 보장에 미치는 영향, 주문 처리 예제와 실패 시 재처리 방법을 다룹니다.

이 글의 핵심

Apache Kafka로 대규모 이벤트 스트리밍을 구축하는 글입니다. Producer, Consumer, Topic, Partition, Replication, Consumer Group까지 실전 예제로 정리했습니다.

실무에서 마주치는 문제들

초당 10만 건의 로그가 발생해요

데이터베이스로 감당 안 됩니다. Kafka로 처리합니다.

여러 서비스가 같은 데이터를 필요로 해요

각각 API를 호출합니다. Kafka로 한 번 발행, 여러 곳에서 소비합니다.

실시간 분석이 필요해요

배치 처리는 느립니다. Kafka Streams로 실시간 처리합니다.

세 경우 모두 핵심은 같습니다. Kafka는 메시지를 받아서 전달하고 지우는 “우체통”이 아니라, 들어온 순서대로 디스크에 추가만 하는 로그를 저장하고 여러 소비자가 각자 원하는 위치부터 읽어 가게 하는 시스템입니다. 그래서 소비자가 느려도 생산자가 막히지 않고, 새 서비스를 나중에 붙여도 보존 기간 안의 과거 이벤트를 처음부터 다시 읽을 수 있습니다. 반대로 “메시지 하나하나를 처리했는지 추적하고, 실패한 것만 다시 보내는” 작업 큐 방식에는 RabbitMQ 같은 전통적인 브로커가 더 자연스럽습니다. 두 방식의 차이는 메시징 시스템 비교 글에서 더 자세히 다룹니다.


Kafka가 로그·이벤트 처리에 쓰이는 이유

핵심 특징

Apache Kafka는 분산 이벤트 스트리밍 플랫폼입니다. 주요 개념:

  • Topic: 메시지 카테고리

  • Partition: 병렬 처리 단위

  • Producer: 메시지 발행

  • Consumer: 메시지 소비

  • Broker: Kafka 서버 성능:

  • 처리량: 클러스터 규모에 따라 초당 수십만~수백만 건

  • 지연 시간: 설정에 따라 수 ms 수준

이 숫자는 브로커 여러 대, 적절한 배치 설정, 작은 메시지를 전제로 한 값이라 그대로 기대하면 안 됩니다. Kafka의 높은 처리량은 순차 디스크 쓰기와 OS 페이지 캐시, 그리고 여러 메시지를 묶어 보내는 배치에서 나옵니다. 프로듀서의 linger.ms를 늘리면 배치가 커져 처리량은 올라가지만 메시지 하나의 지연은 그만큼 늘어나는 트레이드오프가 있습니다. 개념 사이의 관계도 짚어 두면, Topic은 논리적인 이름이고 실제 데이터는 Topic을 나눈 Partition 단위로 브로커에 저장됩니다. 각 Partition 안의 메시지는 0부터 증가하는 offset으로 식별되며, 소비자는 “어느 Partition의 몇 번 offset까지 읽었는지”를 기록하는 방식으로 진행 상황을 관리합니다.


Docker Compose로 로컬 브로커 띄우기

Docker Compose

# docker-compose.yml
version: '3.8'
services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.6.1
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000
  kafka:
    image: confluentinc/cp-kafka:7.6.1
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
docker-compose up -d

이미지 태그를 latest 대신 7.6.1로 고정한 데는 이유가 있습니다. Confluent Platform 8.0부터는 ZooKeeper 지원이 제거되어 cp-kafka:latest가 KRaft 모드 전용 이미지가 되었기 때문에, 위처럼 KAFKA_ZOOKEEPER_CONNECT를 쓰는 설정으로 latest를 받으면 브로커가 시작하지 못합니다. Apache Kafka 자체도 4.0에서 ZooKeeper를 완전히 제거했으므로 새로 구성한다면 KRaft 모드를 쓰는 것이 맞지만, 이 글에서는 여전히 많은 튜토리얼과 기존 시스템에서 볼 수 있는 ZooKeeper 구성을 그대로 두었습니다. 버전을 고정해 두는 습관은 “어제까지 되던 docker-compose가 오늘 안 되는” 문제를 막아 줍니다.

KAFKA_ADVERTISED_LISTENERS는 브로커가 클라이언트에게 “나에게 접속하려면 이 주소로 오라”고 알려 주는 값입니다. 여기서는 localhost:9092라서 호스트에서 실행한 Node.js 앱은 잘 붙지만, 같은 compose 네트워크의 다른 컨테이너에서 접속하면 브로커가 알려 준 localhost가 그 컨테이너 자신을 가리켜 연결이 실패합니다. 처음 Kafka를 도커로 띄울 때 가장 많이 겪는 문제가 이것인데, 첫 연결(bootstrap)은 성공하고 그 뒤 메타데이터로 받은 주소로 재접속하다가 실패하기 때문에 원인을 찾기 어렵습니다. 컨테이너 안팎에서 모두 접속해야 한다면 내부용(kafka:29092)과 외부용(localhost:9092) 리스너를 따로 정의해야 합니다. KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1은 브로커가 한 대뿐이라서 필요한 설정으로, 기본값 3을 그대로 두면 컨슈머 그룹의 offset을 저장할 내부 토픽을 만들지 못해 컨슈머가 동작하지 않습니다.


KafkaJS로 Producer와 Consumer 작성하기

설치

npm install kafkajs

Producer

// producer.ts
import { Kafka } from 'kafkajs';
const kafka = new Kafka({
  clientId: 'my-app',
  brokers: ['localhost:9092'],
});
const producer = kafka.producer();
async function sendMessage() {
  await producer.connect();
  await producer.send({
    topic: 'user-events',
    messages: [
      {
        key: 'user-123',
        value: JSON.stringify({
          type: 'user_created',
          userId: 123,
          email: '[email protected]',
          timestamp: new Date().toISOString(),
        }),
      },
    ],
  });
  console.log('Message sent');
  await producer.disconnect();
}
sendMessage();

send의 key는 이 메시지가 들어갈 Partition을 정하는 값입니다. KafkaJS의 기본 파티셔너는 키를 해시해서 Partition을 고르므로, 같은 user-123 키를 가진 메시지는 항상 같은 Partition에 들어가고 그 안에서 순서가 유지됩니다. 키를 생략하면 여러 Partition에 흩어지므로 같은 사용자의 이벤트라도 순서가 보장되지 않습니다. value는 바이트로 전송되므로 여기서는 JSON.stringify로 문자열을 만들었고, 스키마가 자주 바뀌는 시스템이라면 Avro나 Protobuf와 스키마 레지스트리를 함께 쓰는 경우가 많습니다.

이 예제는 메시지 하나를 보내고 바로 disconnect()하지만, 실제 서비스에서 이렇게 요청마다 연결하고 끊으면 매번 TCP 연결과 메타데이터 조회 비용이 들어 성능이 크게 떨어집니다. 프로듀서는 애플리케이션 시작 시 한 번 연결해 두고 계속 재사용하며, 종료 시그널(SIGTERM)을 받았을 때 disconnect()해서 보내지 못한 메시지를 정리하는 것이 일반적입니다. 참고로 KafkaJS는 널리 쓰였지만 2023년 이후 새 릴리스가 거의 없어서, 새 프로젝트라면 Confluent가 관리하는 @confluentinc/kafka-javascript처럼 유지보수되는 클라이언트도 함께 검토해 보는 것이 좋습니다.

Consumer

// consumer.ts
import { Kafka } from 'kafkajs';
const kafka = new Kafka({
  clientId: 'my-app',
  brokers: ['localhost:9092'],
});
const consumer = kafka.consumer({ groupId: 'user-service' });
async function consumeMessages() {
  await consumer.connect();
  await consumer.subscribe({ topic: 'user-events', fromBeginning: true });
  await consumer.run({
    eachMessage: async ({ topic, partition, message }) => {
      const event = JSON.parse(message.value.toString());
      console.log(`Received from partition ${partition}:`, event);
      // 이벤트 처리
      await handleEvent(event);
    },
  });
}
async function handleEvent(event: any) {
  console.log('Processing:', event);
  // 비즈니스 로직
}
consumeMessages();

fromBeginning: true는 이 컨슈머 그룹이 처음 구독할 때, 즉 저장된 offset이 없을 때만 적용됩니다. 한 번 실행해서 offset이 커밋된 뒤에는 이 옵션과 상관없이 마지막으로 커밋한 위치부터 읽습니다. “fromBeginning을 켰는데 왜 처음부터 안 읽히지?”라는 질문의 답이 대부분 이것이고, 정말 처음부터 다시 읽으려면 새 groupId를 쓰거나 kafka-consumer-groups.sh --reset-offsets로 offset을 되돌려야 합니다.

KafkaJS는 eachMessage가 에러 없이 끝나면 해당 메시지를 처리한 것으로 보고 offset을 주기적으로 자동 커밋합니다. 반대로 eachMessage에서 예외가 나면 같은 메시지를 재시도하고, 재시도가 계속 실패하면 컨슈머가 멈추거나 재시작합니다. 그래서 JSON.parse가 실패하는 잘못된 메시지 하나 때문에 그 Partition 전체의 소비가 막히는 poison pill 문제가 생길 수 있습니다. 뒤의 “에러 처리” 절에서 이 문제를 다룹니다. message.value는 tombstone 메시지(값이 null인 삭제 표시)에서는 null이므로, 실제 코드에서는 .toString() 전에 확인하는 편이 안전합니다.


Consumer Group으로 병렬 소비하기

병렬 처리

// consumer1.ts
const consumer1 = kafka.consumer({ groupId: 'analytics' });
await consumer1.subscribe({ topic: 'events' });
// consumer2.ts
const consumer2 = kafka.consumer({ groupId: 'analytics' });
await consumer2.subscribe({ topic: 'events' });
// 같은 그룹의 Consumer는 Partition을 나눠서 처리
// Partition 0 → Consumer 1
// Partition 1 → Consumer 2

Consumer Group은 Kafka에서 병렬 처리와 브로드캐스트를 동시에 가능하게 하는 장치입니다. 같은 groupId를 가진 컨슈머들은 Topic의 Partition을 나눠 가지므로 각 메시지는 그룹 안에서 한 컨슈머만 처리하고, 다른 groupId를 가진 그룹은 같은 메시지를 독립적으로 모두 받습니다. 앞의 user-service와 여기의 analytics가 같은 Topic을 구독해도 서로 영향을 주지 않는 이유입니다.

중요한 제약은 한 Partition은 그룹 안에서 한 컨슈머에게만 할당된다는 점입니다. Partition이 3개인 Topic에 같은 그룹의 컨슈머를 5개 띄우면 2개는 할당받을 Partition이 없어 놀게 됩니다. 컨슈머를 늘려 처리량을 높이려면 Partition 수가 그만큼 있어야 하므로, Partition 수가 곧 그룹의 최대 병렬도입니다. 컨슈머가 추가되거나 빠질 때는 Partition을 다시 나누는 리밸런스가 일어나고, 그동안 소비가 잠시 멈춥니다. 처리 시간이 긴 메시지 때문에 컨슈머가 제때 응답하지 못하면 브로커가 그 컨슈머를 죽은 것으로 판단해 리밸런스가 반복되는 경우가 있는데, 이때는 처리 시간을 줄이거나 sessionTimeout 같은 설정을 조정해야 합니다.


Topic 파티션 수와 메시지 Key

Topic 생성

# 3개의 Partition으로 생성
kafka-topics.sh --create \
  --topic orders \
  --bootstrap-server localhost:9092 \
  --partitions 3 \
  --replication-factor 1

Partition 수는 나중에 늘릴 수는 있지만 줄일 수는 없습니다. 더 큰 문제는 늘리는 순간 키 해시가 가리키는 Partition이 바뀐다는 점입니다. hash(key) % 3이던 것이 hash(key) % 6이 되므로, 같은 사용자의 이전 이벤트는 Partition 0에 있고 새 이벤트는 Partition 4로 가는 식으로 순서 보장이 깨질 수 있습니다. 순서가 중요한 Topic이라면 예상 최대 병렬도를 고려해 처음부터 여유 있게 잡는 편이 낫습니다. --replication-factor 1은 브로커가 한 대인 로컬 환경용이며, 운영에서는 보통 3으로 두고 min.insync.replicas=2와 프로듀서의 acks=all(KafkaJS에서는 acks: -1)을 함께 설정해 브로커 한 대가 죽어도 데이터가 남도록 합니다.

Key로 Partition 지정

// 같은 userId는 같은 Partition으로
await producer.send({
  topic: 'user-events',
  messages: [
    {
      key: 'user-123',  // Partition 결정
      value: JSON.stringify(event),
    },
  ],
});

키를 무엇으로 정할지는 “어떤 단위로 순서가 보장되어야 하는가”로 판단합니다. 주문 상태 변경이라면 orderId, 사용자 활동이라면 userId가 자연스럽습니다. 다만 특정 키에 트래픽이 몰리면(예: 대형 판매자 한 곳) 그 키가 들어가는 Partition 하나만 과부하가 걸리는 핫 파티션 문제가 생기므로, 키 분포가 고른지도 함께 확인해야 합니다.


주문 서비스와 결제 서비스를 Kafka로 잇기

Producer (주문 서비스)

// services/order-service.ts
import { Kafka } from 'kafkajs';
const kafka = new Kafka({
  clientId: 'order-service',
  brokers: ['localhost:9092'],
});
const producer = kafka.producer();
export async function createOrder(order: any) {
  await producer.connect();
  // 주문 이벤트 발행
  await producer.send({
    topic: 'orders',
    messages: [
      {
        key: order.userId.toString(),
        value: JSON.stringify({
          type: 'order_created',
          orderId: order.id,
          userId: order.userId,
          amount: order.amount,
          timestamp: new Date().toISOString(),
        }),
      },
    ],
  });
  console.log('Order created:', order.id);
}

주문 이벤트의 키를 userId로 잡았기 때문에 같은 사용자의 주문은 같은 Partition에서 순서대로 처리됩니다. 여기서도 createOrder가 호출될 때마다 producer.connect()를 부르는데, KafkaJS는 이미 연결된 상태면 다시 연결하지 않지만 모듈 초기화 시점에 한 번 연결해 두는 쪽이 의도가 분명합니다.

이 구조에서 실무적으로 더 중요한 문제는 DB 저장과 이벤트 발행의 원자성입니다. 주문을 DB에 저장한 뒤 send를 호출하는 사이에 프로세스가 죽으면, 주문은 있는데 이벤트는 없는 상태가 됩니다. 순서를 바꿔 이벤트를 먼저 보내면 반대로 존재하지 않는 주문의 이벤트가 나갑니다. 이 문제를 해결하는 대표적인 방법이 아웃박스 패턴입니다. 주문과 “보낼 이벤트”를 같은 DB 트랜잭션으로 아웃박스 테이블에 저장하고, 별도 프로세스(또는 Debezium 같은 CDC 도구)가 그 테이블을 읽어 Kafka로 발행합니다.

Consumer (결제 서비스)

// services/payment-service.ts
const consumer = kafka.consumer({ groupId: 'payment-service' });
await consumer.connect();
await consumer.subscribe({ topic: 'orders' });
await consumer.run({
  eachMessage: async ({ message }) => {
    const event = JSON.parse(message.value.toString());
    if (event.type === 'order_created') {
      console.log('Processing payment for order:', event.orderId);
      
      // 결제 처리
      const paymentResult = await processPayment(event);
      // 결제 완료 이벤트 발행
      await producer.send({
        topic: 'payments',
        messages: [
          {
            value: JSON.stringify({
              type: 'payment_completed',
              orderId: event.orderId,
              status: paymentResult.status,
            }),
          },
        ],
      });
    }
  },
});

결제 서비스는 orders를 소비하고 결과를 payments에 발행해, 주문 서비스와 결제 서비스가 서로의 API를 직접 호출하지 않고 이벤트로만 연결됩니다. 결제 서비스가 잠시 내려가도 주문은 계속 쌓이고, 다시 올라오면 밀린 이벤트를 처리합니다. (이 코드 조각에서는 kafka와 producer를 앞의 예제처럼 만들어 연결해 두었다고 가정합니다.)

가장 조심할 부분은 같은 주문이 두 번 결제될 수 있다는 점입니다. processPayment가 성공한 직후, offset이 커밋되기 전에 컨슈머가 재시작되거나 리밸런스가 일어나면 그 메시지는 다시 전달됩니다. 결제 같은 작업에서 이것을 막으려면 orderId를 결제 대행사의 멱등성 키로 넘기거나, “이 주문은 결제 처리됨”을 유니크 제약이 있는 테이블에 기록해 두 번째 처리를 건너뛰어야 합니다. Kafka의 트랜잭션 기능으로 “소비 offset 커밋과 payments 발행”을 원자적으로 묶을 수는 있지만, 외부 결제 API 호출은 그 트랜잭션에 포함되지 않기 때문에 멱등성은 여전히 애플리케이션이 책임져야 합니다. payments 메시지에 키가 없는 것도 보완할 점입니다. orderId를 키로 주면 같은 주문의 결제 이벤트가 순서대로 처리됩니다.


실패한 메시지 재시도

재시도

await consumer.run({
  eachMessage: async ({ message }) => {
    try {
      await processMessage(message);
    } catch (error) {
      console.error('Failed to process:', error);
      
      // Dead Letter Queue로 전송 (원본 키·값과 실패 정보를 함께 보존)
      await producer.send({
        topic: 'failed-messages',
        messages: [{
          key: message.key,
          value: message.value,
          headers: {
            ...message.headers,
            'x-error': String(error),
          },
        }],
      });
    }
  },
});

Dead Letter Queue(DLQ)는 처리할 수 없는 메시지를 별도 Topic으로 옮겨 두고 원래 흐름은 계속 진행하게 하는 패턴입니다. 예외를 다시 던지면 KafkaJS가 같은 메시지를 계속 재시도하면서 Partition 전체가 막히지만, DLQ로 보내고 정상 종료하면 offset이 넘어가서 다음 메시지를 처리할 수 있습니다. 원본 메시지 객체를 그대로 넘기는 대신 키·값·헤더만 골라 담고 에러 내용을 헤더에 남긴 이유는, 나중에 DLQ를 들여다보는 사람이 “왜 실패했는지”를 알 수 있어야 재처리 여부를 판단할 수 있기 때문입니다.

모든 실패를 바로 DLQ로 보내는 것이 정답은 아닙니다. 잘못된 JSON처럼 다시 시도해도 절대 성공하지 않는 실패는 즉시 DLQ로 보내는 것이 맞지만, DB 연결 끊김처럼 일시적인 실패까지 DLQ로 보내면 장애가 났던 몇 분 동안의 메시지가 모두 DLQ에 쌓여 수작업으로 재처리해야 합니다. 저는 이런 구조를 처음 운영할 때 이 구분 없이 모든 예외를 DLQ로 보냈다가, DB 점검 시간에 들어온 메시지를 한꺼번에 다시 넣어야 했던 상황이 가장 흔한 실수라고 생각합니다. 일시적 실패는 몇 번 재시도(지수 백오프)한 뒤에도 실패할 때만 DLQ로 보내고, DLQ의 크기를 모니터링해 알림을 거는 것이 일반적입니다. DLQ 전송 자체가 실패하는 경우도 있으므로, 그때는 예외를 던져 원래 메시지를 커밋하지 않게 해야 메시지를 잃지 않습니다.


취업·면접과 연결하기. 파티션·컨슈머 그룹·재처리는 이벤트 스트리밍·시스템 설계 면접과 잘 맞습니다. 개발자 기술 면접 준비: 알고리즘부터 시스템 설계까지와, 이력서에 트래픽·처리량을 쓰는 법은 개발자 이력서·서류·면접 가이드를 참고하세요.


Kafka 요약과 운영 전에 정할 것

  • Kafka: 분산 이벤트 스트리밍
  • Topic: 메시지 카테고리
  • Partition: 병렬 처리
  • Consumer Group: 부하 분산
  • Replication: 고가용성
  • Offset: 컨슈머 그룹별 진행 위치, 처리 후 커밋

운영 전에 결정해야 할 것

로컬에서 돌아가는 예제를 운영으로 옮길 때 가장 먼저 정해야 하는 것은 전달 보장 수준입니다. 대부분의 시스템은 “최소 한 번은 전달되고, 중복은 소비자가 멱등하게 처리한다”를 기준으로 삼고, 프로듀서의 acks, 복제 수, 컨슈머의 커밋 시점을 그 기준에 맞춰 설정합니다. 다음으로 Partition 수와 키 설계를 정하는데, 앞에서 본 것처럼 Partition은 줄일 수 없고 늘리면 키 매핑이 바뀌므로 초기 결정이 오래 영향을 미칩니다. 마지막으로 컨슈머 랙(가장 최신 offset과 커밋된 offset의 차이)을 모니터링해야 합니다. 랙이 꾸준히 늘어난다면 처리 속도가 유입 속도를 따라가지 못한다는 뜻이며, 보존 기간이 지나 아직 읽지 않은 메시지가 삭제되기 전에 조치해야 합니다.


같이 보면 좋은 글


자주 묻는 질문 (FAQ)

Q. Kafka vs RabbitMQ, 어떤 게 나은가요?

A. Kafka는 대용량 이벤트 스트리밍에 유리합니다. RabbitMQ는 작업 큐에 유리합니다. 로그/이벤트는 Kafka, 작업 큐는 RabbitMQ를 권장합니다.

Q. Zookeeper가 뭔가요?

A. Kafka 클러스터의 메타데이터(브로커 목록, 컨트롤러 선출, 토픽 설정)를 관리하던 별도 시스템입니다. Kafka 3.3부터 ZooKeeper 없이 브로커끼리 메타데이터를 관리하는 KRaft 모드가 운영용으로 사용 가능해졌고, Kafka 4.0에서는 ZooKeeper 지원이 완전히 제거되었습니다. 새로 구축한다면 KRaft 모드를 쓰면 됩니다.

Q. 메시지 순서가 보장되나요?

A. 같은 Partition 내에서는 순서가 보장됩니다. 다른 Partition 간에는 보장되지 않습니다. 그래서 순서가 필요한 이벤트는 같은 키로 보내 같은 Partition에 모읍니다. 프로듀서가 재시도할 때 순서가 뒤바뀌는 것을 막으려면 멱등 프로듀서(KafkaJS의 idempotent: true)를 켜는 것이 좋습니다.