학습 목표

시나리오

데이터 팀은 Python으로 ETL을 짜고, 프론트엔드 팀은 Node.js로 BFF를 운영합니다. 둘 다 orders 토픽을 읽어야 합니다.

문제가 된 것은 Java 기준으로 작성된 사내 가이드를 그대로 따랐을 때였습니다. Python 쪽에서 acks=allenable.idempotence=true는 잘 넘어갔지만, producer.produce() 뒤에 flush()를 부르지 않아 스크립트가 종료될 때 버퍼의 레코드가 사라졌습니다. Node 쪽에서는 eachMessage가 자동 커밋을 한다는 사실을 모른 채 예외를 삼켜 유실이 발생했습니다.

"프로토콜이 같으니 코드도 비슷할 것"이라는 가정이 사고의 원인이었습니다.

사전 요구사항

검증 환경 (2026-07 기준)
항목버전비고
Apache Kafka4.3.1예제 1의 클러스터
confluent-kafka (Python)2.15.0PyPI 최신 버전. librdkafka를 번들한 C 확장입니다
Python3.8 이상confluent-kafka 2.15.0의 requires_python
kafkajs (Node.js)2.2.4npm 최신 버전. 순수 JavaScript 구현입니다
Node.js14 이상kafkajs 2.2.4의 engines.node. 실무에서는 LTS를 쓰세요

설정 이름 대응표

librdkafka(Python/C/C++/Go 등)는 대부분 Java와 같은 이름을 쓰지만 일부는 다르고, 기본값이 다른 항목이 있습니다. kafkajs는 camelCase 옵션 객체로 완전히 다른 이름을 씁니다.

Java ↔ librdkafka ↔ kafkajs 대응 (Kafka 4.3 기준 Java 기본값)
개념 Java (기본값) librdkafka / Python kafkajs
부트스트랩 bootstrap.servers bootstrap.servers (또는 별칭 metadata.broker.list) brokers: []
확인 응답 acks (all) acks (별칭 request.required.acks) send({ acks: -1 })숫자로 지정 (-1=all)
멱등성 enable.idempotence (true) enable.idempotence createProducer({ idempotent: true })
배치 지연 linger.ms (5) linger.ms (별칭 queue.buffering.max.ms) 직접 대응 없음 — sendBatch로 애플리케이션이 묶습니다
배치 크기 batch.size (16384) batch.size 직접 대응 없음
전송 총 예산 delivery.timeout.ms (120000) delivery.timeout.ms (별칭 message.timeout.ms) send({ timeout })
그룹 ID group.id group.id createConsumer({ groupId })
오프셋 리셋 auto.offset.reset (latest) auto.offset.reset (기본값도 latest) run({ fromBeginning: true })불리언
자동 커밋 enable.auto.commit (true) enable.auto.commit (기본 true) run({ autoCommit: true }) (기본 true)
poll 간격 상한 max.poll.interval.ms (300000) max.poll.interval.ms createConsumer({ sessionTimeout, heartbeatInterval }) 조합으로 다룹니다
격리 수준 isolation.level (read_uncommitted) isolation.level run({ ... }) 옵션이 아니라 createConsumer({ readUncommitted }) 계열로 다룹니다 — 버전별 문서 확인 필요
압축 compression.type (none) compression.type (별칭 compression.codec) send({ compression: CompressionTypes.GZIP })

Python — confluent-kafka

디렉터리 구조
python-client/
├── requirements.txt
├── producer.py
└── consumer.py
python-client/requirements.txt
# librdkafka 를 번들한 C 확장 바인딩.
# 대부분의 플랫폼에 wheel 이 제공되므로 컴파일러가 필요하지 않습니다.
confluent-kafka==2.15.0

프로듀서

python-client/producer.py
#!/usr/bin/env python3
"""주문 이벤트 프로듀서 (confluent-kafka 2.15.0).

Java 프로듀서와 다른 두 가지:
  1) 전달 콜백은 poll() 또는 flush() 를 호출할 때 실행됩니다.
     produce() 만 반복하면 콜백이 전혀 실행되지 않고 큐만 쌓입니다.
  2) 종료 전에 flush() 를 반드시 호출해야 합니다.
     호출하지 않으면 프로세스가 끝날 때 큐의 레코드가 사라집니다.
"""

import json
import logging
import signal
import sys
import uuid
from typing import Optional

from confluent_kafka import KafkaError, KafkaException, Message, Producer

logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s %(levelname)-5s %(name)s - %(message)s",
)
log = logging.getLogger("producer")

# 예제 1 클러스터의 호스트용 리스너 주소
BOOTSTRAP = "localhost:29092,localhost:39092,localhost:49092"
TOPIC = "orders"

# 전달 실패 건수. 콜백이 다른 스레드가 아니라 poll()/flush() 호출 스레드에서
# 실행되므로 락 없이 카운트해도 안전합니다.
stats = {"delivered": 0, "failed": 0}


def build_producer() -> Producer:
    """무손실 프로듀서 설정.

    설정 키는 librdkafka 이름을 씁니다.
    대부분 Java 와 같지만, 확실하지 않은 값은 문자열로 넘기면 안전합니다.
    """
    conf = {
        "bootstrap.servers": BOOTSTRAP,
        # 브로커 로그와 메트릭에서 이 애플리케이션을 식별합니다.
        "client.id": "orders-producer-py",

        # --- 내구성 -------------------------------------------------------
        # ISR 전부의 기록을 기다립니다. 토픽의 min.insync.replicas=2 와 짝을 이룹니다.
        "acks": "all",
        # 재시도로 인한 중복을 브로커가 제거합니다.
        # 켜면 acks=all, retries 무한, max.in.flight<=5 가 강제됩니다.
        "enable.idempotence": True,
        # 멱등성이 켜지면 5 이하여야 합니다. 6 이상은 설정 오류입니다.
        "max.in.flight.requests.per.connection": 5,

        # --- 시간 예산 ----------------------------------------------------
        # produce() 부터 성공/실패 확정까지의 총 상한(ms).
        # librdkafka 에서는 message.timeout.ms 가 같은 의미의 별칭입니다.
        "delivery.timeout.ms": 120000,
        "request.timeout.ms": 30000,

        # --- 처리량 -------------------------------------------------------
        # 배치를 모으는 시간. Java 4.x 기본값도 5 입니다.
        "linger.ms": 20,
        "batch.size": 65536,
        # 압축. lz4 는 CPU 대비 압축률 균형이 좋습니다.
        "compression.type": "lz4",

        # --- 큐 상한 ------------------------------------------------------
        # librdkafka 는 내부 큐가 가득 차면 produce() 가
        # BufferError 를 발생시킵니다(Java 의 max.block.ms 블로킹과 다릅니다).
        # 이 두 값이 큐의 상한입니다.
        "queue.buffering.max.messages": 100000,
        "queue.buffering.max.kbytes": 1048576,
    }
    return Producer(conf)


def on_delivery(err: Optional[KafkaError], msg: Message) -> None:
    """전달 결과 콜백.

    poll() 또는 flush() 를 호출하는 스레드에서 실행됩니다.
    여기서 블로킹하면 전송이 함께 느려집니다.
    """
    if err is not None:
        stats["failed"] += 1
        # err.retriable() 로 재시도 가능성을 판단할 수 있습니다.
        # 여기까지 온 retriable 오류는 delivery.timeout.ms 를 소진한 것입니다.
        log.error(
            "전달 실패 key=%s code=%s retriable=%s reason=%s",
            msg.key(), err.code(), err.retriable(), err.str(),
        )
        # 실제 서비스에서는 여기서 격리(outbox/디스크 큐)해야 합니다.
        # 아무것도 하지 않으면 그 레코드는 사라집니다.
        return

    stats["delivered"] += 1
    log.debug(
        "전달 성공 key=%s %s-%s@%s",
        msg.key(), msg.topic(), msg.partition(), msg.offset(),
    )


def main(count: int) -> int:
    producer = build_producer()

    # SIGTERM 을 받아도 flush 하도록 플래그를 둡니다.
    stopping = {"flag": False}

    def handle_signal(signum, _frame):
        log.info("시그널 %s 수신 — 남은 큐를 비웁니다", signum)
        stopping["flag"] = True

    signal.signal(signal.SIGTERM, handle_signal)
    signal.signal(signal.SIGINT, handle_signal)

    for i in range(count):
        if stopping["flag"]:
            break

        order_id = f"ORD-{uuid.uuid4()}"
        payload = json.dumps(
            {
                "orderId": order_id,
                "amount": 10000 + i,
                "currency": "KRW",
                "status": "APPROVED",
            },
            ensure_ascii=False,
        )

        try:
            producer.produce(
                topic=TOPIC,
                # 키를 주면 같은 키가 같은 파티션으로 가 순서가 보장됩니다.
                key=order_id.encode("utf-8"),
                value=payload.encode("utf-8"),
                # 헤더는 (이름, 바이트) 튜플 리스트입니다.
                headers=[("event-type", b"OrderApproved"), ("schema-version", b"1")],
                on_delivery=on_delivery,
            )
        except BufferError:
            # 내부 큐가 가득 찼습니다. Java 와 달리 예외로 알려 줍니다.
            # poll() 로 큐를 비우고 다시 시도합니다.
            log.warning("내부 큐 포화 — poll() 로 비우고 재시도합니다")
            producer.poll(1.0)
            producer.produce(
                topic=TOPIC,
                key=order_id.encode("utf-8"),
                value=payload.encode("utf-8"),
                on_delivery=on_delivery,
            )
        except KafkaException as exc:
            # 직렬화·설정 오류 등. 재시도해도 같은 결과입니다.
            log.error("produce 실패 — 재시도 불가 key=%s: %s", order_id, exc)
            stats["failed"] += 1

        # ★ 반드시 주기적으로 poll() 을 호출해야 콜백이 실행됩니다.
        #   0 을 넘기면 블로킹 없이 처리 대기 중인 이벤트만 소화합니다.
        producer.poll(0)

    # ★ 종료 전 flush(). 남은 모든 레코드가 확정될 때까지 블록하고,
    #   반환값은 "아직 전달되지 않은 레코드 수" 입니다.
    remaining = producer.flush(timeout=130)
    if remaining > 0:
        log.error("flush 후에도 %d건이 남았습니다 — 유실 위험", remaining)

    log.info("종료. delivered=%d failed=%d", stats["delivered"], stats["failed"])
    return 1 if stats["failed"] > 0 or remaining > 0 else 0


if __name__ == "__main__":
    n = int(sys.argv[1]) if len(sys.argv) > 1 else 1000
    sys.exit(main(n))

컨슈머

python-client/consumer.py
#!/usr/bin/env python3
"""주문 이벤트 컨슈머 (confluent-kafka 2.15.0).

수동 커밋으로 at-least-once 를 구현합니다.
Java 컨슈머와 다른 점:
  - poll() 이 레코드 하나(Message)를 반환합니다. consume(num_messages=N) 로
    배치를 받을 수도 있습니다.
  - 에러는 예외가 아니라 msg.error() 로 전달되는 경우가 있습니다.
    PARTITION_EOF 는 에러가 아니라 정보성 이벤트입니다.
"""

import json
import logging
import signal
import sys
from typing import List

from confluent_kafka import (
    Consumer,
    KafkaError,
    KafkaException,
    Message,
    TopicPartition,
)

logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s %(levelname)-5s %(name)s - %(message)s",
)
log = logging.getLogger("consumer")

BOOTSTRAP = "localhost:29092,localhost:39092,localhost:49092"
TOPIC = "orders"
GROUP_ID = "orders-etl-py"
BATCH_SIZE = 200

running = {"flag": True}


def build_consumer() -> Consumer:
    conf = {
        "bootstrap.servers": BOOTSTRAP,
        "group.id": GROUP_ID,
        "client.id": "orders-etl-py-1",

        # --- 커밋 ---------------------------------------------------------
        # 자동 커밋을 끕니다. 기본값은 True 이며,
        # 그대로 두면 처리 완료와 무관하게 오프셋이 전진해 유실 창이 생깁니다.
        "enable.auto.commit": False,

        # 커밋된 오프셋이 없을 때만 적용됩니다. 기본값은 latest 입니다.
        "auto.offset.reset": "earliest",

        # --- 배치와 시간 예산 ---------------------------------------------
        # 한 번에 처리할 최대 건수(아래 consume() 에서 씁니다).
        "max.poll.interval.ms": 300000,
        # fetch 튜닝
        "fetch.min.bytes": 65536,
        "fetch.wait.max.ms": 500,

        # 상류가 트랜잭션 프로듀서라면 read_committed 로 바꿉니다.
        # 기본값은 read_uncommitted 입니다.
        "isolation.level": "read_uncommitted",

        # 파티션 끝에 도달하면 PARTITION_EOF 이벤트를 받습니다.
        # 배치 처리에서 "지금 읽을 것이 더 없다" 를 아는 데 유용합니다.
        "enable.partition.eof": True,
    }
    return Consumer(conf)


def on_assign(consumer: Consumer, partitions: List[TopicPartition]) -> None:
    """파티션 할당 콜백."""
    log.info("파티션 할당: %s", [(p.topic, p.partition) for p in partitions])


def on_revoke(consumer: Consumer, partitions: List[TopicPartition]) -> None:
    """파티션 회수 직전 콜백.

    여기서 진행 중인 배치를 처리하고 동기 커밋해야
    다른 인스턴스가 이어받을 때 중복이 최소화됩니다.
    """
    log.info("파티션 회수 예정: %s", [(p.topic, p.partition) for p in partitions])
    try:
        # asynchronous=False → 동기 커밋. 회수 직전에는 반드시 동기여야 합니다.
        consumer.commit(asynchronous=False)
    except KafkaException as exc:
        # 이미 소유권을 잃었으면 실패합니다. 중복 처리가 발생할 뿐 유실은 아닙니다.
        log.warning("회수 전 커밋 실패 — 중복 처리 가능: %s", exc)


def process_batch(batch: List[Message]) -> None:
    """배치 적재. 실패하면 예외를 던져 커밋을 막습니다."""
    rows = []
    for msg in batch:
        try:
            rows.append(json.loads(msg.value().decode("utf-8")))
        except (UnicodeDecodeError, json.JSONDecodeError) as exc:
            # 데이터 자체가 잘못된 경우입니다. 재시도해도 같으므로
            # 실제 서비스라면 DLQ 로 보내고 계속 진행합니다.
            log.error(
                "파싱 실패 → 건너뜁니다 %s-%s@%s: %s",
                msg.topic(), msg.partition(), msg.offset(), exc,
            )

    # 여기서 실제 적재를 수행합니다. 하나의 트랜잭션으로 묶어야 합니다.
    log.info("적재 완료 %d건", len(rows))


def main() -> int:
    consumer = build_consumer()
    consumer.subscribe([TOPIC], on_assign=on_assign, on_revoke=on_revoke)

    def handle_signal(signum, _frame):
        log.info("시그널 %s 수신 — 종료합니다", signum)
        running["flag"] = False

    signal.signal(signal.SIGTERM, handle_signal)
    signal.signal(signal.SIGINT, handle_signal)

    total = 0
    try:
        while running["flag"]:
            # consume() 은 최대 num_messages 개를 timeout 안에 모아 반환합니다.
            # poll() 은 하나씩 반환합니다 — 배치 처리에는 consume() 이 편합니다.
            messages = consumer.consume(num_messages=BATCH_SIZE, timeout=1.0)
            if not messages:
                continue

            batch: List[Message] = []
            for msg in messages:
                err = msg.error()
                if err is None:
                    batch.append(msg)
                    continue

                if err.code() == KafkaError._PARTITION_EOF:
                    # 에러가 아니라 "이 파티션의 끝에 도달했다" 는 정보입니다.
                    log.debug("파티션 끝 도달 %s-%s", msg.topic(), msg.partition())
                    continue

                if err.retriable():
                    log.warning("일시적 오류 — 계속 진행합니다: %s", err.str())
                    continue

                # 복구 불가 오류입니다.
                raise KafkaException(err)

            if not batch:
                continue

            # ★ 처리 후 커밋. 순서를 바꾸면 유실이 됩니다.
            process_batch(batch)
            total += len(batch)

            # 인자 없는 commit() 은 "지금까지 consume 한 위치" 를 커밋합니다.
            # 특정 위치를 커밋하려면 offsets=[TopicPartition(...)] 를 넘깁니다.
            # 커밋 오프셋은 "다음에 읽을 위치" 이므로 마지막 오프셋 + 1 입니다.
            consumer.commit(asynchronous=False)

    except KafkaException as exc:
        log.error("컨슈머 오류로 종료합니다: %s", exc)
        return 1
    finally:
        # close() 는 마지막 오프셋 커밋(자동 커밋일 때)과
        # 그룹 이탈을 처리합니다. 생략하면 리밸런스가
        # session.timeout.ms 만큼 늦게 일어납니다.
        consumer.close()

    log.info("종료. 처리 %d건", total)
    return 0


if __name__ == "__main__":
    sys.exit(main())
실행
cd python-client
python3 -m venv .venv && source .venv/bin/activate
pip install -r requirements.txt

# 프로듀서
python producer.py 1000

# 컨슈머 (다른 터미널)
python consumer.py

Node.js — kafkajs

디렉터리 구조
node-client/
├── package.json
├── producer.mjs
└── consumer.mjs
node-client/package.json
{
  "name": "kafka-node-client",
  "version": "1.0.0",
  "private": true,
  "type": "module",
  "engines": {
    "node": ">=20"
  },
  "scripts": {
    "produce": "node producer.mjs",
    "consume": "node consumer.mjs"
  },
  "dependencies": {
    "kafkajs": "2.2.4"
  }
}

프로듀서

node-client/producer.mjs
// 주문 이벤트 프로듀서 (kafkajs 2.2.4).
//
// Java/librdkafka 와 다른 점:
//   - 설정이 camelCase 옵션 객체입니다.
//   - acks 를 문자열이 아니라 숫자로 지정합니다 (-1 = all).
//   - linger.ms 에 대응하는 설정이 없습니다. 배치는 sendBatch 로 직접 묶습니다.
//   - send() 가 Promise 를 반환하므로 await 하면 그것이 동기 전송입니다.

import { Kafka, CompressionTypes, logLevel } from 'kafkajs';
import { randomUUID } from 'node:crypto';

const BOOTSTRAP = ['localhost:29092', 'localhost:39092', 'localhost:49092'];
const TOPIC = 'orders';

const kafka = new Kafka({
  // 브로커 로그와 메트릭에서 이 애플리케이션을 식별합니다.
  clientId: 'orders-producer-node',
  // 초기 연결 이중화. 부하 분산 목적이 아닙니다.
  brokers: BOOTSTRAP,
  logLevel: logLevel.INFO,
  // 연결 재시도 정책. 지수 백오프로 최대 5회 시도합니다.
  retry: {
    initialRetryTime: 300,
    retries: 5,
    maxRetryTime: 30_000,
  },
  // 요청 타임아웃(ms). Java 의 request.timeout.ms 에 대응합니다.
  requestTimeout: 30_000,
  // 연결 타임아웃
  connectionTimeout: 10_000,
});

const producer = kafka.producer({
  // 멱등 프로듀서. 켜면 kafkajs 가 acks=-1 과 maxInFlightRequests<=5 를 강제합니다.
  // 명시적으로 acks 를 0 이나 1 로 주면 오류가 발생합니다.
  idempotent: true,
  // 멱등성이 켜진 상태의 상한입니다.
  maxInFlightRequests: 5,
  // 메타데이터를 얻지 못했을 때 토픽 자동 생성을 시도하지 않습니다.
  // 브로커에서 auto.create.topics.enable 을 껐다면 이 값도 false 여야
  // 오류가 조기에 드러납니다.
  allowAutoTopicCreation: false,
});

let delivered = 0;
let failed = 0;

/** 레코드 한 건을 만듭니다. */
function buildMessage(index) {
  const orderId = `ORD-${randomUUID()}`;
  return {
    // 키를 주면 같은 키가 같은 파티션으로 가 순서가 보장됩니다.
    key: orderId,
    value: JSON.stringify({
      orderId,
      amount: 10_000 + index,
      currency: 'KRW',
      status: 'APPROVED',
    }),
    // 헤더 값은 문자열 또는 Buffer 입니다.
    headers: {
      'event-type': 'OrderApproved',
      'schema-version': '1',
    },
  };
}

/**
 * 배치 전송.
 *
 * kafkajs 에는 linger.ms 가 없으므로 애플리케이션이 직접 묶습니다.
 * 한 건씩 await send() 하면 매번 브로커 왕복을 기다려 처리량이 급락합니다.
 */
async function sendBatch(messages) {
  try {
    const result = await producer.send({
      topic: TOPIC,
      // -1 = all (ISR 전부). 0 이나 1 을 주면 유실 가능성이 생기고,
      // idempotent: true 와 함께 쓰면 오류가 됩니다.
      acks: -1,
      // 전송 타임아웃(ms). Java 의 delivery.timeout.ms 에 가까운 역할입니다.
      timeout: 30_000,
      // 압축. GZIP 은 내장이고 LZ4/SNAPPY/ZSTD 는 별도 패키지가 필요합니다.
      compression: CompressionTypes.GZIP,
      messages,
    });

    delivered += messages.length;
    // result 는 파티션별 baseOffset 등을 담은 배열입니다.
    for (const r of result) {
      console.debug(
        `전송 성공 partition=${r.partition} baseOffset=${r.baseOffset} count=${messages.length}`,
      );
    }
  } catch (err) {
    failed += messages.length;
    // 여기서 아무것도 하지 않으면 이 배치는 사라집니다.
    // 실제 서비스에서는 격리(outbox/디스크 큐)가 필요합니다.
    console.error(`전송 실패 count=${messages.length}`, err);
    throw err;
  }
}

async function main() {
  const total = Number.parseInt(process.argv[2] ?? '1000', 10);
  const batchSize = 200;

  await producer.connect();

  // SIGTERM 을 받으면 진행 중인 전송을 마치고 연결을 닫습니다.
  // disconnect() 는 대기 중인 요청이 끝날 때까지 기다립니다.
  let stopping = false;
  const shutdown = async (signal) => {
    if (stopping) return;
    stopping = true;
    console.info(`${signal} 수신 — 연결을 닫습니다`);
    await producer.disconnect();
    console.info(`종료. delivered=${delivered} failed=${failed}`);
    process.exit(failed > 0 ? 1 : 0);
  };
  process.on('SIGTERM', () => void shutdown('SIGTERM'));
  process.on('SIGINT', () => void shutdown('SIGINT'));

  try {
    let buffer = [];
    for (let i = 0; i < total && !stopping; i++) {
      buffer.push(buildMessage(i));
      if (buffer.length >= batchSize) {
        await sendBatch(buffer);
        buffer = [];
      }
    }
    if (buffer.length > 0) {
      await sendBatch(buffer);
    }
  } finally {
    // ★ 반드시 disconnect() 를 호출합니다.
    //   호출하지 않으면 프로세스가 종료되지 않거나 대기 중 요청이 유실됩니다.
    await producer.disconnect();
  }

  console.info(`종료. delivered=${delivered} failed=${failed}`);
  process.exitCode = failed > 0 ? 1 : 0;
}

main().catch((err) => {
  console.error('치명적 오류', err);
  process.exit(1);
});

컨슈머

node-client/consumer.mjs
// 주문 이벤트 컨슈머 (kafkajs 2.2.4).
//
// ★ 가장 중요한 사실: eachMessage/eachBatch 는 기본적으로 자동 커밋합니다.
//   핸들러가 정상 반환하면 kafkajs 가 오프셋을 전진시킵니다.
//   그래서 핸들러 안에서 try-catch 로 예외를 삼키면
//   "실패했는데 커밋된" 상태가 되어 유실이 됩니다.
//   예외를 throw 하면 kafkajs 가 커밋하지 않고 재시도합니다.

import { Kafka, logLevel } from 'kafkajs';

const BOOTSTRAP = ['localhost:29092', 'localhost:39092', 'localhost:49092'];
const TOPIC = 'orders';
const GROUP_ID = 'orders-bff-node';

const kafka = new Kafka({
  clientId: 'orders-bff-node-1',
  brokers: BOOTSTRAP,
  logLevel: logLevel.INFO,
  retry: { initialRetryTime: 300, retries: 5 },
});

const consumer = kafka.consumer({
  groupId: GROUP_ID,
  // 하트비트가 이 시간 동안 끊기면 그룹에서 축출됩니다(ms).
  // Java 의 session.timeout.ms 에 대응합니다.
  sessionTimeout: 45_000,
  // 하트비트 전송 간격. sessionTimeout 의 1/3 이하가 권장입니다.
  heartbeatInterval: 3_000,
  // 리밸런스 타임아웃. 배치 처리 시간이 이 값을 넘으면 축출됩니다.
  // Java 의 max.poll.interval.ms 에 해당하는 역할입니다.
  rebalanceTimeout: 300_000,
  // 존재하지 않는 토픽을 자동 생성하지 않습니다(오타 조기 발견).
  allowAutoTopicCreation: false,
});

let processed = 0;
let failedRecords = 0;

/** 실제 처리 로직. 실패하면 예외를 던집니다. */
function handleOrder(payload) {
  const order = JSON.parse(payload);
  if (!order.orderId) {
    // 데이터 자체가 잘못된 경우입니다. 재시도해도 같습니다.
    throw new Error(`orderId 없음: ${payload.slice(0, 120)}`);
  }
  // 실제로는 여기서 캐시 갱신, 알림 발송 등을 수행합니다.
  processed += 1;
}

async function main() {
  await consumer.connect();

  // fromBeginning 은 "커밋된 오프셋이 없을 때만" 적용됩니다.
  // Java 의 auto.offset.reset=earliest 와 같은 의미이며,
  // 기존 그룹의 재처리 수단이 아닙니다.
  await consumer.subscribe({ topic: TOPIC, fromBeginning: true });

  await consumer.run({
    // autoCommit 기본값은 true 입니다.
    // false 로 두고 직접 commitOffsets 를 호출하면 커밋 시점을 통제할 수 있습니다.
    autoCommit: false,

    // eachBatch 는 파티션 단위 배치를 받습니다.
    // eachMessage 보다 제어가 넓고, 수동 커밋과 함께 쓰기 좋습니다.
    eachBatch: async ({
      batch,
      resolveOffset,
      heartbeat,
      commitOffsetsIfNecessary,
      isRunning,
      isStale,
    }) => {
      for (const message of batch.messages) {
        // 리밸런스가 일어났거나 종료 중이면 즉시 멈춥니다.
        // 이 검사를 빼면 소유권을 잃은 파티션을 계속 처리해 중복이 커집니다.
        if (!isRunning() || isStale()) break;

        try {
          handleOrder(message.value.toString('utf8'));
          // 이 레코드까지 처리했다고 표시합니다.
          // resolveOffset 만으로는 브로커에 커밋되지 않습니다.
          resolveOffset(message.offset);
        } catch (err) {
          failedRecords += 1;
          console.error(
            `처리 실패 ${batch.topic}-${batch.partition}@${message.offset}`,
            err,
          );
          // 데이터 오류는 건너뛰고 진행합니다(실제로는 DLQ 로 보냅니다).
          // 예외를 그대로 던지면 이 배치 전체가 재시도되어
          // 같은 레코드에서 영구히 막힙니다(포이즌 필).
          resolveOffset(message.offset);
        }

        // 긴 배치를 처리하는 동안 하트비트를 보내
        // sessionTimeout 초과로 축출되는 것을 막습니다.
        await heartbeat();
      }

      // resolveOffset 으로 표시한 위치를 실제로 커밋합니다.
      // autoCommit: false 이므로 이 호출이 없으면 오프셋이 전진하지 않습니다.
      await commitOffsetsIfNecessary();
    },
  });

  const shutdown = async (signal) => {
    console.info(`${signal} 수신 — 컨슈머를 정상 종료합니다`);
    // disconnect() 는 진행 중인 배치를 마치고 그룹에서 이탈합니다.
    // 생략하면 sessionTimeout 만큼 리밸런스가 늦어집니다.
    await consumer.disconnect();
    console.info(`종료. processed=${processed} failedRecords=${failedRecords}`);
    process.exit(0);
  };
  process.on('SIGTERM', () => void shutdown('SIGTERM'));
  process.on('SIGINT', () => void shutdown('SIGINT'));
}

main().catch((err) => {
  console.error('치명적 오류', err);
  process.exit(1);
});
실행
cd node-client
npm install

# 프로듀서
npm run produce -- 1000

# 컨슈머 (다른 터미널)
npm run consume

검증 방법

1. 건수가 맞는가

브로커가 커밋한 건수를 확인
cd kafka-lab
./kcli kafka-get-offsets.sh --topic orders --time -1 \
  | awk -F: '{s+=$3} END {print "총 " s "건"}'

Python은 delivered=1000, Node는 delivered=1000이 나오고 위 합계가 그만큼 증가해야 합니다. 애플리케이션 카운터와 브로커 오프셋이 어긋나면 어딘가에서 실패를 삼킨 것입니다.

2. flush()를 빼면 실제로 유실되는가

Python에서 flush()를 제거해 확인
# producer.py 의 flush 호출을 임시로 주석 처리하고 실행합니다.
# 큐에 남은 레코드가 프로세스 종료와 함께 사라집니다.
python producer.py 5000

# 오프셋 증가량이 5000보다 작습니다.
./kcli kafka-get-offsets.sh --topic orders --time -1 \
  | awk -F: '{s+=$3} END {print s}'

linger.ms=20이므로 마지막 배치가 큐에 남은 채 프로세스가 끝납니다. 에러도 로그도 남지 않습니다on_delivery 콜백이 실행될 기회가 없기 때문입니다. 시나리오의 사고가 정확히 이것입니다.

3. 두 언어의 컨슈머 그룹이 독립적인가

그룹 목록 확인
./kcli kafka-consumer-groups.sh --list | grep -E 'py|node'
# → orders-etl-py
#   orders-bff-node

# 각 그룹이 독립적으로 전체를 읽습니다.
./kcli kafka-consumer-groups.sh --describe --group orders-etl-py
./kcli kafka-consumer-groups.sh --describe --group orders-bff-node

서로 다른 group.id이므로 같은 데이터를 각자 전부 읽습니다. 같은 그룹 ID를 쓰면 파티션을 나눠 갖게 되어 각 팀이 데이터의 일부만 보게 됩니다 — 흔한 사고입니다.

4. 키 파티셔닝이 언어 간에 호환되는가

같은 키를 두 언어에서 보내 파티션을 비교
./kcli kafka-console-consumer.sh --topic orders --from-beginning \
  --timeout-ms 10000 \
  --property print.key=true --property print.partition=true \
  --property print.value=false 2>/dev/null \
  | sort | uniq -c | head

프로덕션 고려사항

로컬 예제와 프로덕션의 차이
항목이 예제프로덕션
Python 종료 처리 flush(timeout=130) 컨테이너의 terminationGracePeriodSecondsdelivery.timeout.ms보다 크게. flush()반환값(남은 건수)을 반드시 확인하세요
Node 종료 처리 disconnect() SIGTERM 핸들러에서 disconnect()await합니다. Node는 이벤트 루프에 대기 작업이 남으면 프로세스가 종료되지 않기도 합니다
에러 격리 로그만 양쪽 모두 격리 저장소가 필요합니다. 콜백/catch에서 로그만 남기면 유실입니다(예제 3)
직렬화 JSON 문자열 Avro/Protobuf. confluent-kafkaSerializingProducer와 Schema Registry 클라이언트를 제공합니다. kafkajs는 별도 패키지가 필요합니다(예제 7)
트랜잭션 사용하지 않음 confluent-kafka-pythoninit_transactions·send_offsets_to_transaction을 제공합니다. kafkajs의 지원 범위는 버전 문서를 확인하세요(예제 5)
Python GIL 단일 스레드 librdkafka의 I/O는 C 스레드에서 돌아 GIL 영향을 덜 받지만, 처리 로직은 GIL에 묶입니다. CPU 바운드면 프로세스를 여러 개 띄우세요(파티션 수까지)
Node 단일 스레드 단일 프로세스 이벤트 루프를 막는 동기 작업이 있으면 하트비트가 지연되어 축출됩니다. 무거운 처리는 워커 스레드나 별도 서비스로 분리하세요
보안 PLAINTEXT Python은 security.protocol·sasl.mechanisms·sasl.username/password·ssl.ca.location. Node는 ssl·sasl 옵션 객체입니다
모니터링 없음 Python은 statistics.interval.ms + stats_cb로 통계 JSON을 받습니다. kafkajs는 consumer.on(consumer.events.*) 이벤트를 씁니다. 양쪽 모두 Java의 JMX 메트릭과 이름이 다릅니다(예제 10)

자주 하는 실수

이어서 볼 곳

공식 문서 출처

Kafka 설정 기본값은 Apache Kafka 4.3.1 문서에서, 라이브러리 버전은 PyPI와 npm 레지스트리에서 직접 확인했습니다.