C++에서 Kafka 쓰기: librdkafka 프로듀서·컨슈머 그룹, 오프셋 관리, 리밸런싱과 메시지 유실

들어가며: C++에서 Kafka를 왜 쓰나요?

여러 서비스가 로그를 중앙 DB에 직접 INSERT하면, 트래픽이 늘 때 DB 연결과 쓰기 지연이 그대로 서비스 응답 시간에 얹힙니다. 주문 완료 이벤트를 재고·포인트·알림 서비스에 HTTP로 차례로 호출하면 가장 느린 서비스가 주문 처리 전체를 붙잡고, 그중 하나가 죽으면 주문도 실패합니다. Kafka는 프로듀서가 토픽에 이벤트를 쓰기만 하고, 각 소비 서비스가 자기 속도로 읽어 가게 만들어 이 결합을 끊습니다. 컨슈머 그룹은 어디까지 읽었는지를 오프셋으로 커밋해 두므로 재시작해도 이어서 처리할 수 있고, 파티션 수만큼 컨슈머를 늘려 처리량을 확장할 수 있습니다.

대신 C++에서 librdkafka로 붙일 때는 새로 신경 쓸 문제가 생깁니다. 브로커 주소를 넣었는데 연결이 안 되는 문제, 처리 도중 크래시 후 재시작하면 같은 메시지를 다시 처리하는 문제, 컨슈머를 추가하면 파티션이 재할당(리밸런싱)되면서 처리 중이던 메시지가 중복되는 문제가 대표적입니다. 이 글은 이런 상황을 코드와 설정 수준에서 다룹니다.

flowchart LR
  subgraph Producer["프로듀서 (C++)"]
    P1[앱 로그]
    P2[주문 이벤트]
    P3[센서 데이터]
  end
  subgraph Kafka[Apache Kafka]
    T1[logs 토픽]
    T2[orders 토픽]
    T3[events 토픽]
  end
  subgraph Consumer["컨슈머 (C++)"]
    C1[로그 저장]
    C2[재고 차감]
    C3[실시간 분석]
  end
  P1 --> T1
  P2 --> T2
  P3 --> T3
  T1 --> C1
  T2 --> C2
  T3 --> C3

Kafka 프로듀서-컨슈머 흐름

sequenceDiagram
  participant P as 프로듀서
  participant B as 브로커
  participant C as 컨슈머
  P->>B: produce(topic, key, value)
  B->>B: 파티션에 저장
  B->>P: dr_cb (delivery report)
  C->>B: subscribe(topic)
  B->>C: 파티션 할당
  loop consume
    C->>B: poll() → 메시지 수신
    C->>C: 비즈니스 로직 처리
    C->>B: commit() 오프셋
  end

RabbitMQ vs Kafka 비교

항목RabbitMQKafka
모델큐, Exchange토픽, 파티션
메시지 보존소비 후 삭제 (기본)보존 기간 동안 유지
재처리별도 구현오프셋 이동으로 가능
확장 방식큐·컨슈머 추가파티션 추가, 순차 로그 쓰기로 대량 처리에 유리
C++ 클라이언트rabbitmq-c, AMQP-CPPlibrdkafka

librdkafka 설치와 Docker 브로커

필수 의존성

항목버전비고
C++C++14 이상C++17 권장
librdkafka2.0+vcpkg, Homebrew, 또는 소스 빌드
Apache Kafka2.8+브로커 (Docker 권장)
CMake3.16+FindPackage 지원

Kafka 브로커 실행 (Docker)

# 단일 브로커 (개발용, KRaft 모드). 공식 apache/kafka 이미지는 3.7.0부터 제공되며
# 환경 변수 없이 실행하면 localhost:9092로 접속 가능한 단일 노드 설정으로 뜬다
docker run -d --name kafka -p 9092:9092 apache/kafka:3.7.0
# 토픽 생성 (이 이미지는 스크립트가 /opt/kafka/bin에 있음)
docker exec kafka /opt/kafka/bin/kafka-topics.sh --create --topic app-logs \
  --partitions 3 --replication-factor 1 --bootstrap-server localhost:9092

KAFKA_CFG_* 형식의 환경 변수는 Bitnami 이미지의 규칙이라 공식 이미지에는 적용되지 않습니다. 설정을 바꾸려면 공식 이미지에서는 KAFKA_NODE_ID, KAFKA_LISTENERS처럼 KAFKA_ 접두사를 씁니다.

librdkafka 설치

# vcpkg (포트 이름은 librdkafka)
vcpkg install librdkafka
# macOS (Homebrew)
brew install librdkafka
# Ubuntu/Debian
sudo apt-get install librdkafka-dev

CMakeLists.txt 기본 설정

cmake_minimum_required(VERSION 3.16)
project(kafka_example LANGUAGES CXX)
set(CMAKE_CXX_STANDARD 17)
find_package(RdKafka CONFIG REQUIRED)
add_executable(kafka_producer producer.cpp)
target_link_libraries(kafka_producer PRIVATE RdKafka::rdkafka++)  # C++ API(rdkafkacpp.h)
add_executable(kafka_consumer consumer.cpp)
target_link_libraries(kafka_consumer PRIVATE RdKafka::rdkafka++)

이 글의 예제는 C++ API(rdkafkacpp.h)를 쓰므로 C 라이브러리 타깃(RdKafka::rdkafka)이 아니라 RdKafka::rdkafka++에 링크해야 합니다. Ubuntu 패키지로 설치했다면 CMake 설정 파일이 없을 수 있어 pkg-config의 rdkafka++ 모듈을 쓰는 편이 간단합니다. 주의: vcpkg 사용 시 -DCMAKE_TOOLCHAIN_FILE=[vcpkg root]/scripts/buildsystems/vcpkg.cmake를 CMake에 전달해야 합니다.


Delivery Report와 RAII 래퍼를 갖춘 프로듀서

최소 동작 프로듀서

// producer.cpp
// 컴파일: g++ -std=c++17 -o producer producer.cpp -lrdkafka++ -lrdkafka
// (rdkafka++가 rdkafka에 의존하므로 링크 순서는 rdkafka++ 먼저)
#include <librdkafka/rdkafkacpp.h>
#include <iostream>
#include <memory>
#include <string>
int main(int argc, char** argv) {
    std::string brokers = argc > 1 ? argv[1] : "localhost:9092";
    std::string topic = argc > 2 ? argv[2] : "app-logs";
    std::string errstr;
    // 1. 전역 설정
    RdKafka::Conf* conf = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
    conf->set("bootstrap.servers", brokers, errstr);
    // 2. 프로듀서 생성
    RdKafka::Producer* producer = RdKafka::Producer::create(conf, errstr);
    delete conf;
    if (!producer) {
        std::cerr << "프로듀서 생성 실패: " << errstr << std::endl;
        return 1;
    }
    // 3. 메시지 발행
    std::string payload = "Hello, Kafka from C++!";
    RdKafka::ErrorCode err = producer->produce(
        topic,
        RdKafka::Topic::PARTITION_UA,  // 파티션 자동 할당
        RdKafka::Producer::RK_MSG_COPY,
        const_cast<char*>(payload.data()),
        payload.size(),
        "key1", 4,  // 키와 길이 (NUL 문자는 포함하지 않음)
        0, nullptr);
    if (err != RdKafka::ERR_NO_ERROR) {
        std::cerr << "발행 실패: " << RdKafka::err2str(err) << std::endl;
    } else {
        std::cout << "메시지 발행 완료" << std::endl;
    }
    // 4. 대기 중인 delivery report 처리
    while (producer->outq_len() > 0) {
        producer->poll(100);
    }
    delete producer;
    RdKafka::wait_destroyed(5000);
    return 0;
}

PARTITION_UA를 주면 파티셔너가 파티션을 고릅니다. 키가 있으면 키 해시로, 없으면 버전에 따라 무작위 또는 sticky 방식으로 정합니다. RK_MSG_COPY는 librdkafka가 페이로드를 내부로 복사하게 하므로 호출 직후 버퍼를 바꾸거나 해제해도 됩니다. produce()가 성공을 반환한 것은 내부 큐에 들어갔다는 뜻일 뿐이므로, outq_len()이 0이 될 때까지 poll()을 호출하거나 flush()로 실제 전송을 기다려야 합니다.

Delivery Report 콜백 (실전용)

발행 결과를 콜백으로 받아 로깅·재시도할 수 있습니다.

// producer_with_dr_cb.cpp
#include <librdkafka/rdkafkacpp.h>
#include <iostream>
#include <memory>
#include <string>
class DeliveryReportCb : public RdKafka::DeliveryReportCb {
public:
    void dr_cb(RdKafka::Message& msg) override {
        if (msg.err()) {
            std::cerr << "[DR] 전달 실패: " << msg.errstr() << std::endl;
        } else {
            std::cout << "[DR] 전달 완료: " << msg.topic_name() << "["
                      << msg.partition() << "] @ " << msg.offset() << std::endl;
        }
    }
};
int main() {
    std::string errstr;
    RdKafka::Conf* conf = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
    conf->set("bootstrap.servers", "localhost:9092", errstr);
    DeliveryReportCb dr_cb;
    conf->set("dr_cb", &dr_cb, errstr);
    RdKafka::Producer* producer = RdKafka::Producer::create(conf, errstr);
    delete conf;
    if (!producer) {
        std::cerr << "프로듀서 생성 실패: " << errstr << std::endl;
        return 1;
    }
    std::string payload = "Test message with delivery report";
    producer->produce("app-logs", RdKafka::Topic::PARTITION_UA,
                      RdKafka::Producer::RK_MSG_COPY,
                      const_cast<char*>(payload.data()), payload.size(),
                      nullptr, 0, 0, nullptr);
    // poll을 호출해야 dr_cb가 실행됨
    for (int i = 0; i < 10 && producer->outq_len() > 0; ++i) {
        producer->poll(100);
    }
    producer->flush(5000);
    delete producer;
    RdKafka::wait_destroyed(5000);
    return 0;
}

주의: poll()을 주기적으로 호출하지 않으면 delivery report 콜백이 실행되지 않습니다. 프로듀서 스레드에서 루프로 poll(100)을 호출하는 것이 일반적입니다.

RAII 래퍼 클래스 (재사용 가능)

// kafka_producer.hpp
#pragma once
#include <librdkafka/rdkafkacpp.h>
#include <memory>
#include <stdexcept>
#include <string>
class KafkaProducer {
public:
    KafkaProducer(const std::string& brokers,
                  const std::string& client_id = "cpp-producer") {
        std::string errstr;
        conf_.reset(RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL));
        conf_->set("bootstrap.servers", brokers, errstr);
        conf_->set("client.id", client_id, errstr);
        conf_->set("dr_cb", &dr_cb_, errstr);
        producer_.reset(RdKafka::Producer::create(conf_.get(), errstr));
        if (!producer_) {
            throw std::runtime_error("프로듀서 생성 실패: " + errstr);
        }
    }
    void produce(const std::string& topic, const std::string& key,
                 const std::string& value) {
        RdKafka::ErrorCode err = producer_->produce(
            topic, RdKafka::Topic::PARTITION_UA,
            RdKafka::Producer::RK_MSG_COPY,
            const_cast<char*>(value.data()), value.size(),
            key.empty() ? nullptr : key.data(), key.size(),
            0, nullptr);
        if (err != RdKafka::ERR_NO_ERROR) {
            throw std::runtime_error("발행 실패: " + RdKafka::err2str(err));
        }
        producer_->poll(0);
    }
    void flush(int timeout_ms = 10000) {
        producer_->flush(timeout_ms);
    }
    void poll(int timeout_ms = 0) {
        producer_->poll(timeout_ms);
    }
private:
    struct DeliveryReportCb : public RdKafka::DeliveryReportCb {
        void dr_cb(RdKafka::Message& msg) override {
            if (msg.err()) {
                std::cerr << "[DR] 실패: " << msg.errstr() << std::endl;
            }
        }
    } dr_cb_;
    std::unique_ptr<RdKafka::Conf> conf_;
    std::unique_ptr<RdKafka::Producer> producer_;
};

리밸런싱·이벤트 콜백을 갖춘 컨슈머

최소 동작 컨슈머

// consumer.cpp
#include <librdkafka/rdkafkacpp.h>
#include <iostream>
#include <memory>
#include <string>
#include <vector>
#include <csignal>
static volatile sig_atomic_t run = 1;
void sigterm_handler(int) { run = 0; }
int main(int argc, char** argv) {
    std::string brokers = argc > 1 ? argv[1] : "localhost:9092";
    std::string group_id = argc > 2 ? argv[2] : "cpp-consumer-group";
    std::string topic = argc > 3 ? argv[3] : "app-logs";
    std::string errstr;
    signal(SIGINT, sigterm_handler);
    signal(SIGTERM, sigterm_handler);
    RdKafka::Conf* conf = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
    conf->set("bootstrap.servers", brokers, errstr);
    conf->set("group.id", group_id, errstr);
    conf->set("enable.auto.commit", "false", errstr);  // 수동 커밋
    RdKafka::KafkaConsumer* consumer =
        RdKafka::KafkaConsumer::create(conf, errstr);
    delete conf;
    if (!consumer) {
        std::cerr << "컨슈머 생성 실패: " << errstr << std::endl;
        return 1;
    }
    std::vector<std::string> topics = {topic};
    RdKafka::ErrorCode err = consumer->subscribe(topics);
    if (err) {
        std::cerr << "구독 실패: " << RdKafka::err2str(err) << std::endl;
        delete consumer;
        return 1;
    }
    std::cout << "메시지 수신 대기 중... (Ctrl+C로 종료)" << std::endl;
    while (run) {
        RdKafka::Message* msg = consumer->consume(1000);
        switch (msg->err()) {
            case RdKafka::ERR_NO_ERROR:
                std::cout << "[" << msg->topic_name() << ":" << msg->partition()
                          << "@" << msg->offset() << "] "
                          << std::string(static_cast<const char*>(msg->payload()),
                                         msg->len())
                          << std::endl;
                consumer->commitSync(msg);  // msg->offset() + 1을 커밋
                break;
            case RdKafka::ERR__TIMED_OUT:
                break;
            case RdKafka::ERR__PARTITION_EOF:  // enable.partition.eof=true일 때만 발생
                break;
            default:
                std::cerr << "소비 에러: " << msg->errstr() << std::endl;
                run = 0;
                break;
        }
        delete msg;
    }
    consumer->close();
    delete consumer;
    RdKafka::wait_destroyed(5000);
    return 0;
}

enable.auto.commit=false로 두면 메시지를 처리한 뒤에 직접 커밋하므로, 처리 중 크래시가 나도 커밋되지 않은 메시지는 재시작 후 다시 읽힙니다. consume(1000)은 최대 1초 기다리고, 메시지가 없으면 ERR__TIMED_OUT을 담은 메시지 객체를 돌려줍니다. C++ API의 KafkaConsumer에는 commit()이라는 메서드가 없고 commitSync(msg)·commitAsync(msg)를 씁니다. 둘 다 “다음에 읽을 위치”인 msg->offset() + 1을 커밋합니다. 메시지마다 commitSync를 부르면 브로커 왕복이 메시지 수만큼 생기므로, 처리량이 중요하면 일정 개수마다 커밋하거나 commitAsync를 씁니다.

리밸런싱 콜백 (컨슈머 그룹)

파티션 재할당 시 리밸런싱 콜백에서 assign/unassign를 처리해야 합니다.

// consumer_rebalance.cpp
#include <librdkafka/rdkafkacpp.h>
#include <iostream>
#include <string>
#include <vector>
#include <csignal>
static volatile sig_atomic_t run = 1;
void sig_handler(int) { run = 0; }
class RebalanceCb : public RdKafka::RebalanceCb {
public:
    void rebalance_cb(RdKafka::KafkaConsumer* consumer,
                     RdKafka::ErrorCode err,
                     std::vector<RdKafka::TopicPartition*>& partitions) override {
        if (err == RdKafka::ERR__ASSIGN_PARTITIONS) {
            consumer->assign(partitions);
            std::cerr << "[리밸런싱] 파티션 할당: ";
            for (auto* p : partitions)
                std::cerr << p->topic() << "[" << p->partition() << "] ";
            std::cerr << std::endl;
        } else {
            consumer->unassign();
            std::cerr << "[리밸런싱] 파티션 해제" << std::endl;
        }
    }
};
int main() {
    std::string errstr;
    RdKafka::Conf* conf = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
    conf->set("bootstrap.servers", "localhost:9092", errstr);
    conf->set("group.id", "cpp-group", errstr);
    conf->set("enable.auto.commit", "false", errstr);
    RebalanceCb rebalance_cb;
    conf->set("rebalance_cb", &rebalance_cb, errstr);
    RdKafka::KafkaConsumer* consumer =
        RdKafka::KafkaConsumer::create(conf, errstr);
    delete conf;
    if (!consumer) {
        std::cerr << "컨슈머 생성 실패: " << errstr << std::endl;
        return 1;
    }
    signal(SIGINT, sig_handler);
    signal(SIGTERM, sig_handler);
    consumer->subscribe({"app-logs"});
    while (run) {
        RdKafka::Message* msg = consumer->consume(1000);
        if (msg->err() == RdKafka::ERR_NO_ERROR) {
            std::cout << std::string(static_cast<const char*>(msg->payload()),
                                     msg->len()) << std::endl;
            consumer->commitSync(msg);
        }
        delete msg;
    }
    consumer->close();
    delete consumer;
    RdKafka::wait_destroyed(5000);
    return 0;
}

이벤트 콜백 (에러·통계)

// EventCb: 연결 에러, throttle 등 수신
class EventCb : public RdKafka::EventCb {
public:
    void event_cb(RdKafka::Event& event) override {
        switch (event.type()) {
            case RdKafka::Event::EVENT_ERROR:
                if (event.fatal()) {
                    std::cerr << "치명적 에러: " << event.str() << std::endl;
                } else {
                    std::cerr << "에러: " << event.str() << std::endl;
                }
                break;
            case RdKafka::Event::EVENT_THROTTLE:
                std::cerr << "Throttled: " << event.throttle_time() << "ms by "
                          << event.broker_name() << std::endl;
                break;
            default:
                break;
        }
    }
};
// 설정 시
EventCb event_cb;
conf->set("event_cb", &event_cb, errstr);

세 콜백이 실제로 실행되는 시점과 수명

앞에서 DeliveryReportCb, RebalanceCb, EventCb를 각각 예제로 봤습니다. 코드 모양은 비슷하지만 운영에서 문제가 되는 지점은 코드가 아니라 언제, 어느 스레드에서 불리느냐와 콜백 객체가 얼마나 오래 살아야 하느냐입니다.

콜백은 어느 스레드에서 실행될까요? librdkafka는 내부 브로커 스레드에서 네트워크 I/O를 처리하지만, 등록한 콜백은 그 스레드에서 바로 부르지 않습니다. 결과를 큐에 쌓아 두었다가 애플리케이션이 프로듀서의 poll()/flush(), 컨슈머의 consume()을 호출할 때 그 호출 스레드에서 실행합니다. 그래서 콜백 안에서는 별도 락 없이 호출 스레드의 상태를 만질 수 있는 대신, 콜백이 오래 걸리면 poll()/consume() 루프 전체가 멈춥니다. 컨슈머에서 콜백이나 메시지 처리가 max.poll.interval.ms(기본 300초)를 넘기면 이 컨슈머는 그룹에서 빠지고 리밸런싱이 시작됩니다. 무거운 처리는 작업 큐로 넘기고 콜백은 짧게 끝내는 편이 안전합니다.

콜백 객체의 수명도 중요합니다. conf->set("dr_cb", &dr_cb, errstr)는 포인터만 저장합니다. 콜백 객체가 프로듀서·컨슈머보다 먼저 파괴되면 다음 poll()에서 해제된 객체를 호출하게 되고, 대개 종료 직전의 flush()에서 원인 모를 크래시로 나타납니다. 위 RAII 래퍼에서 dr_cb_를 producer_보다 먼저 선언한 것도 이 때문입니다. 멤버는 선언 역순으로 파괴되므로 프로듀서가 먼저 사라집니다. 같은 이유로, 래퍼 소멸자에서 flush()를 부르지 않으면 큐에 남은 메시지는 전달 보고도 없이 버려진다는 점도 기억해 둬야 합니다.

리밸런싱 프로토콜에 따라서도 콜백 구현이 달라집니다. 위 RebalanceCb 예제는 eager 방식(range, roundrobin) 기준입니다. partition.assignment.strategy=cooperative-sticky로 바꾸면 리밸런싱 때 전체 파티션이 아니라 바뀐 파티션만 넘어오므로 assign()/unassign() 대신 incremental_assign()/incremental_unassign()을 써야 합니다. 콜백 안에서 consumer->rebalance_protocol()이 "COOPERATIVE"인지 확인해 분기하는 것이 librdkafka 예제의 방식입니다. 이 분기를 빼먹으면 cooperative 전략에서 할당 상태가 어긋나 에러가 나거나 일부 파티션을 소비하지 않게 됩니다. 또 수동 커밋을 쓴다면 ERR__REVOKE_PARTITIONS 분기에서 처리를 끝낸 오프셋을 먼저 커밋하고 해제해야, 파티션을 넘겨받은 다른 컨슈머가 같은 메시지를 다시 처리하는 일이 줄어듭니다.

EventCb의 EVENT_STATS는 statistics.interval.ms를 0보다 크게 설정해야 발생합니다(기본 0 = 꺼짐). 켜면 event.str()로 브로커별 RTT, 큐 길이, 재전송 횟수가 담긴 큰 JSON이 주기적으로 들어옵니다. 이 JSON을 매번 통째로 로그에 남기면 그 자체가 부하가 되므로, 필요한 몇 개 필드만 뽑아 메트릭으로 내보내는 식으로 쓰는 것이 보통입니다.

dr_cb에 에러가 담겨 오는 것은 librdkafka가 message.timeout.ms(기본 300초) 안에서 내부 재시도를 모두 소진한 뒤입니다. 그래서 콜백에서 같은 메시지를 곧바로 다시 produce()하면 이미 긴 재시도를 거친 메시지를 또 보내는 셈이고, 멱등성이 꺼져 있다면 브로커에는 실제로 저장됐는데 응답만 늦은 경우 중복이 생깁니다. 실패 메시지는 로컬 파일이나 별도 토픽(dead letter)에 남기고 원인을 조사하는 편이 안전합니다.


트랜잭션·정확히 한 번 전달

전달 시맨틱 비교

시맨틱설명중복유실
at-most-once커밋 후 처리없음가능
at-least-once처리 후 커밋가능없음
exactly-once멱등성·트랜잭션없음없음

exactly-once는 Kafka 안에서(토픽에서 읽고 → 가공해 → 다른 토픽에 쓰고 → 오프셋을 커밋하는 과정이 한 트랜잭션일 때) 성립하는 보장입니다. 처리 중에 외부 DB에 쓰거나 결제 API를 호출하는 부수 효과는 트랜잭션에 포함되지 않으므로, 그 부분은 아래 멱등성 키처럼 소비자 쪽에서 중복을 걸러야 합니다.

프로듀서 멱등성 (enable.idempotence)

conf->set("enable.idempotence", "true", errstr);
// retries, acks, max.in.flight.requests.per.connection가 자동으로 안전한 값으로 설정됨

트랜잭션 프로듀서

여러 토픽에 원자적으로 쓰거나 consume-transform-produce를 원자적으로 처리할 때 사용합니다. C++ API의 트랜잭션 메서드는 ErrorCode가 아니라 RdKafka::Error*를 반환하며, 성공이면 nullptr이고 실패면 호출자가 delete해야 합니다. 실패 객체의 txn_requires_abort()로 중단이 필요한지, is_retriable()로 재시도 가능한지를 판단합니다.

// transactional_producer.cpp
#include <librdkafka/rdkafkacpp.h>
#include <iostream>
#include <string>
int main() {
    std::string errstr;
    RdKafka::Conf* conf = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
    conf->set("bootstrap.servers", "localhost:9092", errstr);
    conf->set("transactional.id", "my-txn-id", errstr);  // 유일값 필수
    RdKafka::Producer* producer = RdKafka::Producer::create(conf, errstr);
    delete conf;
    if (!producer) {
        std::cerr << "프로듀서 생성 실패: " << errstr << std::endl;
        return 1;
    }
    // 1. 트랜잭션 초기화 (반환값은 RdKafka::Error*, 성공 시 nullptr)
    if (RdKafka::Error* e = producer->init_transactions(5000)) {
        std::cerr << "init_transactions 실패: " << e->str() << std::endl;
        delete e;
        delete producer;
        return 1;
    }
    // 2. 트랜잭션 시작
    if (RdKafka::Error* e = producer->begin_transaction()) {
        std::cerr << "begin_transaction 실패: " << e->str() << std::endl;
        delete e;
        delete producer;
        return 1;
    }
    // 3. 여러 토픽에 발행
    std::string msg1 = "order-123";
    std::string msg2 = "event-order-123";
    producer->produce("orders", RdKafka::Topic::PARTITION_UA,
                     RdKafka::Producer::RK_MSG_COPY,
                     const_cast<char*>(msg1.data()), msg1.size(),
                     nullptr, 0, 0, nullptr);
    producer->produce("order-events", RdKafka::Topic::PARTITION_UA,
                     RdKafka::Producer::RK_MSG_COPY,
                     const_cast<char*>(msg2.data()), msg2.size(),
                     nullptr, 0, 0, nullptr);
    // 4. 트랜잭션 커밋 (commit_transaction이 내부적으로 flush까지 수행)
    if (RdKafka::Error* e = producer->commit_transaction(5000)) {
        std::cerr << "commit_transaction 실패: " << e->str() << std::endl;
        if (e->txn_requires_abort()) {
            if (RdKafka::Error* ae = producer->abort_transaction(5000)) delete ae;
        }
        delete e;
    }
    delete producer;
    RdKafka::wait_destroyed(5000);
    return 0;
}

transactional.id는 프로듀서 인스턴스마다 고유해야 합니다. 같은 ID로 새 프로듀서가 init_transactions()를 호출하면 이전 인스턴스는 펜싱(fencing)되어 더 이상 트랜잭션을 커밋할 수 없습니다. 예제에서는 produce()의 반환값 검사를 생략했지만, 실제 코드에서는 실패 시 트랜잭션을 중단해야 합니다.

컨슈머 read_committed

트랜잭션으로 커밋된 메시지만 읽으려면:

conf->set("isolation.level", "read_committed", errstr);

멱등성 키로 중복 제거

// 멱등성 캐시 (Redis/DB 권장)
class IdempotencyCache {
public:
    bool try_acquire(const std::string& key) {
        std::lock_guard<std::mutex> lock(mutex_);
        if (seen_.count(key)) return false;
        seen_.insert(key);
        if (seen_.size() > 100000) seen_.erase(seen_.begin());  // ⚠️ unordered_set의 begin()은 '가장 오래된 키'가 아님
        return true;
    }
private:
    std::unordered_set<std::string> seen_;
    std::mutex mutex_;
};
// 메시지 처리 시
std::string idempotency_key = msg->key() ?
    std::string(msg->key()->data(), msg->key()->size()) :
    std::string(msg->topic_name()) + "-" + std::to_string(msg->partition()) +
    "-" + std::to_string(msg->offset());
if (idempotency_cache.try_acquire(idempotency_key)) {
    do_business_logic(msg);
}
consumer->commitSync(msg);

위 캐시는 개념을 보여 주는 예시입니다. unordered_set의 begin()은 삽입 순서와 무관해서, 크기 제한에 걸렸을 때 방금 넣은 키가 지워질 수도 있고, 프로세스가 재시작하면 기록이 모두 사라집니다. 실제로는 키에 TTL을 둔 Redis SET key 1 NX EX 86400이나, 처리 결과를 저장하는 DB 테이블의 유니크 제약으로 중복을 막는 것이 일반적입니다.


topic-a → topic-b 스트리밍 파이프라인

변환 단계 구성

C++에서는 컨슈머로 읽고 → 변환 → 프로듀서로 쓰는 패턴으로 스트림 처리를 구현합니다.

sequenceDiagram
    participant C as 컨슈머
    participant App as C++ 앱
    participant P as 프로듀서
    C->>App: consume(topic-a)
    App->>App: transform/filter/aggregate
    App->>P: produce(topic-b)
    P-->>App: delivery report
    App->>C: commit (전달 확인 후)

전체 파이프라인 예제

// stream_pipeline.cpp
#include <librdkafka/rdkafkacpp.h>
#include <iostream>
#include <string>
#include <csignal>
static volatile sig_atomic_t run = 1;
void sig_handler(int) { run = 0; }
int main() {
    std::string errstr;
    std::string brokers = "localhost:9092";
    RdKafka::Conf* c_conf = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
    c_conf->set("bootstrap.servers", brokers, errstr);
    c_conf->set("group.id", "stream-pipeline-group", errstr);
    c_conf->set("enable.auto.commit", "false", errstr);
    RdKafka::KafkaConsumer* consumer = RdKafka::KafkaConsumer::create(c_conf, errstr);
    delete c_conf;
    RdKafka::Conf* p_conf = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
    p_conf->set("bootstrap.servers", brokers, errstr);
    RdKafka::Producer* producer = RdKafka::Producer::create(p_conf, errstr);
    delete p_conf;
    consumer->subscribe({"raw-logs"});
    signal(SIGINT, sig_handler);
    signal(SIGTERM, sig_handler);
    while (run) {
        RdKafka::Message* msg = consumer->consume(1000);
        switch (msg->err()) {
            case RdKafka::ERR_NO_ERROR: {
                std::string payload(static_cast<const char*>(msg->payload()),
                                    msg->len());
                std::string transformed;
                for (char c : payload) {
                    transformed += (c >= 'a' && c <= 'z') ? (c - 32) : c;
                }
                RdKafka::ErrorCode err = producer->produce(
                    "processed-logs",
                    RdKafka::Topic::PARTITION_UA,
                    RdKafka::Producer::RK_MSG_COPY,
                    const_cast<char*>(transformed.data()),
                    transformed.size(),
                    msg->key() ? msg->key()->data() : nullptr,
                    msg->key() ? msg->key()->size() : 0,
                    0, nullptr);
                if (err == RdKafka::ERR_NO_ERROR) {
                    // 주의: produce() 성공은 '큐에 들어감'일 뿐 전달 확인이 아님
                    consumer->commitSync(msg);
                }
                break;
            }
            case RdKafka::ERR__TIMED_OUT:
            case RdKafka::ERR__PARTITION_EOF:
                break;
            default:
                run = 0;
                break;
        }
        delete msg;
        producer->poll(0);
    }
    producer->flush(5000);
    consumer->close();
    delete producer;
    delete consumer;
    RdKafka::wait_destroyed(5000);
    return 0;
}

이 예제는 produce()가 성공을 반환하면 원본 오프셋을 커밋합니다. 하지만 produce()의 성공은 메시지가 librdkafka 내부 큐에 들어갔다는 뜻이지 브로커에 저장됐다는 뜻이 아닙니다. 커밋한 뒤에 전송이 최종 실패하면(dr_cb에 에러) 그 메시지는 출력 토픽에도 없고 다시 읽히지도 않아 유실됩니다. 진짜 at-least-once가 필요하면 dr_cb에서 전달이 확인된 메시지의 원본 오프셋만 커밋하거나, 트랜잭션 프로듀서의 send_offsets_to_transaction()으로 출력 메시지와 입력 오프셋을 한 트랜잭션에 묶어야 합니다. 후자가 Kafka 안에서의 exactly-once입니다.

필터링 파이프라인

// 에러 로그만 필터링
bool should_forward(const std::string& payload) {
    return payload.find("ERROR") != std::string::npos ||
           payload.find("FATAL") != std::string::npos;
}
if (msg->err() == RdKafka::ERR_NO_ERROR) {
    std::string payload(static_cast<const char*>(msg->payload()), msg->len());
    if (should_forward(payload)) {
        producer->produce("error-logs", RdKafka::Topic::PARTITION_UA,
                         RdKafka::Producer::RK_MSG_COPY,
                         const_cast<char*>(payload.data()), payload.size(),
                         nullptr, 0, 0, nullptr);
    }
    consumer->commitSync(msg);  // 위 파이프라인과 같은 유실 가능성이 있음
}

Connection refused, 메타데이터 누락, Queue full: 에러 해결

”Connection refused” / “Broker: Connection refused”

Producer::create()는 브로커에 연결하지 않고도 성공하므로, 연결 문제는 발행·소비 시점의 로그나 event_cb 에러로 드러납니다. 브로커가 실행 중인지, 주소와 포트가 맞는지, 방화벽이 막는지부터 확인합니다. Docker에서 특히 흔한 원인은 advertised listener입니다. 클라이언트는 bootstrap.servers로 처음 접속한 뒤 브로커가 알려 주는 advertised.listeners 주소로 다시 접속하므로, 브로커가 컨테이너 내부 호스트명을 광고하면 호스트의 클라이언트는 처음 접속에 성공하고도 그다음 연결에서 실패합니다.

# 브로커 실행 확인
docker ps | grep kafka
# Docker 내부에서 접속 시
# bootstrap.servers = localhost:9092 (호스트에서)
# bootstrap.servers = host.docker.internal:9092 (Docker 컨테이너에서 호스트 접속)
// ✅ 여러 브로커 지정 (고가용성)
conf->set("bootstrap.servers", "broker1:9092,broker2:9092,broker3:9092", errstr);
// ✅ 연결 타임아웃 설정
conf->set("socket.connection.setup.timeout.ms", "5000", errstr);

”Topic ‘xxx’ not present in metadata”

토픽이 아직 없거나, 브로커의 auto.create.topics.enable이 꺼져 있는 경우입니다. librdkafka는 메타데이터에 없는 토픽으로 보낸 메시지를 message.timeout.ms 동안 붙잡고 있다가, 그때까지 토픽이 나타나지 않으면 dr_cb로 ERR__UNKNOWN_TOPIC이나 타임아웃을 보고합니다. 운영 환경에서는 자동 생성에 기대지 말고 파티션 수와 복제 수를 정해 미리 만드는 것이 원칙입니다.

# 토픽 사전 생성
kafka-topics --create --topic app-logs --partitions 3 --replication-factor 1 \
  --bootstrap-server localhost:9092

”Message delivery failed: Broker: Message size too large”

메시지가 브로커의 message.max.bytes(토픽별로는 max.message.bytes)보다 크면 발생합니다. 프로듀서, 브로커, 토픽, 컨슈머(fetch.max.bytes) 설정을 함께 맞춰야 하며, 수 MB 이상의 데이터라면 오브젝트 스토리지에 올리고 경로만 메시지로 보내는 편이 낫습니다.

// 프로듀서: 메시지 크기 제한
conf->set("message.max.bytes", "10485760", errstr);  // 10MB
// 브로커 설정: message.max.bytes, replica.fetch.max.bytes 동일하게

”Commit failed: Local: No offset stored”

인자 없는 commitSync()는 librdkafka가 저장해 둔(store된) 오프셋을 커밋하는데, 아직 소비한 메시지가 없거나 enable.auto.offset.store=false로 두고 offsets_store()를 호출하지 않았다면 커밋할 것이 없어 ERR__NO_OFFSET이 납니다. 정상적인 상황에서도 나올 수 있는 에러이므로 무시해도 되는 경우가 많고, 특정 메시지 기준으로 커밋하려면 commitSync(msg)를 씁니다.

consumer->commitSync(msg);  // 이 메시지 다음 위치를 커밋
// 또는
consumer->commitSync();     // 저장된 오프셋 전체 커밋 (없으면 ERR__NO_OFFSET)

메시지 중복 처리

처리 후 커밋 전에 크래시가 나면 재시작 후 같은 메시지를 다시 처리합니다. 처리 후 커밋하는 at-least-once 방식에서는 이 중복을 피할 수 없으므로, 앞의 멱등성 키처럼 소비자 쪽에서 중복을 걸러야 합니다. 반대로 처리 전에 커밋하면 중복은 없지만 처리 중 크래시가 나면 그 메시지는 유실됩니다(at-most-once).

// ✅ 처리 완료 후 커밋 (at-least-once, 중복 가능)
doBusinessLogic(msg);
consumer->commitSync(msg);
// ❌ 처리 전 커밋 (at-most-once, 처리 중 크래시 시 유실)
consumer->commitSync(msg);
doBusinessLogic(msg);

Kafka 토픽 간 처리에서 정확히 한 번이 필요하면 프로듀서의 enable.idempotence=true와 트랜잭션을 사용합니다. Kafka 고급(#52-6)에서 다룹니다.

”Out of memory” / “Queue full”

발행 속도가 전송 속도보다 빨라 내부 큐(queue.buffering.max.messages, queue.buffering.max.kbytes)가 가득 차면 produce()가 ERR__QUEUE_FULL을 반환합니다. 큐를 키우면 순간적인 폭주는 흡수하지만 근본 해결은 아니고 메모리 사용량만 늘어납니다. poll()로 전달 보고를 처리해 큐를 비우면서 재시도하는 것이 기본적인 백프레셔입니다.

RdKafka::ErrorCode err;
do {
    err = producer->produce(topic, RdKafka::Topic::PARTITION_UA,
                            RdKafka::Producer::RK_MSG_COPY,
                            const_cast<char*>(value.data()), value.size(),
                            nullptr, 0, 0, nullptr);
    if (err == RdKafka::ERR__QUEUE_FULL) producer->poll(100);  // 큐가 빌 때까지 대기
} while (err == RdKafka::ERR__QUEUE_FULL);

”Broker: Not leader for partition”

브로커 장애나 재배치로 파티션 리더가 바뀌면 일시적으로 발생합니다. librdkafka는 메타데이터를 갱신해 새 리더로 자동 재시도하며, 재시도 횟수(retries) 기본값이 사실상 무제한이고 전체 시간은 message.timeout.ms로 제한됩니다. retries를 작은 값으로 낮추면 오히려 일시적 리더 변경에도 메시지가 실패하므로, 보통은 기본값을 그대로 두고 message.timeout.ms로 포기 시점을 정합니다.

”Consumer group is rebalancing”

리밸런싱 중에 커밋하면 ERR__REBALANCE_IN_PROGRESS가 날 수 있습니다. 이 커밋은 실패한 것이므로 해당 메시지는 새 담당 컨슈머가 다시 처리할 수 있다고 가정해야 합니다. 앞의 콜백 절에서 설명한 대로 ERR__REVOKE_PARTITIONS 분기에서 처리를 끝낸 오프셋을 먼저 커밋해 두면 이런 중복이 줄어듭니다.

DeliveryReportCb가 호출되지 않음

dr_cb가 한 번도 실행되지 않는다면 poll()을 호출하지 않은 것입니다. delivery report는 poll()이나 flush()를 호출한 스레드에서 실행됩니다.

// ✅ 프로듀서 스레드에서 주기적 poll
void producer_loop() {
    while (running) {
        producer->produce(...);
        producer->poll(100);  // 필수!
    }
    producer->flush(10000);
}

Linger·압축·파티션 수로 처리량 올리기

배치 발행 (Linger)

메시지를 모아서 한 번에 보내면 요청 수와 요청당 오버헤드가 줄어듭니다.

conf->set("linger.ms", "20", errstr);       // 배치를 채우려고 최대 20ms 대기 (librdkafka 기본 5ms)
conf->set("batch.size", "1000000", errstr); // 배치 최대 바이트 (librdkafka 기본 1000000)

Java 클라이언트의 batch.size 기본값이 16384라서 이 값을 그대로 옮겨 쓰는 경우가 있는데, librdkafka 기본값은 약 1MB라 16KB로 설정하면 오히려 배치를 작게 쪼개게 됩니다. linger.ms를 늘리면 배치가 커져 처리량이 오르지만 메시지당 지연도 그만큼 늘어납니다.

압축

네트워크 대역폭을 줄입니다.

conf->set("compression.type", "gzip", errstr);
// none, gzip, snappy, lz4, zstd
압축속도압축률CPU
none빠름없음낮음
snappy빠름중간중간
lz4빠름중간낮음
gzip느림높음높음
zstd중간매우 높음중간

프로듀서 풀링

// ❌ 나쁜 예: 매 요청마다 새 프로듀서
void handle_request() {
    auto producer = create_producer();
    producer->produce(...);
    delete producer;
}
// ✅ 좋은 예: 프로듀서 재사용 (싱글톤 또는 풀)
static KafkaProducer* get_producer() {
    static KafkaProducer producer("localhost:9092");
    return &producer;
}

파티션 수와 컨슈머 수

한 파티션은 그룹 안에서 컨슈머 하나에만 할당되므로, 컨슈머 수가 파티션 수보다 많으면 남는 컨슈머는 놉니다. 파티션을 늘리면 병렬성이 늘지만 순서는 파티션 안에서만 보장되고, 운영 중에 파티션 수를 늘리면 같은 키가 다른 파티션으로 가게 되어 키별 순서가 깨질 수 있습니다.

// 키가 같은 메시지는 같은 파티션으로
producer->produce(topic, partition, ..., key, key_len, ...);
// partition = PARTITION_UA면 키 해시로 파티션 선택

Graceful Shutdown·키 설계·SSL/SASL

Graceful Shutdown

static std::atomic<bool> running{true};
void sig_handler(int) {
    running = false;
}
int main() {
    signal(SIGINT, sig_handler);
    signal(SIGTERM, sig_handler);
    // ... 프로듀서 생성
    while (running) {
        producer->produce(...);
        producer->poll(100);
    }
    // 종료 전 모든 메시지 전송 완료 대기
    producer->flush(10000);
    delete producer;
    RdKafka::wait_destroyed(5000);
    return 0;
}

메시지 키 설계

// 파티션 내 순서 보장: 같은 키 → 같은 파티션
// 예: 사용자별 이벤트 순서
std::string key = "user:" + std::to_string(user_id);
producer->produce(topic, RdKafka::Topic::PARTITION_UA, ..., key.data(), key.size(), ...);

헬스 체크

// 프로듀서: 메타데이터 조회로 브로커 연결 확인
RdKafka::Metadata* metadata = nullptr;
RdKafka::ErrorCode err = producer->metadata(true, nullptr, &metadata, 5000);
bool healthy = (err == RdKafka::ERR_NO_ERROR && metadata);
if (metadata) delete metadata;

환경 변수 기반 설정

struct KafkaConfig {
    std::string brokers = "localhost:9092";
    std::string group_id = "cpp-consumer";
    int linger_ms = 5;
};
KafkaConfig load_config() {
    KafkaConfig c;
    if (const char* b = std::getenv("KAFKA_BROKERS")) c.brokers = b;
    if (const char* g = std::getenv("KAFKA_GROUP_ID")) c.group_id = g;
    return c;
}

로깅 (debug)

// 문제 디버깅 시
conf->set("debug", "broker,topic,msg", errstr);
// broker, topic, msg, protocol, cgrp 등

SSL/TLS (프로덕션)

conf->set("security.protocol", "ssl", errstr);
conf->set("ssl.ca.location", "/path/to/ca-cert", errstr);
conf->set("ssl.certificate.location", "/path/to/client-cert", errstr);
conf->set("ssl.key.location", "/path/to/client-key", errstr);

SASL 인증

conf->set("security.protocol", "sasl_plaintext", errstr);
conf->set("sasl.mechanisms", "PLAIN", errstr);
conf->set("sasl.username", "user", errstr);
conf->set("sasl.password", "pass", errstr);

다음 글: Kafka 고급: 스트림 처리·트랜잭션·정확히 한 번(#52-6) 이전 글: C++ 시리즈 목차


참고 자료


같이 보면 좋은 글