C++에서 RabbitMQ·Kafka 연동하기: 메시지 큐 기반 비동기 처리 구현과 에러 대응

들어가며: “동기 호출이 병목이라 서비스가 느려요”

메시지 큐가 필요한 이유

REST API로 서비스 A가 B를 호출하면, B가 응답할 때까지 A가 블로킹됩니다. B가 DB 쿼리·외부 API·무거운 연산으로 5초 걸리면 A도 5초를 기다려야 합니다. 사용자가 “주문 완료” 버튼을 눌렀을 때 결제·재고 차감·이메일 발송·로그 기록을 모두 동기로 처리하면 응답이 10초 넘게 걸릴 수 있습니다. 메시지 큐는 “메시지를 보내고 바로 다음 일을 한다”는 비동기 패턴으로, 느슨한 결합과 부하 분산을 가능하게 합니다.

SimpleAmqpClient로 RabbitMQ 프로듀서/컨슈머를, librdkafka로 Kafka 프로듀서/컨슈머를 구현하고, JSON·Protobuf 메시지 직렬화와 운영 중에 겪는 문제를 다룹니다.

다만 메시지 큐가 공짜는 아닙니다. 동기 호출에서는 “결제가 실패했다”는 사실을 즉시 사용자에게 돌려줄 수 있지만, 큐에 넣고 나면 실패는 나중에 다른 프로세스에서 일어나고, 그 결과를 사용자에게 알리는 경로(상태 조회 API, 알림, 보상 트랜잭션)를 따로 설계해야 합니다. 또 브로커라는 운영 대상이 하나 늘고, “정확히 한 번” 처리는 브로커가 보장해 주지 않으므로 컨슈머 쪽 멱등성까지 직접 챙겨야 합니다. 응답에 결과가 꼭 필요한 호출(잔액 조회, 로그인)은 동기로 남기고, 결과를 나중에 알려줘도 되는 부수 효과(메일, 통계, 썸네일 생성)부터 큐로 빼는 것이 현실적인 순서입니다.

동기 호출 병목으로 메시지 큐가 필요한 상황

주문 처리 지연

이커머스에서 주문이 들어오면 결제·재고·이메일·로그를 순차 처리합니다. 동기 호출이면 하나만 느려도 전체가 지연됩니다. 메시지 큐에 작업을 넣고 각 워커가 비동기로 처리하면, 사용자는 즉시 “주문 접수됨” 응답을 받고 백엔드는 여유 있게 처리할 수 있습니다.

트래픽 급증 시 서비스 다운

블랙 프라이데이에 주문이 100배로 늘면 동기 API는 DB·외부 API에 과부하를 주고, 결국 타임아웃으로 연쇄 장애가 납니다. 메시지 큐는 버퍼 역할을 해서 피크 시 메시지를 쌓아 두고, 컨슈머가 처리 가능한 속도로 소비합니다.

마이크로서비스 간 통신

서비스 A가 B, C, D를 각각 호출하는 구조에서는 B가 다운되면 A도 실패합니다. 이벤트 기반으로 A가 “주문 생성됨” 이벤트를 큐에 발행하면 B·C·D는 각자 구독해 독립적으로 처리합니다. B가 잠시 다운돼도 메시지는 큐에 남아 있다가 복구 후 처리됩니다.

로그·메트릭 수집

수십 대 서버에서 로그를 중앙으로 보낼 때 동기 HTTP를 쓰면 수집 서버가 병목이 됩니다. Kafka 같은 스트리밍 플랫폼에 로그를 발행하면 여러 컨슈머가 병렬로 소비해 실시간 분석·저장을 할 수 있습니다.

작업 큐 (Task Queue)

이미지 리사이징, 비디오 인코딩처럼 무거운 작업을 워커 풀에 분배할 때, 메시지 큐가 작업을 공정하게 분배하며, 워커 장애 시 다른 워커가 재시도할 수 있게 합니다.


RabbitMQ와 Kafka 아키텍처 비교

RabbitMQ vs Kafka 비교

항목RabbitMQKafka
모델메시지 브로커 (Push)이벤트 스트리밍 (Pull)
메시지 보존컨슈머 ACK 후 삭제 (기본)로그 보존 (일정 기간)
처리량수만 msg/s수십만~수백만 msg/s
순서큐 단위 보장파티션 단위 보장
적합 용도작업 큐, RPC, 라우팅로그, 스트리밍, 대용량 이벤트

처리량 수치는 하드웨어, 메시지 크기, 영속화·복제 설정에 따라 크게 달라지므로 “대략 몇 자릿수 차이가 날 수 있다” 정도로만 받아들이는 편이 좋습니다. 두 브로커를 가르는 더 본질적인 차이는 메시지가 소비된 뒤에도 남는가입니다. RabbitMQ 큐는 ACK된 메시지를 지우므로 “해야 할 일 목록”에 가깝고, 새 컨슈머가 과거 메시지를 다시 읽을 수 없습니다. Kafka 토픽은 보존 기간 동안 로그를 남기므로 새 컨슈머 그룹이 처음부터 다시 읽거나, 버그를 고친 뒤 오프셋을 되돌려 재처리할 수 있습니다. 반대로 Kafka에는 메시지별 ACK·재시도·지연 전달 같은 기능이 없어 작업 큐 용도로 쓰면 그 부분을 애플리케이션이 직접 만들어야 합니다. 작업을 워커에 분배하고 개별 실패를 재시도해야 하면 RabbitMQ, 같은 이벤트를 여러 시스템이 각자의 속도로 읽고 재생해야 하면 Kafka가 자연스럽습니다.

전체 아키텍처

flowchart TB
    subgraph Producer["프로듀서 (C++)"]
        P1[메시지 생성]
        P2[직렬화]
        P3[발행]
        P1 --> P2 --> P3
    end
    subgraph RabbitMQ[RabbitMQ]
        R1[Exchange]
        R2[Queue]
        R1 --> R2
    end
    subgraph Kafka[Kafka]
        K1[Topic]
        K2[Partition 0]
        K3[Partition 1]
        K1 --> K2
        K1 --> K3
    end
    subgraph Consumer["컨슈머 (C++)"]
        C1[구독/폴링]
        C2[역직렬화]
        C3[처리]
        C1 --> C2 --> C3
    end
    Producer -->|AMQP| RabbitMQ
    Producer -->|Kafka Protocol| Kafka
    RabbitMQ -->|Push/Pull| Consumer
    Kafka -->|Pull| Consumer

시퀀스 다이어그램 (RabbitMQ)

sequenceDiagram
    participant P as C++ 프로듀서
    participant R as RabbitMQ
    participant C as C++ 컨슈머
    P->>R: Connection + Channel 생성
    P->>R: BasicPublish(exchange, routing_key, body)
    R->>R: Exchange → Queue 라우팅
    C->>R: BasicConsume(queue)
    loop 메시지 처리
        R->>C: BasicConsumeMessage() → 메시지 전달
        C->>C: 비즈니스 로직 처리
        C->>R: BasicAck(delivery_tag)
    end

시퀀스 다이어그램 (Kafka)

sequenceDiagram
    participant P as C++ 프로듀서
    participant K as Kafka Broker
    participant C as C++ 컨슈머
    P->>K: produce(topic, partition, key, value)
    K->>K: 파티션에 append
    C->>K: subscribe(topic)
    loop poll()
        C->>K: fetch(partition, offset)
        K-->>C: 메시지 배치 반환
        C->>C: 처리 후 commit(offset)
    end

RabbitMQ 프로듀서·컨슈머 구현

환경 설정

의존성: SimpleAmqpClient (rabbitmq-c 래퍼)

# vcpkg로 설치
vcpkg install simpleamqpclient
# CMakeLists.txt
find_package(SimpleAmqpClient REQUIRED)
target_link_libraries(myapp SimpleAmqpClient::SimpleAmqpClient)

프로듀서 구현

// rabbitmq_producer.hpp
#pragma once
#include <SimpleAmqpClient/SimpleAmqpClient.h>
#include <string>
#include <stdexcept>
class RabbitMQProducer {
public:
    RabbitMQProducer(const std::string& host = "localhost",
                     int port = 5672,
                     const std::string& user = "guest",
                     const std::string& password = "guest")
        : host_(host), port_(port), user_(user), password_(password) {}
    void connect(const std::string& queue = "task_queue") {
        channel_ = AmqpClient::Channel::Create(host_, port_, user_, password_);
        channel_->DeclareQueue(queue, /*passive=*/false,
                              /*durable=*/true, /*exclusive=*/false,
                              /*auto_delete=*/false);
    }
    void publish(const std::string& queue, const std::string& message) {
        if (!channel_) throw std::runtime_error("Not connected");
        auto msg = AmqpClient::BasicMessage::Create(message);
        msg->DeliveryMode(2);  // persistent (디스크에 저장)
        channel_->BasicPublish("", queue, msg);
    }
private:
    std::string host_, user_, password_;
    int port_;
    AmqpClient::Channel::ptr_t channel_;
};

주의점:

  • DeliveryMode(2): 브로커 재시작 시에도 메시지 유지
  • DeclareQueue의 durable=true: 큐 자체도 영구 저장

두 설정은 반드시 짝으로 써야 합니다. durable 큐에 non-persistent 메시지를 넣으면 브로커 재시작 후 큐는 남아 있지만 비어 있고, persistent 메시지를 non-durable 큐에 넣으면 큐와 함께 메시지도 사라집니다. 그리고 persistent여도 BasicPublish가 반환됐다는 것은 “소켓에 썼다”는 뜻일 뿐, 브로커가 디스크에 기록했다는 보장이 아닙니다. 발행 직후 브로커가 죽는 상황까지 막으려면 뒤에서 다루는 publisher confirm을 켜야 합니다.

DeclareQueue는 같은 인자로 여러 번 호출해도 안전한 멱등 연산이지만, 이미 존재하는 큐를 다른 속성으로 선언하면 채널이 닫히며 PRECONDITION_FAILED - inequivalent arg 'durable' for queue 'task_queue' 같은 에러가 납니다. 개발 중 durable=false로 만들어 둔 큐를 나중에 true로 바꾸면 이 에러를 만나므로, 관리 UI(15672 포트)에서 기존 큐를 지우고 다시 선언해야 합니다. 이 글의 connect()가 큐 이름을 인자로 받는 것도, 프로듀서와 컨슈머가 같은 이름·같은 속성으로 큐를 선언하도록 맞추기 위해서입니다.

컨슈머 구현

// rabbitmq_consumer.hpp
#pragma once
#include <SimpleAmqpClient/SimpleAmqpClient.h>
#include <string>
#include <functional>
#include <atomic>
// 타입 정의
class RabbitMQConsumer {
public:
    using MessageHandler = std::function<bool(const std::string&)>;
    RabbitMQConsumer(const std::string& host = "localhost", int port = 5672,
                     const std::string& user = "guest",
                     const std::string& password = "guest")
        : host_(host), port_(port), user_(user), password_(password) {}
    void connect(const std::string& queue = "task_queue") {
        channel_ = AmqpClient::Channel::Create(host_, port_, user_, password_);
        channel_->DeclareQueue(queue, false, true, false, false);
        // QoS: 한 번에 하나씩 처리 (fair dispatch)
        channel_->BasicQos(0, 1, false);
    }
    void consume(const std::string& queue, MessageHandler handler) {
        if (!channel_) throw std::runtime_error("Not connected");
        // 인자: queue, consumer_tag, no_local, no_ack, exclusive
        std::string consumer_tag = channel_->BasicConsume(queue, "", true,
                                                          /*no_ack=*/false, false);
        running_.store(true);
        while (running_.load()) {
            AmqpClient::Envelope::ptr_t envelope;
            if (channel_->BasicConsumeMessage(consumer_tag, envelope, 1000)) {
                std::string body = envelope->Message()->Body();
                if (handler(body)) {
                    channel_->BasicAck(envelope);
                } else {
                    channel_->BasicReject(envelope, /*requeue=*/true);
                }
            }
        }
    }
    void stop() { running_.store(false); }
private:
    std::string host_, user_, password_;
    int port_;
    AmqpClient::Channel::ptr_t channel_;
    std::atomic<bool> running_{true};
};

핵심:

  • BasicQos(0, 1, false): prefetch=1로 한 번에 하나씩만 가져와 워커 간 공정 분배
  • BasicAck: 처리 성공 시 ACK, 실패 시 BasicReject(..., true)로 재큐잉

SimpleAmqpClient의 BasicConsume은 인자 순서가 (queue, consumer_tag, no_local, no_ack, exclusive, prefetch)이고 no_ack의 기본값이 true입니다. 인자를 생략하고 BasicConsume(queue)만 호출하면 자동 ACK 모드가 되어, 아래 코드에서 BasicAck를 불러도 의미가 없고 컨슈머가 죽는 순간 받아 둔 메시지가 사라집니다. 세 번째 인자를 no_ack로 착각해 true를 넘기는 실수도 흔하므로 위 코드처럼 주석으로 인자 이름을 적어 두는 것을 권합니다.

BasicReject(..., true)의 재큐잉에는 함정이 있습니다. JSON 파싱 실패처럼 몇 번을 다시 시도해도 실패할 메시지(poison message)를 requeue=true로 돌려보내면 메시지가 곧바로 큐 앞쪽으로 되돌아와 같은 컨슈머가 다시 받고, 다시 실패하는 루프가 CPU를 점유합니다. 로그가 초당 수천 줄씩 같은 에러로 도배된다면 이 상황입니다. 일시적 오류(DB 타임아웃)와 영구적 오류(형식 오류)를 구분해, 영구적 오류는 requeue=false로 거절해 뒤에서 설명할 DLQ로 보내야 합니다.


Kafka 프로듀서·컨슈머 구현

환경 설정

의존성: librdkafka

# vcpkg로 설치
vcpkg install rdkafka
# CMakeLists.txt
find_package(RdKafka REQUIRED)
target_link_libraries(myapp RdKafka::rdkafka)

프로듀서 구현

// kafka_producer.hpp
#pragma once
#include <librdkafka/rdkafkacpp.h>
#include <string>
#include <iostream>
class KafkaProducer {
public:
    KafkaProducer(const std::string& brokers) {
        std::string errstr;
        auto conf = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
        conf->set("bootstrap.servers", brokers, errstr);
        conf->set("dr_cb", &dr_cb_, errstr);
        producer_.reset(RdKafka::Producer::create(conf, errstr));
        if (!producer_) {
            throw std::runtime_error("Producer create failed: " + 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.c_str(), key.size(),
            0, nullptr);
        if (err != RdKafka::ERR_NO_ERROR) {
            std::cerr << "Produce failed: " << RdKafka::err2str(err) << "\n";
            return;
        }
        producer_->poll(0);
    }
    void flush(int timeout_ms = 10000) {
        // 내부 큐가 빌 때까지(또는 timeout까지) 전송하고 delivery report 콜백 실행
        if (producer_->flush(timeout_ms) != RdKafka::ERR_NO_ERROR)
            std::cerr << producer_->outq_len() << " message(s) not delivered\n";
    }
private:
    struct DeliveryReportCb : public RdKafka::DeliveryReportCb {
        void dr_cb(RdKafka::Message& msg) override {
            if (msg.err())
                std::cerr << "Delivery failed: " << msg.errstr() << "\n";
        }
    } dr_cb_;
    std::unique_ptr<RdKafka::Producer> producer_;
};

librdkafka의 produce()는 네트워크로 바로 보내지 않고 내부 큐에 넣기만 하고 즉시 반환합니다. 실제 전송은 백그라운드 스레드가 배치로 처리하고, 성공·실패 결과는 poll()을 호출하는 스레드에서 dr_cb로 전달됩니다. 그래서 produce()가 ERR_NO_ERROR를 돌려줬다고 해서 메시지가 브로커에 도착한 것이 아닙니다. 전송 실패(브로커 다운, 토픽 없음, message.timeout.ms 초과)는 오직 delivery report로만 알 수 있으므로, dr_cb에서 실패를 로그만 찍고 버리는 위 코드는 실서비스라면 재발행이나 보관 로직으로 바꿔야 합니다.

자주 겪는 두 가지 증상도 이 구조에서 나옵니다. 첫째, poll()을 주기적으로 부르지 않으면 delivery report가 쌓이기만 하고 내부 큐가 가득 차서 produce()가 ERR__QUEUE_FULL(Local: Queue full)을 반환합니다. 이때는 poll(100)으로 잠시 비운 뒤 재시도하면 됩니다. 둘째, 짧게 실행되는 프로그램이 flush() 없이 종료하면 큐에 남은 메시지가 전송되지 않은 채 사라집니다. 로그 수집 예제에서 produce() 뒤에 flush()를 부르는 이유가 이것입니다.

컨슈머 구현

// kafka_consumer.hpp
#pragma once
#include <librdkafka/rdkafkacpp.h>
#include <string>
#include <functional>
#include <atomic>
class KafkaConsumer {
public:
    using MessageHandler = std::function<bool(const std::string&)>;
    KafkaConsumer(const std::string& brokers, const std::string& group_id) {
        std::string errstr;
        auto 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);  // 수동 commit
        consumer_.reset(RdKafka::KafkaConsumer::create(conf, errstr));
        if (!consumer_) {
            throw std::runtime_error("Consumer create failed: " + errstr);
        }
    }
    void subscribe(const std::string& topic) {
        std::vector<std::string> topics = {topic};
        RdKafka::ErrorCode err = consumer_->subscribe(topics);
        if (err) {
            throw std::runtime_error("Subscribe failed: " +
                                    RdKafka::err2str(err));
        }
    }
    void consume(MessageHandler handler) {
        running_.store(true);
        while (running_.load()) {
            RdKafka::Message* msg = consumer_->consume(1000);
            if (!msg) continue;
            if (msg->err()) {
                // 타임아웃·파티션 끝은 정상 흐름, 나머지만 에러로 기록
                if (msg->err() != RdKafka::ERR__TIMED_OUT &&
                    msg->err() != RdKafka::ERR__PARTITION_EOF)
                    std::cerr << "Consume error: " << msg->errstr() << "\n";
                delete msg;
                continue;
            }
            std::string payload(static_cast<const char*>(msg->payload()),
                               msg->len());
            if (handler(payload)) {
                consumer_->commit(msg);
            }
            delete msg;
        }
    }
    void stop() { running_.store(false); }
    ~KafkaConsumer() { if (consumer_) consumer_->close(); }  // 그룹에서 정상 탈퇴
private:
    std::unique_ptr<RdKafka::KafkaConsumer> consumer_;
    std::atomic<bool> running_{true};
};

consume(1000)은 1초 동안 메시지가 없으면 nullptr가 아니라 ERR__TIMED_OUT 에러가 담긴 메시지 객체를 돌려줍니다. 이 경우를 따로 거르지 않으면 조용한 토픽에서 “Consume error: Local: Timed out”이 1초마다 찍히고, 에러 분기에서 delete를 빠뜨리면 메모리가 계속 샙니다. 또 close()를 부르지 않고 프로세스를 끝내면 브로커는 session.timeout.ms가 지날 때까지 이 컨슈머가 살아 있다고 보고 파티션을 다른 컨슈머에 넘기지 않으므로, 재배포 때마다 그만큼 소비가 멈춥니다.

여기서 가장 오해하기 쉬운 부분은 커밋의 의미입니다. Kafka의 커밋은 “이 메시지를 처리했다”가 아니라 “이 파티션은 오프셋 N까지 읽었다”는 기록입니다. 위 코드에서 handler가 false를 반환하면 커밋을 건너뛰지만, 다음 메시지가 성공해 커밋되는 순간 실패한 메시지의 오프셋도 함께 넘어갑니다. 결국 실패한 메시지는 재시작하거나 리밸런스될 때만 다시 올 뿐, 정상 흐름에서는 조용히 건너뛰어집니다. RabbitMQ의 reject처럼 메시지 하나만 되돌리는 기능이 없으므로, 재처리가 필요하면 seek()으로 오프셋을 되돌리거나 실패 메시지를 재시도 토픽으로 다시 발행해야 합니다. 메시지마다 동기 commit()을 부르는 것도 브로커 왕복이 매번 생겨 처리량을 크게 떨어뜨리므로, 운영에서는 일정 개수마다 커밋하거나 commitAsync()를 쓰는 편이 일반적입니다.


주문 처리·로그 수집 파이프라인 예제

주문 처리 파이프라인 (RabbitMQ)

주문 생성 → 큐에 발행 → 워커가 결제·재고·이메일 처리

// order_processor.cpp
#include "rabbitmq_producer.hpp"
#include "rabbitmq_consumer.hpp"
#include <nlohmann/json.hpp>
#include <iostream>
using json = nlohmann::json;
int main(int argc, char* argv[]) {
    if (argc < 2) {
        std::cerr << "Usage: " << argv[0] << " [producer|consumer]\n";
        return 1;
    }
    std::string mode = argv[1];
    if (mode == "producer") {
        RabbitMQProducer producer("localhost", 5672, "guest", "guest");
        producer.connect("order_queue");  // 컨슈머와 같은 이름으로 선언
        json order = {
            {"order_id", "ORD-001"},
            {"user_id", "user_123"},
            {"amount", 29900},
            {"items", {"item_a", "item_b"}}
        };
        producer.publish("order_queue", order.dump());
        std::cout << "Order published\n";
    } else if (mode == "consumer") {
        RabbitMQConsumer consumer("localhost", 5672, "guest", "guest");
        consumer.connect("order_queue");
        consumer.consume("order_queue", [](const std::string& body) {
            try {
                auto order = json::parse(body);
                std::cout << "Processing order: " << order["order_id"] << "\n";
                // 결제, 재고 차감, 이메일 발송 등
                return true;  // 성공 → ACK
            } catch (const std::exception& e) {
                std::cerr << "Error: " << e.what() << "\n";
                return false;  // 실패 → NACK, 재큐잉
            }
        });
    }
    return 0;
}

이 예제에서 주문 ID를 메시지에 넣는 것은 단순한 식별용이 아닙니다. 수동 ACK 구조에서는 “결제는 끝났는데 ACK 직전에 워커가 재시작됨” 같은 상황에서 같은 주문 메시지가 한 번 더 전달됩니다. 결제·재고처럼 두 번 실행되면 안 되는 작업은 order_id를 기준으로 이미 처리했는지 확인할 수 있어야 하므로, 메시지에는 반드시 재처리 여부를 판별할 수 있는 안정적인 키를 담아야 합니다. 또 하나의 워커가 결제·재고·메일을 모두 처리하면 메일 서버 장애 때문에 결제까지 재시도되는 일이 생기므로, 실무에서는 “주문 생성됨” 이벤트를 fanout exchange로 발행하고 결제·재고·메일 큐를 각각 두어 실패를 서로 격리하는 구성이 더 흔합니다.

로그 수집 파이프라인 (Kafka)

여러 서버에서 로그를 Kafka에 발행, 컨슈머가 Elasticsearch/파일에 저장

// log_collector.cpp
#include "kafka_producer.hpp"
#include "kafka_consumer.hpp"
#include <nlohmann/json.hpp>
#include <chrono>
#include <iostream>
int main(int argc, char* argv[]) {
    if (argc < 2) {
        std::cerr << "Usage: " << argv[0] << " [producer|consumer] [brokers]\n";
        return 1;
    }
    std::string mode = argv[1];
    std::string brokers = argc > 2 ? argv[2] : "localhost:9092";
    if (mode == "producer") {
        KafkaProducer producer(brokers);
        json log = {
            {"timestamp", std::time(nullptr)},
            {"level", "INFO"},
            {"service", "api-server"},
            {"message", "Request processed"}
        };
        producer.produce("app-logs", "api-server", log.dump());
        producer.flush();
    } else if (mode == "consumer") {
        KafkaConsumer consumer(brokers, "log-consumer-group");
        consumer.subscribe("app-logs");
        consumer.consume([](const std::string& body) {
            auto log = json::parse(body);
            std::cout << "[" << log["level"] << "] " << log["message"] << "\n";
            return true;
        });
    }
    return 0;
}

로그 프로듀서에서 키로 "api-server"를 넘기는 것은 순서와 분산 사이의 선택입니다. 같은 키는 같은 파티션으로 가므로 한 서비스의 로그는 순서가 보장되지만, 서비스가 몇 개뿐이면 트래픽이 소수 파티션에 몰리고 나머지 컨슈머는 놀게 됩니다. 로그처럼 전역 순서가 중요하지 않은 데이터라면 키를 비워 파티션에 고르게 퍼지게 하고, 주문 이벤트처럼 같은 엔티티의 순서가 중요하면 order_id를 키로 쓰는 식으로 데이터의 성격에 맞게 정합니다.

메시지 직렬화 (Protobuf)

// order.proto
syntax = "proto3";
package order;
message OrderEvent {
    string order_id = 1;
    string user_id = 2;
    int64 amount = 3;
    repeated string items = 4;
}
// protobuf_serialization.cpp
#include "order.pb.h"
#include "rabbitmq_producer.hpp"
#include <iostream>
int main() {
    order::OrderEvent evt;
    evt.set_order_id("ORD-002");
    evt.set_user_id("user_456");
    evt.set_amount(15000);
    evt.add_items("item_x");
    std::string serialized;
    evt.SerializeToString(&serialized);
    RabbitMQProducer producer("localhost");
    producer.connect("order_queue");
    producer.publish("order_queue", serialized);
    return 0;
}

Protobuf를 메시지 포맷으로 쓰면 크기와 파싱 비용이 줄어드는 대신, 큐에 쌓인 바이트를 관리 UI나 kafka-console-consumer로 봐도 사람이 읽을 수 없게 됩니다. 장애 조사 때 이 차이가 꽤 크게 느껴집니다. 그리고 큐는 프로듀서와 컨슈머가 서로 다른 시점에 배포되는 구조라서 스키마 호환성이 중요합니다. 필드 번호를 재사용하거나 타입을 바꾸면 이미 큐에 쌓인 옛 메시지를 새 컨슈머가 잘못 해석하므로, 필드는 추가만 하고 삭제한 번호는 reserved로 막아 두는 규칙을 지켜야 합니다.

Docker Compose로 전체 스택 실행

# docker-compose.mq.yml
version: '3.8'
services:
  rabbitmq:
    image: rabbitmq:3.12-management
    ports:
      - "5672:5672"
      - "15672:15672"
    environment:
      RABBITMQ_DEFAULT_USER: guest
      RABBITMQ_DEFAULT_PASS: guest
  zookeeper:
    image: confluentinc/cp-zookeeper:7.5.0
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000
  kafka:
    image: confluentinc/cp-kafka:7.5.0
    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   # 브로커 1대 환경 필수

단일 브로커로 띄울 때 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1을 빠뜨리면, 컨슈머 그룹이 오프셋을 저장할 내부 토픽 __consumer_offsets를 기본 복제 계수 3으로 만들려다 실패합니다. 프로듀서는 잘 동작하는데 컨슈머만 아무 메시지도 받지 못하고 멈춰 있다면 먼저 이 설정을 확인합니다. KAFKA_ADVERTISED_LISTENERS의 localhost:9092는 호스트에서 실행한 C++ 프로그램 기준 주소이므로, 컨슈머를 같은 Compose 네트워크의 다른 컨테이너에서 실행하면 localhost가 그 컨테이너 자신을 가리켜 연결이 실패합니다. 이때는 내부용·외부용 리스너를 따로 광고해야 합니다.


연결 실패, ACK 전 크래시, 리밸런스: 메시지 큐 에러 해결

”Connection refused” — RabbitMQ/Kafka에 연결 실패

증상: Connection refused 또는 ECONNREFUSED 원인:

  1. 브로커가 실행 중이 아님
  2. 호스트/포트 오류
  3. 방화벽 차단 해결법:
# RabbitMQ 상태 확인
docker ps | grep rabbitmq
# 또는
rabbitmqctl status
# Kafka 상태 확인
docker exec -it kafka kafka-broker-api-versions --bootstrap-server localhost:9092
// ✅ 연결 재시도 로직 추가
void connect_with_retry(int max_retries = 5) {
    for (int i = 0; i < max_retries; ++i) {
        try {
            channel_ = AmqpClient::Channel::Create(host_, port_, user_, password_);
            return;
        } catch (const std::exception& e) {
            std::cerr << "Connect failed, retry " << (i+1) << "/" << max_retries << "\n";
            std::this_thread::sleep_for(std::chrono::seconds(2));
        }
    }
    throw std::runtime_error("Connection failed after retries");
}

“Queue not found” — RabbitMQ 큐 미선언

증상: 컨슈머 쪽에서 NOT_FOUND - no queue 'task_queue'로 채널이 닫힘. 프로듀서 쪽은 에러 없이 성공한 것처럼 보이지만 메시지가 어디에도 없음 원인: 큐를 Declare하지 않은 상태에서 소비하거나 발행함

두 증상이 비대칭이라는 점이 핵심입니다. 존재하지 않는 큐를 BasicConsume하면 브로커가 NOT_FOUND로 채널을 닫아 바로 알 수 있습니다. 하지만 기본 exchange("")로 없는 큐 이름을 라우팅 키로 발행하면 AMQP는 에러 없이 메시지를 버립니다. 컨슈머보다 프로듀서를 먼저 배포했는데 큐 선언 코드가 컨슈머에만 있으면, 그 사이에 발행된 메시지는 흔적 없이 사라집니다. 처음 RabbitMQ를 붙일 때 “발행 로그는 찍히는데 큐가 비어 있다”는 상황에서 가장 먼저 확인할 것이 이 부분입니다.

해결법:

// ❌ 나쁜 예: 큐 선언 없이 발행 → 큐가 없으면 에러 없이 메시지 폐기
channel_->BasicPublish("", "task_queue", msg);
// ✅ 좋은 예: 프로듀서·컨슈머 양쪽에서 발행/소비 전 큐 선언 (멱등)
channel_->DeclareQueue("task_queue", false, true, false, false);
channel_->BasicPublish("", "task_queue", msg);
// 라우팅 실패를 감지하려면 mandatory=true로 발행 (SimpleAmqpClient는
// 반환된 메시지를 MessageReturnedException으로 알려줌)
channel_->BasicPublish("", "task_queue", msg, /*mandatory=*/true);

메시지 유실 — ACK 전 컨슈머 크래시

증상: 메시지가 처리됐는데 큐에서 사라짐 (유실) 원인: BasicConsume의 no_ack=true로 설정하면 메시지가 전달 즉시 삭제됩니다. 컨슈머가 크래시하면 메시지가 유실됩니다. SimpleAmqpClient는 no_ack의 기본값이 true이므로 인자를 생략해도 같은 문제가 생깁니다. 해결법:

// 인자 순서: queue, consumer_tag, no_local, no_ack, exclusive
// ❌ 나쁜 예: no_ack=true → 전달 즉시 삭제
channel_->BasicConsume(queue, "", true, /*no_ack=*/true, false);
// ✅ 좋은 예: no_ack=false, 처리 후 명시적 ACK
channel_->BasicConsume(queue, "", true, /*no_ack=*/false, false);
// ... 처리 ...
channel_->BasicAck(envelope);

반대 방향의 실수도 있습니다. 수동 ACK 모드로 바꿔 놓고 특정 분기(예외 처리 경로)에서 ACK도 reject도 하지 않으면, 그 메시지는 관리 UI에서 “Unacked” 상태로 남아 prefetch 슬롯을 계속 차지합니다. prefetch=1이면 컨슈머가 아무것도 받지 못하고 멈춘 것처럼 보이고, 연결이 끊겨야 비로소 다른 컨슈머에게 재전달됩니다. 모든 경로에서 ACK나 reject가 정확히 한 번 호출되는지 확인해야 합니다.


Kafka “Broker: Topic authorization failed”

증상: TopicAuthorizationException 또는 Not authorized to access topic 원인: ACL 설정 또는 SASL 인증 오류 해결법:

// SASL/SSL 설정 (librdkafka)
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 ACL 확인
kafka-acls --bootstrap-server localhost:9092 --list

“Consumer group rebalance” — 처리 중 오프셋 커밋 실패

증상: 같은 메시지가 여러 번 처리되거나, 일부 메시지 건너뜀 원인: enable.auto.commit=true이면 폴링 시점에 자동 커밋됩니다. 처리 완료 전에 리밸런스되면 다른 컨슈머가 같은 메시지를 가져갑니다. 해결법:

// ✅ 수동 커밋: 처리 성공 후에만 commit
conf->set("enable.auto.commit", "false", errstr);
// ...
if (handler(payload)) {
    consumer_->commit(msg);  // 처리 성공 시에만
}

수동 커밋은 유실을 막아 주지만 중복을 막지는 못합니다. 처리를 마치고 커밋하기 직전에 리밸런스가 일어나면, 새로 파티션을 받은 컨슈머는 마지막 커밋 오프셋부터 읽으므로 같은 메시지를 다시 처리합니다. 처리 시간이 max.poll.interval.ms(기본 5분)를 넘는 경우도 조심해야 합니다. 그 시간 안에 다음 consume()을 부르지 못하면 브로커는 컨슈머가 멈췄다고 판단해 그룹에서 내보내고, 이후 커밋은 Broker: Specified group generation id is not valid나 Local: Assignment lost류의 에러로 실패합니다. 무거운 작업은 처리 시간을 줄이거나 이 값을 늘리고, 어떤 경우든 컨슈머를 멱등하게 만드는 것이 최종 방어선입니다.


RabbitMQ 메모리 부족 — 메시지 적체

증상: RabbitMQ 메모리 사용량 급증, memory_alarm 발생 원인: 프로듀서가 컨슈머보다 빠르게 발행, 큐에 메시지 적체 해결법:

// 1. 프로듀서 발행 속도 제한 (throttling)
// 2. publisher confirm 사용해 브로커 수용 여부 확인
channel_->ConfirmSelect();
channel_->BasicPublish("", "task_queue", msg);
channel_->WaitForConfirms();  // 브로커가 수신할 때까지 대기
// 3. 큐 max-length 설정 (RabbitMQ 정책)
// rabbitmqctl set_policy max_len "task_queue" '{"max-length":10000}' --apply-to queues

배치 발행·압축·Prefetch로 처리량 올리기

RabbitMQ: 배치 발행 (Batch Publish)

// ❌ 메시지마다 네트워크 왕복
for (const auto& msg : messages) {
    channel_->BasicPublish("", queue, msg);
}
// ✅ 배치로 발행 (가능한 경우)
// SimpleAmqpClient는 배치 API가 없음 → 채널 재사용으로 오버헤드 감소
// 또는 publisher confirms를 비동기로 처리해 블로킹 최소화

Kafka: 배치 크기 및 압축

conf->set("batch.size", "16384", errstr);      // 16KB 배치
conf->set("linger.ms", "5", errstr);          // 5ms 대기 후 발행
conf->set("compression.type", "snappy", errstr);  // 압축

Kafka: 파티션 수 조정

파티션 수 ≥ 컨슈머 수일 때 병렬 처리 극대화

# 토픽 생성 시 파티션 수 지정
kafka-topics --create --topic orders --partitions 8 --replication-factor 1

RabbitMQ: Prefetch 조정

// prefetch=1: 공정 분배, 처리 시간 긴 작업에 적합
channel_->BasicQos(0, 1, false);
// prefetch=10: 처리량 우선, 짧은 작업에 적합
channel_->BasicQos(0, 10, false);

메시지 직렬화 최적화

// JSON vs Protobuf: Protobuf는 필드 이름을 싣지 않는 바이너리 포맷이라
// 보통 더 작고 파싱이 빠름 (차이는 메시지 구조에 따라 다르므로 직접 측정)
// 대용량/고빈도 메시지는 Protobuf 권장

이 절의 설정들은 모두 지연 시간과 처리량을 맞바꾸는 손잡이입니다. linger.ms를 늘리면 배치가 커져 처리량이 오르지만 메시지 하나가 전송되기까지 최대 그만큼 더 기다립니다. prefetch를 늘리면 네트워크 왕복이 줄지만, 한 컨슈머가 메시지를 잔뜩 가져간 채 느리게 처리하는 동안 다른 컨슈머는 놀 수 있고, 그 컨슈머가 죽으면 가져간 메시지가 모두 재전달됩니다. 튜닝 전에 현재 병목이 브로커·네트워크인지, 컨슈머의 처리 로직(DB 쓰기 등)인지부터 확인하는 편이 좋습니다. 경험상 컨슈머 쪽 DB 쓰기가 병목인 경우가 훨씬 많고, 그때는 브로커 설정을 아무리 바꿔도 적체가 줄지 않습니다.


DLQ·Exponential Backoff·Idempotent Consumer

Dead Letter Queue (DLQ)

처리 실패 메시지를 별도 큐로 보내 나중에 수동 처리

# RabbitMQ 정책으로 DLQ 설정 (권장: 앱 재배포 없이 변경 가능)
rabbitmqctl set_policy dlx_policy "task_queue" \
  '{"dead-letter-exchange":"dlx","dead-letter-routing-key":"failed"}' \
  --apply-to queues
// 또는 큐 선언 시 arguments로 설정 (rabbitmq-c/AMQP-CPP 등)
// x-dead-letter-exchange, x-dead-letter-routing-key를 Table에 포함

재시도 (Exponential Backoff)

bool process_with_retry(const std::string& body) {
    const int max_retries = 5;
    for (int i = 0; i < max_retries; ++i) {
        try {
            return handle_message(body);
        } catch (const std::exception& e) {
            if (i == max_retries - 1) return false;
            std::this_thread::sleep_for(
                std::chrono::milliseconds(100 * (1 << i)));
        }
    }
    return false;
}

Idempotent Consumer

같은 메시지가 여러 번 전달돼도 한 번만 처리

std::unordered_set<std::string> processed_ids;
bool handle_order(const std::string& body) {
    auto order = json::parse(body);
    std::string id = order["order_id"];
    if (processed_ids.count(id)) return true;  // 이미 처리됨 → ACK
    // ... 실제 처리 ...
    processed_ids.insert(id);
    return true;
}

위 코드는 개념을 보여 주기 위한 것이고, 그대로 운영에 쓰면 세 가지 문제가 생깁니다. 메모리에만 있으므로 컨슈머가 재시작하면 기록이 사라지고(중복은 바로 재시작 직후에 가장 많이 생깁니다), 컨슈머가 여러 대면 서로의 기록을 모르며, 집합이 끝없이 커집니다. 실무에서는 처리 결과를 쓰는 DB 트랜잭션 안에서 processed_messages(message_id PRIMARY KEY) 테이블에 함께 INSERT하고, 중복 키 에러가 나면 이미 처리된 것으로 보고 ACK하는 방식이 가장 단단합니다. 비즈니스 쓰기와 “처리됨” 기록이 같은 트랜잭션에 묶이기 때문에 둘 중 하나만 반영되는 틈이 없습니다.

헬스 체크

// RabbitMQ 연결 상태 확인
bool is_connected() {
    try {
        channel_->GetChannelId();
        return true;
    } catch (...) {
        return false;
    }
}

메트릭 수집

// 발행/소비 메트릭 (Prometheus 등)
void record_published(const std::string& queue, size_t size) {
    metrics_.counter("mq_messages_published_total", {{"queue", queue}}).inc();
    metrics_.histogram("mq_message_size_bytes", buckets_, {{"queue", queue}})
        .observe(static_cast<double>(size));
}

브로커별 구현 점검 항목

환경 설정

  • RabbitMQ/Kafka Docker 또는 로컬 설치
  • SimpleAmqpClient, librdkafka (vcpkg) 설치
  • 연결 정보 (host, port, credentials) 환경 변수화

RabbitMQ

  • 큐 선언 (durable=true)
  • 메시지 persistent (DeliveryMode=2)
  • no_ack=false, 처리 후 BasicAck
  • BasicQos(prefetch) 설정
  • DLQ 설정 (선택)

Kafka

  • enable.auto.commit=false, 수동 commit
  • 처리 성공 후에만 commit
  • consumer group.id 설정
  • 파티션 수 및 replication factor 검토

공통

  • 연결 재시도 로직
  • 메시지 직렬화 (JSON/Protobuf) 일관성
  • 에러 처리 및 로깅
  • 메트릭 수집 (발행/소비 수, 지연)

보안

  • 프로덕션에서 기본 guest/guest 비밀번호 변경
  • TLS/SSL 적용 (필요 시)
  • Kafka SASL 인증 (필요 시)

RabbitMQ·Kafka 선택 요약

항목설명
RabbitMQ작업 큐, RPC, 라우팅, Push 모델
Kafka로그, 스트리밍, 대용량, Pull 모델
직렬화JSON(간편), Protobuf(성능)
에러 처리ACK/NACK, DLQ, 재시도
프로덕션수동 commit, DLQ, 멱등성, 메트릭

핵심 원칙:

  1. 비동기로 서비스 간 결합도 낮추기
  2. 메시지 유실 방지: ACK, persistent, 수동 commit
  3. 멱등성으로 중복 처리 방지
  4. 모니터링으로 적체·지연 감지

자주 묻는 질문 (FAQ)

Q. RabbitMQ와 Kafka 중 뭘 써야 하나요?

A. 작업 큐, RPC, 복잡한 라우팅이 필요하면 RabbitMQ. 로그 수집, 이벤트 스트리밍, 초당 수십만 메시지가 필요하면 Kafka. 둘 다 쓰는 하이브리드 구성도 많습니다.

Q. C++에서 쓸 수 있는 다른 메시지 큐 라이브러리는 무엇인가요?

A. RabbitMQ는 SimpleAmqpClient, AMQP-CPP 등이 있습니다. Kafka는 librdkafka가 사실상 표준입니다. ZeroMQ는 브로커 없이 P2P 메시징에 적합합니다.

Q. 메시지 순서를 보장하려면 어떻게 해야 하나요?

A. RabbitMQ는 단일 큐 + 단일 컨슈머로 순서를 보장합니다. Kafka는 같은 키가 같은 파티션으로 가므로 파티션 내 순서가 보장됩니다. 파티션 키를 메시지별로 바꾸면 순서가 깨질 수 있습니다.


같이 보면 좋은 글