C++ Kafka 고급 활용 | librdkafka·Exactly-Once Semantics (EOS)·트랜잭션·스트림 처리
들어가며: “스트림 처리·정확히 한 번이 막막해요”
Kafka 기본(#52-5)에서 프로듀서·컨슈머·오프셋·리밸런싱을 다뤘다면, 이 글에서는 고급 기능을 다룹니다. 실무에서 자주 맞닥뜨리는 문제와 해결 방법을 제시합니다.
시나리오 1: 토픽 A → 처리 → 토픽 B 파이프라인
상황: 클릭 로그를 읽어 집계 후 요약 토픽에 저장하는 파이프라인
문제: C++에서 Kafka Streams/KSQL 없이 어떻게 구현하나?
결과: 컨슈머-프로듀서 패턴으로 read → transform → produce
시나리오 2: 주문 이벤트 중복 처리로 재고가 음수
상황: 컨슈머 재시작 시 같은 메시지를 다시 처리해 재고가 두 번 차감됨
문제: at-least-once만으로는 중복 제거 불가
결과: 멱등성 키 설계 또는 정확히 한 번(idempotence + 트랜잭션) 적용
시나리오 3: 여러 토픽에 원자적으로 쓰기
상황: 주문 생성 시 orders 토픽 + order-events 토픽에 동시에 저장해야 함
문제: 중간에 실패하면 일부만 반영되어 데이터 불일치
결과: 트랜잭션 프로듀서로 여러 토픽에 원자적 발행
시나리오 4: 실시간 윈도우 집계
상황: 1분 단위로 API 호출 수를 집계해 대시보드에 표시
문제: C++에서 Kafka Streams 없이 윈도우 집계를 어떻게?
결과: 시간 윈도우 버퍼 + 주기적 flush로 시뮬레이션
시나리오 5: 컨슈머 처리 중 오프셋 커밋 타이밍
상황: 배치 처리 후 커밋 vs 메시지별 커밋 시 리밸런싱 시 중복 유실
문제: 리밸런싱 시 처리 중인 메시지가 다른 컨슈머로 넘어감
결과: 처리 완료 직후 커밋, 리밸런싱 콜백에서 상태 정리
flowchart TB
subgraph 문제[실무 문제]
P1[스트림 파이프라인] --> S1[컨슈머-프로듀서]
P2[중복 처리] --> S2[멱등성·정확히 한 번]
P3[원자적 다중 토픽] --> S3[트랜잭션]
P4[윈도우 집계] --> S4[버퍼·주기 flush]
P5[리밸런싱] --> S5[상태·커밋 전략]
end
목표:
- 스트림 처리: read → transform → produce
- 트랜잭션: 여러 토픽 원자적 발행
- 정확히 한 번: idempotence·멱등성·트랜잭션
- 윈도우 집계: C++에서 시뮬레이션
- 프로덕션: 모니터링·백프레셔·재시도 요구 환경: C++17 이상, librdkafka 1.6 이상(예제는 2.x 기준), 트랜잭션 예제는 Kafka 브로커 2.5 이상, Kafka 기본(#52-5) 선행
파이프라인·필터·시간 윈도우 스트림 처리
C++에서 스트림 처리란?
Kafka Streams나 KSQL은 JVM 기반입니다. 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)
단순 파이프라인: topic-a → topic-b
// stream_pipeline.cpp
// 컴파일: g++ -std=c++17 -o stream_pipeline stream_pipeline.cpp -lrdkafka -lrdkafka++
#include <librdkafka/rdkafkacpp.h>
#include <iostream>
#include <memory>
#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";
// 1. 컨슈머 설정
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;
if (!consumer) {
std::cerr << "컨슈머 생성 실패: " << errstr << std::endl;
return 1;
}
// 2. 프로듀서 설정
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;
if (!producer) {
std::cerr << "프로듀서 생성 실패: " << errstr << std::endl;
delete consumer;
return 1;
}
consumer->subscribe({"raw-logs"});
signal(SIGINT, sig_handler);
signal(SIGTERM, sig_handler);
std::cout << "스트림 파이프라인 시작 (raw-logs → processed-logs)" << std::endl;
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) {
std::cerr << "발행 실패: " << RdKafka::err2str(err) << std::endl;
} else {
consumer->commit(msg); // 처리 완료 후 커밋
}
break;
}
case RdKafka::ERR__TIMED_OUT:
case RdKafka::ERR__PARTITION_EOF:
break;
default:
std::cerr << "소비 에러: " << msg->errstr() << std::endl;
run = 0;
break;
}
delete msg;
producer->poll(0);
}
producer->flush(5000);
consumer->close();
delete producer;
delete consumer;
RdKafka::wait_destroyed(5000);
return 0;
}
코드 설명:
raw-logs→ 읽기 → 변환 →processed-logs발행produce()가 에러를 반환하지 않았을 때만commit()producer->poll(0): delivery report 처리
주의 (이 코드는 아직 at-least-once가 아닙니다): librdkafka의 produce()는 비동기라서, ERR_NO_ERROR 반환은 “브로커에 도착했다”가 아니라 “내부 전송 큐에 들어갔다”는 뜻일 뿐입니다. 위 코드처럼 그 직후 입력 오프셋을 커밋하면, 큐에 있던 메시지가 실제로 전송되기 전에 프로세스가 죽었을 때 입력은 소비된 것으로 기록되고 출력은 사라집니다. 제가 이런 파이프라인을 처음 만들 때 가장 늦게 깨달은 부분이 이것이었습니다. 평소에는 전혀 드러나지 않다가, 배포 중 강제 종료가 몇 번 겹치고 나서야 출력 토픽의 건수가 입력보다 적다는 것을 발견하게 됩니다.
올바른 at-least-once를 원한다면 두 가지 중 하나를 선택합니다. (1) delivery report 콜백(dr_cb)에서 전송 성공이 확인된 메시지의 입력 오프셋만 커밋 대상으로 기록하거나, (2) 일정 개수마다 producer->flush()로 전송 완료를 기다린 뒤 그때까지의 오프셋을 한꺼번에 커밋합니다. 메시지마다 동기 commit(msg)를 호출하는 것도 브로커 왕복이 매번 생겨 처리량이 크게 떨어지므로, 실무에서는 N개 또는 N밀리초마다 commitAsync로 묶는 것이 일반적입니다. 입력과 출력 오프셋을 원자적으로 묶고 싶다면 뒤에서 볼 트랜잭션이 정답입니다.
필터링: 조건에 맞는 메시지만 전달
// stream_filter.cpp - 에러 로그만 필터링
bool should_forward(const std::string& payload) {
return payload.find("ERROR") != std::string::npos ||
payload.find("FATAL") != std::string::npos;
}
// consume 루프 내부
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->commit(msg); // 필터 통과 여부와 관계없이 소비 완료
}
시간 윈도우 집계 (1분 단위)
// stream_window_aggregate.cpp - 1분 윈도우 카운트
#include <librdkafka/rdkafkacpp.h>
#include <chrono>
#include <map>
#include <string>
class WindowAggregator {
public:
using WindowKey = std::string;
using CountMap = std::map<WindowKey, int64_t>;
void add(const std::string& key, int64_t count = 1) {
auto window = current_window();
counts_[window][key] += count;
}
void flush_if_needed(RdKafka::Producer* producer) {
auto now = std::chrono::system_clock::now();
auto current = current_window();
for (auto it = counts_.begin(); it != counts_.end();) {
if (it->first < current) {
for (const auto& [k, v] : it->second) {
std::string msg = k + ":" + std::to_string(v);
producer->produce("aggregated", RdKafka::Topic::PARTITION_UA,
RdKafka::Producer::RK_MSG_COPY,
const_cast<char*>(msg.data()), msg.size(),
nullptr, 0, 0, nullptr);
}
it = counts_.erase(it);
} else {
++it;
}
}
}
private:
std::string current_window() {
auto now = std::chrono::system_clock::now();
auto sec = std::chrono::duration_cast<std::chrono::seconds>(
now.time_since_epoch()).count();
return std::to_string((sec / 60) * 60); // 1분 단위
}
std::map<std::string, CountMap> counts_;
};
이 집계기는 처리 시각(wall clock) 기준이라는 점을 알고 써야 합니다. 컨슈머가 밀려 있다가 따라잡는 동안에는 10분 전에 발생한 이벤트가 “지금” 윈도우로 들어가 집계가 왜곡되고, 재처리하면 결과가 달라집니다. 이벤트가 실제로 발생한 시각 기준(event time)으로 집계하려면 msg->timestamp().timestamp(밀리초)로 윈도우를 정하고, 늦게 도착하는 이벤트를 얼마나 기다릴지(워터마크)를 정해야 합니다. 또 메모리에만 있는 counts_는 재시작하면 사라지는데, 입력 오프셋은 이미 커밋되었을 수 있으므로 진행 중이던 윈도우의 부분 집계가 유실됩니다. 이것이 Kafka Streams가 상태를 changelog 토픽에 백업하는 이유이고, C++로 직접 구현할 때 가장 품이 많이 드는 부분이기도 합니다.
멱등성으로 정확히 한 번 전달하기
전달 시맨틱 비교
| 시맨틱 | 설명 | 중복 | 유실 |
|---|---|---|---|
| at-most-once | 커밋 후 처리 | 없음 | 가능 |
| at-least-once | 처리 후 커밋 | 가능 | 없음 |
| exactly-once | 멱등성·트랜잭션 | 없음 | 없음 |
멱등성 키로 중복 제거
비즈니스 로직에서 멱등성 키(idempotency key)를 사용해 중복 처리를 방지합니다.
// exactly_once_idempotent.cpp - 멱등성 키 기반 중복 제거
#include <unordered_set>
#include <string>
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() > max_size_) {
// FIFO: 오래된 키 제거 (간단히 처음 것 제거)
seen_.erase(seen_.begin());
}
return true;
}
private:
std::unordered_set<std::string> seen_;
std::mutex mutex_;
size_t max_size_ = 100000;
};
// 메시지 처리 시
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->commit(msg); // 커밋은 항상 (재처리 시 멱등성으로 스킵)
주의: idempotency_key는 Redis·DB에 저장해 재시작 후에도 유지하는 것이 좋습니다. 메모리 캐시는 재시작 시 초기화됩니다.
이 캐시에는 코드 자체의 함정도 두 가지 있습니다. 첫째, unordered_set::begin()은 가장 오래된 원소가 아니라 해시 순서상 첫 원소이므로 주석과 달리 FIFO가 아니라 임의의 키를 지웁니다. 최근에 처리한 키가 지워지면 바로 다음 중복을 놓칠 수 있으므로, 순서가 중요하다면 std::deque와 함께 쓰거나 시간 기반 만료를 둬야 합니다. 둘째, 키를 먼저 기록하고 나서 처리하므로, do_business_logic이 실패하거나 도중에 크래시하면 그 메시지는 “처리됨”으로 표시된 채 다시는 처리되지 않습니다. 가장 견고한 방식은 멱등성 키 기록과 비즈니스 효과(재고 차감 등)를 같은 DB 트랜잭션에 넣는 것입니다. processed_messages(key PRIMARY KEY) 테이블에 INSERT가 유니크 제약으로 실패하면 이미 처리된 메시지이고, 성공하면 같은 트랜잭션에서 재고를 차감한 뒤 커밋합니다. 이렇게 하면 크래시 시점과 무관하게 “정확히 한 번의 효과”가 보장됩니다.
프로듀서 멱등성 (enable.idempotence)
프로듀서가 메시지를 중복 발행하지 않도록 브로커 설정을 합니다.
// 프로듀서 설정
conf->set("enable.idempotence", "true", errstr);
// 이 설정 시 retries, acks, max.in.flight.requests.per.connection가 자동으로
// 안전한 값으로 설정됨
효과: 네트워크 오류로 재시도 시 브로커가 중복을 제거합니다. 브로커는 프로듀서 ID와 파티션별 시퀀스 번호로 “이미 받은 배치”를 알아보므로, 보장 범위는 파티션 단위, 같은 프로듀서 인스턴스가 살아 있는 동안입니다. 프로세스가 재시작되면 새 프로듀서 ID를 받기 때문에, 재시작 전후에 같은 메시지를 다시 보내는 것(애플리케이션 수준의 재전송)은 막지 못합니다. 또 acks=all이 강제되므로 acks=1로 지연을 줄이던 설정과는 함께 쓸 수 없습니다. librdkafka에서는 enable.idempotence=true인데 acks를 all이 아닌 값으로 명시하면 프로듀서 생성이 설정 에러로 실패합니다.
트랜잭션 프로듀서
트랜잭션이 필요한 경우
- 여러 토픽에 원자적으로 쓰기
- consume → transform → produce를 원자적으로 처리 (read-process-write)
sequenceDiagram
participant C as 컨슈머
participant App as 앱
participant P as 트랜잭션 프로듀서
P->>P: init_transactions()
P->>P: begin_transaction()
C->>App: consume
App->>P: produce (topic-a)
App->>P: produce (topic-b)
P->>P: commit_transaction()
App->>C: commit (offsets)
librdkafka 트랜잭션 API (C++)
librdkafka의 C++ 래퍼에서도 RdKafka::Producer의 init_transactions() 등으로 트랜잭션을 씁니다. 주의할 점은 이 함수들이 ErrorCode가 아니라 RdKafka::Error*를 반환한다는 것입니다. 성공이면 nullptr이고, 실패하면 호출자가 delete해야 하는 에러 객체가 돌아오며, 이 객체의 is_fatal(), is_retriable(), txn_requires_abort()로 다음 행동을 결정합니다.
// 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", "order-writer-0", errstr); // 인스턴스별로 고정된 ID
RdKafka::Producer* producer = RdKafka::Producer::create(conf, errstr);
delete conf;
if (!producer) {
std::cerr << "프로듀서 생성 실패: " << errstr << std::endl;
return 1;
}
// 1. 트랜잭션 초기화 (이전 인스턴스의 미완료 트랜잭션 정리, 좀비 펜싱)
if (RdKafka::Error* error = producer->init_transactions(5000)) {
std::cerr << "init_transactions 실패: " << error->str() << std::endl;
delete error;
delete producer;
return 1;
}
// 2. 트랜잭션 시작
if (RdKafka::Error* error = producer->begin_transaction()) {
std::cerr << "begin_transaction 실패: " << error->str() << std::endl;
delete error;
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. 트랜잭션 커밋 (내부적으로 남은 메시지를 flush한 뒤 커밋)
if (RdKafka::Error* error = producer->commit_transaction(5000)) {
std::cerr << "commit_transaction 실패: " << error->str() << std::endl;
if (error->txn_requires_abort()) {
if (RdKafka::Error* e2 = producer->abort_transaction(5000)) delete e2;
// 이 트랜잭션의 입력을 처음부터 다시 처리
} else if (error->is_retriable()) {
// commit_transaction을 다시 호출해도 됨 (타임아웃 등)
} else if (error->is_fatal()) {
// 펜싱 등: 이 프로듀서 인스턴스는 더 이상 쓸 수 없음 → 재생성 또는 종료
}
delete error;
}
delete producer;
RdKafka::wait_destroyed(5000);
return 0;
}
주의: transactional.id는 동시에 실행되는 인스턴스끼리는 서로 달라야 하지만, 같은 인스턴스가 재시작할 때는 같은 값이어야 합니다. 브로커는 같은 ID로 init_transactions()를 호출한 새 인스턴스가 나타나면 이전 인스턴스(네트워크가 끊긴 채 살아 있는 “좀비”일 수 있음)를 펜싱해, 두 인스턴스가 동시에 같은 출력을 쓰는 것을 막습니다. 그래서 transactional.id에 PID나 난수를 넣으면 재시작할 때마다 새 ID가 되어 이 보호가 사라지고, 미완료 트랜잭션이 타임아웃될 때까지 read_committed 컨슈머가 막히는 문제도 생깁니다. 쿠버네티스 StatefulSet의 파드 이름처럼 재시작해도 유지되는 식별자를 쓰는 것이 정석입니다. enable.idempotence는 트랜잭션 사용 시 자동으로 활성화됩니다.
read-process-write를 하나의 트랜잭션으로
위 예제는 출력만 원자적으로 묶습니다. 컨슈머에서 읽어 변환한 뒤 쓰는 파이프라인에서 진짜 exactly-once를 얻으려면, 입력 오프셋 커밋도 같은 트랜잭션에 넣어야 합니다. 컨슈머의 commit()을 따로 부르는 대신, 프로듀서의 send_offsets_to_transaction()에 다음에 읽을 오프셋(처리한 오프셋 + 1)과 컨슈머 그룹 메타데이터를 넘깁니다.
// consume → transform → produce 루프 안, 배치 하나를 처리한 뒤
std::vector<RdKafka::TopicPartition*> offsets;
offsets.push_back(RdKafka::TopicPartition::create(
msg->topic_name(), msg->partition(), msg->offset() + 1)); // "다음에 읽을" 오프셋
RdKafka::ConsumerGroupMetadata* cgmd = consumer->groupMetadata();
RdKafka::Error* error = producer->send_offsets_to_transaction(offsets, cgmd, 5000);
delete cgmd;
RdKafka::TopicPartition::destroy(offsets);
if (!error) error = producer->commit_transaction(5000);
// error 처리 규칙은 위와 동일 (txn_requires_abort면 abort 후 마지막 커밋 지점부터 재처리)
이렇게 하면 출력 메시지와 입력 오프셋이 함께 커밋되거나 함께 버려지므로, 크래시가 어디서 나든 다운스트림의 read_committed 컨슈머는 각 입력에 대한 출력을 정확히 한 번만 봅니다. 이때 컨슈머는 enable.auto.commit=false여야 하고, 트랜잭션을 abort했다면 컨슈머 위치도 마지막 커밋 지점으로 되돌려야(seek 또는 재할당) 합니다. 트랜잭션마다 코디네이터와 여러 번 통신하므로 메시지 하나마다 트랜잭션을 여닫으면 처리량이 크게 떨어집니다. 보통 수백~수천 개 메시지 또는 수백 밀리초 단위로 배치를 묶습니다.
컨슈머 read_committed
트랜잭션으로 커밋된 메시지만 읽으려면:
conf->set("isolation.level", "read_committed", errstr);
read_committed 컨슈머는 아직 커밋되지 않은 트랜잭션 메시지를 건너뛰는 것이 아니라, 그 트랜잭션이 끝날 때까지 그 지점 이후를 읽지 않고 기다립니다(LSO, Last Stable Offset). 그래서 프로듀서가 트랜잭션을 오래 열어 두거나, 크래시로 트랜잭션이 방치되면 컨슈머 지연(lag)이 갑자기 커집니다. 방치된 트랜잭션은 transaction.timeout.ms(기본 60초)가 지나야 브로커가 abort하므로, 이 값을 처리 배치 시간에 맞게 줄이는 것도 고려할 만합니다. 참고로 librdkafka 컨슈머의 기본값이 이미 read_committed입니다.
트랜잭션 미지원 브로커
트랜잭션은 Kafka 0.11부터 지원되며, send_offsets_to_transaction에 컨슈머 그룹 메타데이터를 쓰는 방식은 2.5 이상이 필요합니다. 로컬 개발용 단일 브로커에서 가장 흔히 막히는 원인은 버전이 아니라 transaction.state.log.replication.factor(기본 3)와 transaction.state.log.min.isr(기본 2)입니다. 브로커가 하나뿐이면 트랜잭션 상태 토픽을 만들 수 없어 init_transactions()가 코디네이터를 찾지 못한 채 타임아웃됩니다. 개발 환경에서는 두 값을 1로 낮춰야 합니다. 트랜잭션을 쓸 수 없는 환경이면 멱등성 키 + at-least-once로 대체합니다.
로그 집계·주문 이벤트·Graceful Shutdown 예제
로그 집계 파이프라인 (에러율 계산)
// log_aggregator.cpp - 에러 로그 비율 집계
#include <librdkafka/rdkafkacpp.h>
#include <iostream>
#include <chrono>
#include <atomic>
struct LogStats {
std::atomic<int64_t> total{0};
std::atomic<int64_t> errors{0};
};
void run_log_aggregator(const std::string& brokers) {
std::string errstr;
RdKafka::Conf* c_conf = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
c_conf->set("bootstrap.servers", brokers, errstr);
c_conf->set("group.id", "log-aggregator", errstr);
c_conf->set("enable.auto.commit", "false", errstr);
auto* 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);
auto* producer = RdKafka::Producer::create(p_conf, errstr);
delete p_conf;
consumer->subscribe({"app-logs"});
LogStats stats;
while (true) {
RdKafka::Message* msg = consumer->consume(1000);
if (msg->err() == RdKafka::ERR_NO_ERROR) {
std::string payload(static_cast<const char*>(msg->payload()), msg->len());
stats.total++;
if (payload.find("ERROR") != std::string::npos) stats.errors++;
// 1분마다 요약 발행
auto now = std::chrono::system_clock::now();
static auto last_flush = now;
if (std::chrono::duration_cast<std::chrono::milliseconds>(now - last_flush).count() > 60000) {
double rate = (stats.total > 0) ? (100.0 * stats.errors / stats.total) : 0;
std::string summary = "{\"total\":" + std::to_string(stats.total) +
",\"errors\":" + std::to_string(stats.errors) +
",\"error_rate\":" + std::to_string(rate) + "}";
producer->produce("log-summary", RdKafka::Topic::PARTITION_UA,
RdKafka::Producer::RK_MSG_COPY,
const_cast<char*>(summary.data()), summary.size(),
nullptr, 0, 0, nullptr);
last_flush = now;
}
consumer->commit(msg);
}
delete msg;
producer->poll(0);
}
delete producer;
delete consumer;
}
주문 이벤트 처리 (멱등성)
// order_processor.cpp - 멱등성 키로 중복 주문 방지
#include <librdkafka/rdkafkacpp.h>
#include <unordered_set>
#include <mutex>
#include <string>
class OrderProcessor {
public:
void process(RdKafka::Message* msg, RdKafka::KafkaConsumer* consumer) {
if (msg->err() != RdKafka::ERR_NO_ERROR) return;
std::string order_id = extract_order_id(msg);
if (order_id.empty()) {
consumer->commit(msg);
return;
}
{
std::lock_guard<std::mutex> lock(mutex_);
if (processed_.count(order_id)) {
consumer->commit(msg); // 이미 처리되며, 스킵
return;
}
processed_.insert(order_id);
if (processed_.size() > 100000) {
processed_.erase(processed_.begin());
}
}
// 실제 비즈니스 로직: 재고 차감, DB 저장 등
do_order_processing(order_id, msg);
consumer->commit(msg);
}
private:
std::string extract_order_id(RdKafka::Message* msg) {
if (msg->key()) {
return std::string(msg->key()->data(), msg->key()->size());
}
return "";
}
void do_order_processing(const std::string& order_id, RdKafka::Message* msg) {
(void)order_id;
(void)msg;
// 재고 차감, 주문 저장 등
}
std::unordered_set<std::string> processed_;
std::mutex mutex_;
};
Graceful Shutdown + Flush
// graceful_shutdown.cpp
static std::atomic<bool> running{true};
void shutdown_handler(int sig) {
(void)sig;
running = false;
}
int main() {
signal(SIGINT, shutdown_handler);
signal(SIGTERM, shutdown_handler);
// ... 프로듀서·컨슈머 생성
while (running) {
RdKafka::Message* msg = consumer->consume(1000);
if (msg->err() == RdKafka::ERR_NO_ERROR) {
process(msg);
consumer->commit(msg);
}
delete msg;
}
// 1. 프로듀서: 대기 중인 메시지 전송 완료
producer->flush(10000);
delete producer;
// 2. 컨슈머: 정리
consumer->close();
delete consumer;
RdKafka::wait_destroyed(5000);
return 0;
}
transaction coordinator 부재, Queue full, 리밸런스 중복: 에러 해결
init_transactions() 타임아웃 / “Coordinator not available”
증상: init_transactions()가 타임아웃되거나 코디네이터를 찾지 못했다는 에러로 실패.
원인: 대부분 단일 브로커 개발 환경에서 transaction.state.log.replication.factor(기본 3)와 transaction.state.log.min.isr(기본 2)를 충족할 수 없어 트랜잭션 상태 토픽이 만들어지지 않은 경우입니다. 운영 환경이라면 브로커 ACL에서 해당 transactional.id에 대한 권한이 빠진 경우도 있습니다.
해결법:
# 단일 브로커 개발 환경의 server.properties
transaction.state.log.replication.factor=1
transaction.state.log.min.isr=1
// ✅ 트랜잭션을 쓸 수 없는 환경이면 fallback: 멱등성 키 + at-least-once
if (RdKafka::Error* error = producer->init_transactions(5000)) {
std::cerr << "트랜잭션 사용 불가(" << error->str() << "), 멱등성 모드로 전환" << std::endl;
delete error;
use_transaction_ = false;
}
“Producer fenced” / “Broker: Producer attempted an operation with an old epoch”
증상: 잘 돌던 트랜잭션 프로듀서가 갑자기 치명적(fatal) 에러를 내고 더 이상 동작하지 않음.
원인: 같은 transactional.id를 쓰는 다른 인스턴스가 init_transactions()를 호출해 이 인스턴스가 펜싱된 것입니다. 롤링 배포 중 새 파드와 옛 파드가 잠시 공존하면 정상적으로 일어나는 일이고, 의도치 않게 두 인스턴스에 같은 ID를 설정한 경우에도 발생합니다.
해결법:
// ✅ 인스턴스마다 다르되, 재시작해도 유지되는 ID (예: StatefulSet 파드 이름)
std::string id = std::string("order-writer-") + std::getenv("POD_NAME");
conf->set("transactional.id", id, errstr);
// ❌ getpid()·난수는 재시작마다 바뀌어 좀비 펜싱이 동작하지 않음
펜싱된 프로듀서는 복구할 수 없으므로 객체를 폐기하고 프로세스를 종료하거나 새로 만들어야 합니다.
”Broker: Group rebalance in progress” / “Broker: Specified group generation id is not valid”
증상: 리밸런싱 직후 commit()이 ERR_REBALANCE_IN_PROGRESS 또는 ERR_ILLEGAL_GENERATION을 반환.
원인: 커밋하려는 파티션이 리밸런스로 이미 다른 컨슈머에게 넘어갔거나 넘어가는 중이라, 이전 세대(generation)의 커밋이 거부된 것입니다. 이 에러는 consume()이 반환한 메시지의 에러가 아니라 commit()의 반환값으로 나타납니다.
해결법:
// ✅ commit 반환값을 확인하고, 리밸런스 관련 에러는 재처리를 전제로 로그만 남김
RdKafka::ErrorCode cerr = consumer->commitSync(msg);
if (cerr == RdKafka::ERR_REBALANCE_IN_PROGRESS || cerr == RdKafka::ERR_ILLEGAL_GENERATION) {
// 이 메시지는 새 소유자가 다시 처리할 수 있음 → 처리가 멱등해야 하는 이유
} else if (cerr != RdKafka::ERR_NO_ERROR) {
std::cerr << "commit 실패: " << RdKafka::err2str(cerr) << std::endl;
}
근본적인 대책은 파티션이 회수되기 전에 커밋을 끝내는 것입니다. 아래 rebalance_cb의 회수(revoke) 분기에서 처리 중인 작업을 마무리하고 동기 커밋하십시오. 참고로 Local: No offset stored(ERR__NO_OFFSET)는 이와 다른 에러로, 커밋할 오프셋이 없을 때(아직 아무것도 소비하지 않았거나 이미 커밋한 뒤) 인자 없는 commitSync()를 호출하면 나오며 대개 무시해도 됩니다.
”Producer: Queue full”
증상: produce() 호출 시 ERR__QUEUE_FULL.
원인: 발행 속도 > 브로커 수신 속도.
해결법:
// ✅ 백프레셔: 큐가 비워질 때까지 대기
while (producer->outq_len() > 10000) {
producer->poll(100);
std::this_thread::sleep_for(std::chrono::milliseconds(10));
}
producer->produce(...);
더 단순하고 흔한 형태는 produce()가 ERR__QUEUE_FULL을 반환하면 poll(100)으로 delivery report를 처리해 큐에 자리를 만든 뒤 같은 메시지로 다시 produce()하는 재시도 루프입니다. 큐 크기 자체는 queue.buffering.max.messages와 queue.buffering.max.kbytes로 정해지는데, 이 값을 키우는 것은 근본적인 해결이 아니라 증상을 늦추는 것일 뿐입니다. 큐가 계속 찬다면 브로커 쪽 처리량(파티션 수, 네트워크, acks 설정)이나 배치·압축 설정을 먼저 봐야 합니다. 그리고 poll()을 전혀 호출하지 않는 코드는 delivery report가 쌓여 결국 큐가 가득 차므로, 프로듀서 전용 스레드든 메인 루프든 어딘가에서 주기적으로 poll()을 부르는 것이 필수입니다.
스트림 파이프라인에서 메시지 유실
증상: 컨슈머가 commit() 후 프로듀서 발행 전에 크래시.
원인: 처리 후 커밋 → 발행 순서가 잘못되면, 커밋은 됐는데 발행이 안 된 상태로 재시작됩니다. produce()가 성공을 반환한 직후 커밋하는 것도 같은 문제입니다. produce() 성공은 로컬 큐에 들어갔다는 뜻일 뿐이고, poll(0)은 그 시점까지 도착한 delivery report를 처리할 뿐 전송 완료를 기다리지 않습니다.
해결법:
// ✅ 전송 완료를 확인한 뒤 커밋 (at-least-once)
RdKafka::ErrorCode err = producer->produce(...);
if (err == RdKafka::ERR_NO_ERROR) {
// 배치 단위로: 전송 완료까지 대기 후 그 배치의 입력 오프셋 커밋
if (producer->flush(10000) == RdKafka::ERR_NO_ERROR /* && dr_cb에서 실패 없음 */) {
consumer->commitSync(msg);
}
}
// ✅ 정확히 한 번이 필요하면 트랜잭션 + send_offsets_to_transaction
“Broker: Not leader for partition” 반복
증상: 리더가 자주 바뀌는 파티션에서 지속적 실패. 원인: 브로커 불안정, 네트워크 지연. 해결법:
// ✅ 재시도·백오프 (librdkafka에서 retries는 message.send.max.retries의 별칭이므로 하나만 설정)
conf->set("retry.backoff.ms", "500", errstr);
conf->set("message.timeout.ms", "120000", errstr); // 재시도를 포함해 메시지 하나에 허용할 총 시간
리더 변경 자체는 브로커 재시작이나 파티션 재배치 때 정상적으로 일어나는 일이고, librdkafka는 메타데이터를 갱신해 새 리더로 자동 재시도합니다. 멱등성 프로듀서는 재시도 횟수가 사실상 무제한이 기본이므로, 재시도 횟수보다 message.timeout.ms(전체 전달 제한 시간)로 제어하는 것이 맞습니다. 에러가 계속된다면 설정보다 브로커 로그에서 컨트롤러 불안정, ISR 축소, 디스크 문제를 먼저 확인하십시오.
”Consumer group is rebalancing” 후 메시지 중복
증상: 리밸런싱 시 처리 중이던 메시지가 다른 컨슈머로 재할당되어 중복 처리.
원인: rebalance_cb에서 unassign 시 처리 중인 메시지 정리 전에 새 컨슈머가 할당받음.
해결법:
// ✅ rebalance_cb에서 파티션 해제 시 처리 중인 작업 완료 대기
void rebalance_cb(RdKafka::KafkaConsumer* consumer,
RdKafka::ErrorCode err,
std::vector<RdKafka::TopicPartition*>& partitions) override {
if (err == RdKafka::ERR__REVOKE_PARTITIONS) {
// 처리 중인 메시지 완료 대기 후, 회수되기 전에 오프셋 동기 커밋
wait_for_pending_processing();
consumer->commitSync();
consumer->unassign();
} else if (err == RdKafka::ERR__ASSIGN_PARTITIONS) {
consumer->assign(partitions);
}
}
이 콜백은 기본(eager) 리밸런스 프로토콜 기준입니다. eager 방식에서는 리밸런스가 일어날 때마다 그룹의 모든 컨슈머가 모든 파티션을 내려놓았다가 다시 받으므로, 컨슈머 하나가 추가되기만 해도 그룹 전체가 잠시 멈춥니다. partition.assignment.strategy=cooperative-sticky로 바꾸면 실제로 이동하는 파티션만 회수되지만, 콜백에서 assign/unassign 대신 incremental_assign(partitions)/incremental_unassign(partitions)을 호출해야 합니다. 이를 섞어 쓰면 Local: Erroneous state 같은 에러가 나거나 파티션 할당이 꼬이므로, 전략을 바꿀 때는 콜백도 함께 고쳐야 합니다(consumer->rebalance_protocol()로 현재 프로토콜을 확인할 수 있습니다). wait_for_pending_processing()이 max.poll.interval.ms(기본 5분)보다 오래 걸리면 컨슈머가 그룹에서 쫓겨나므로, 처리 시간이 긴 작업은 이 값을 조정하거나 작업 단위를 쪼개야 합니다.
”Message delivery failed: Broker: Message size too large”
증상: 큰 메시지 발행 시 실패. 해결법:
// 프로듀서
conf->set("message.max.bytes", "10485760", errstr); // 10MB
// 브로커: message.max.bytes, replica.fetch.max.bytes 동일하게
배치 처리·병렬 프로듀서·linger.ms 조정
스트림 파이프라인 배치 처리
메시지를 모아서 한 번에 처리하면 오버헤드가 줄어듭니다.
// 배치 크기만큼 모아서 처리
const size_t BATCH_SIZE = 100;
std::vector<RdKafka::Message*> batch;
while (batch.size() < BATCH_SIZE) {
RdKafka::Message* msg = consumer->consume(100);
if (msg->err() == RdKafka::ERR_NO_ERROR) {
batch.push_back(msg);
} else {
bool timed_out = (msg->err() == RdKafka::ERR__TIMED_OUT);
delete msg;
if (timed_out && !batch.empty()) break;
}
}
for (auto* m : batch) {
process_and_produce(m);
consumer->commit(m);
delete m;
}
프로듀서 병렬화
// ✅ 여러 스레드에서 동일 프로듀서 사용 (librdkafka는 스레드 안전)
std::thread t1([&]() { producer->produce(...); producer->poll(100); });
std::thread t2([&]() { producer->produce(...); producer->poll(100); });
압축
conf->set("compression.type", "lz4", errstr); // 속도·압축률 균형
linger.ms vs 처리량
linger.ms는 프로듀서가 배치를 채우기 위해 기다리는 최대 시간입니다. 값을 키우면 한 번에 더 많은 메시지를 묶어 보내 요청 수와 브로커 부하가 줄고 압축 효율도 좋아지지만, 그만큼 메시지 하나의 전송 지연이 늘어납니다. librdkafka의 기본값은 버전에 따라 다른데 최근 버전은 5ms입니다. 처리량과 지연의 실제 관계는 메시지 크기와 발행 속도에 크게 좌우되므로, 대상 환경에서 statistics.interval.ms 통계의 배치 크기와 요청 지연을 보며 조정하는 것이 좋습니다.
conf->set("linger.ms", "5", errstr);
conf->set("batch.size", "32768", errstr); // 32KB
파티션 수와 컨슈머 수
- 파티션 수 ≥ 컨슈머 수: 모든 컨슈머가 활용됨
- 파티션 수 < 컨슈머 수: 일부 컨슈머는 idle
# 토픽 생성 시 파티션 수 결정
kafka-topics --create --topic events --partitions 12 --replication-factor 2 \
--bootstrap-server localhost:9092
헬스 체크·메트릭·SSL/SASL 설정
헬스 체크
bool health_check(RdKafka::Producer* producer) {
RdKafka::Metadata* meta = nullptr;
RdKafka::ErrorCode err = producer->metadata(true, nullptr, &meta, 5000);
bool ok = (err == RdKafka::ERR_NO_ERROR && meta);
if (meta) delete meta;
return ok;
}
bool health_check(RdKafka::KafkaConsumer* consumer) {
std::vector<RdKafka::TopicPartition*> partitions;
RdKafka::ErrorCode err = consumer->assignment(partitions);
bool ok = (err == RdKafka::ERR_NO_ERROR);
for (auto* p : partitions) delete p;
return ok;
}
메트릭 수집
// librdkafka 통계 (JSON)
conf->set("statistics.interval.ms", "10000", errstr);
class StatsCb : public RdKafka::EventCb {
public:
void event_cb(RdKafka::Event& event) override {
if (event.type() == RdKafka::Event::EVENT_STATS) {
// event.str()에 JSON 통계
std::cout << event.str() << std::endl;
}
}
};
환경 변수 기반 설정
struct KafkaConfig {
std::string brokers = "localhost:9092";
std::string group_id = "cpp-consumer";
std::string transactional_id;
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;
if (const char* t = std::getenv("KAFKA_TRANSACTIONAL_ID")) c.transactional_id = t;
return c;
}
SSL/TLS + SASL
conf->set("security.protocol", "sasl_ssl", errstr);
conf->set("sasl.mechanisms", "PLAIN", errstr);
conf->set("sasl.username", std::getenv("KAFKA_USER"), errstr);
conf->set("sasl.password", std::getenv("KAFKA_PASSWORD"), errstr);
conf->set("ssl.ca.location", "/etc/ssl/certs/ca-certificates.crt", errstr);
std::getenv는 변수가 없으면 nullptr을 반환하는데, 이것을 std::string 인자에 그대로 넘기면 미정의 동작(보통 크래시)입니다. 필수 자격 증명은 시작 시점에 존재 여부를 검사하고 없으면 명확한 에러로 종료하게 하십시오. 또 conf->set()의 반환값(RdKafka::Conf::CONF_OK)을 확인하지 않으면 설정 이름 오타나 빌드에 SASL/SSL 지원이 빠진 librdkafka를 쓰는 경우(Unsupported value "sasl_ssl" for configuration property "security.protocol": OpenSSL not available at build time 같은 메시지)를 놓치게 됩니다.
멀티 브로커
conf->set("bootstrap.servers",
"broker1:9092,broker2:9092,broker3:9092", errstr);
디버그 로깅
conf->set("debug", "broker,topic,msg", errstr);
// 프로덕션에서는 제거
스트림·트랜잭션 구현 점검 항목
스트림 처리
- consume → transform → produce 순서 준수
- delivery report(또는 flush)로 전송 완료 확인 후 commit (유실 방지)
- producer->poll() 주기적 호출
정확히 한 번
- 멱등성 키 설계 (키 또는 topic-partition-offset)
- 멱등성 캐시 저장소 (Redis/DB 권장)
- enable.idempotence (프로듀서)
트랜잭션
- transactional.id 설정 (인스턴스별로 다르고, 재시작해도 유지되는 값)
- init_transactions → begin → produce → commit
- 실패 시 abort_transaction
- isolation.level=read_committed (컨슈머)
에러 처리
- produce() 반환값 검사
- msg->err() 검사
- ERR__QUEUE_FULL 시 백프레셔
- rebalance_cb 구현
프로덕션
- Graceful Shutdown (flush, close)
- SSL/TLS, SASL (필요 시)
- 환경 변수로 설정 외부화
- 헬스 체크·메트릭
기능별 요약
| 항목 | 요약 |
|---|---|
| 스트림 처리 | consume → transform → produce, 발행 후 커밋 |
| 정확히 한 번 | 멱등성 키 + enable.idempotence + 트랜잭션 |
| 트랜잭션 | transactional.id, init_transactions, begin/commit/abort |
| 에러 | Transaction coordinator, Queue full, Rebalance |
| 성능 | 배치 처리, linger.ms, 압축, 파티션 수 |
| 프로덕션 | Graceful Shutdown, SSL/SASL, 헬스 체크 |
핵심 원칙:
- 발행 후 커밋: 처리 완료 후 프로듀서 발행 성공 확인 후 commit
- 멱등성: 비즈니스 키로 중복 처리 방지
- 트랜잭션: 다중 토픽 원자적 쓰기
- 리밸런싱: rebalance_cb에서 상태 정리·재할당 처리
자주 묻는 질문 (FAQ)
Q. produce() 호출에서 ERR__QUEUE_FULL이 나면 어떻게 해야 하나요?
A. librdkafka는 produce()를 비동기로 처리하므로, 발행 속도가 브로커로 보내는 속도를 앞지르면 내부 큐가 가득 차서 ERR__QUEUE_FULL을 반환합니다. 본문처럼 outq_len()이 임계값을 넘으면 poll()을 호출해 전송 완료 콜백을 처리하면서 큐가 비워질 때까지 기다리는 백프레셔를 넣어야 합니다. 에러를 무시하고 다음 메시지로 넘어가면 그 메시지는 조용히 유실되므로 반드시 재시도 루프로 감싸는 것이 안전합니다.
Q. C++에서 Kafka Streams 대신 뭘 쓰나요?
A. Kafka Streams는 JVM 전용입니다. C++에서는 컨슈머-프로듀서 패턴으로 read → transform → produce를 구현합니다. 복잡한 윈도우·조인은 별도 상태 저장소(Redis, DB)와 함께 구현합니다.
Q. 정확히 한 번 전달이 꼭 필요한가요?
A. 금융·주문처럼 중복 효과가 곧 사고인 시스템에서는 필요하지만, 그 보장을 Kafka 트랜잭션만으로 얻을 수 있는지는 따로 따져야 합니다. Kafka EOS는 Kafka 토픽에서 읽어 Kafka 토픽에 쓰는 구간만 보장하므로, 처리 과정에서 외부 DB를 갱신한다면 멱등성 키를 그 DB 트랜잭션에 함께 기록하는 방식이 실제 보장을 만듭니다. 로그·메트릭은 at-least-once로도 충분한 경우가 많습니다.
Q. 트랜잭션 프로듀서가 실패하면?
A. 반환된 RdKafka::Error의 종류에 따라 다릅니다. txn_requires_abort()면 abort_transaction() 후 마지막 커밋 지점부터 입력을 다시 처리하고 begin_transaction()부터 재시작합니다. is_retriable()이면 같은 호출을 다시 시도하면 되고, is_fatal()(펜싱 등)이면 그 프로듀서 인스턴스는 폐기하고 새로 만들어야 합니다.
스트림 처리·트랜잭션·정확히 한 번 전달을 C++로 구현해 실무에 적용할 수 있습니다.
다음 글: C++ 시리즈 목차
이전 글: C++에서 Kafka 쓰기
참고 자료
- librdkafka 공식 문서
- Apache Kafka 트랜잭션
- C++에서 Kafka 쓰기 — 프로듀서·컨슈머 기초