C++에서 RabbitMQ 쓰기: SimpleAmqpClient·rabbitmq-c, 토픽 라우팅, Dead Letter, 메시지 TTL

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

회원가입이나 주문 완료 시 이메일을 동기적으로 보내면 SMTP 서버가 느릴 때 API 응답 전체가 함께 느려집니다. 주문 서비스가 재고·결제·알림 서비스를 차례로 HTTP로 호출하는 구조도 하나가 멈추면 전체가 멈춥니다. 업로드된 이미지를 여러 워커가 HTTP 폴링으로 가져가 처리하면 어떤 워커가 어떤 작업을 가져갔는지 추적하기 어렵고, 처리 도중 워커가 죽으면 작업이 사라지거나 두 번 처리됩니다.

메시지 큐를 두면 API는 작업을 큐에 넣고 바로 반환하고, 워커가 큐에서 꺼내 처리합니다. RabbitMQ에서는 prefetch로 워커마다 가져갈 수 있는 미확인 메시지 수를 제한해 느린 워커에 일이 몰리지 않게 하고, durable 큐와 persistent 메시지로 브로커가 재시작해도 메시지를 유지하며, 처리를 마친 뒤에만 ACK를 보내 크래시한 워커의 메시지가 다른 워커에게 다시 전달되게 합니다. 로그 레벨별 라우팅(topic exchange), 긴급 작업 우선 처리(우선순위 큐), 실패 메시지 수집(Dead Letter Exchange), 오래된 메시지 자동 만료(TTL)도 브로커 설정으로 해결할 수 있습니다.

flowchart LR
  subgraph Producer["프로듀서 (C++)"]
    P1[주문 이벤트]
    P2[이메일 작업]
    P3[이미지 처리]
  end
  subgraph RabbitMQ[RabbitMQ]
    Q1[task_queue]
    Q2[email_queue]
  end
  subgraph Consumer["컨슈머 (C++)"]
    C1[주문 처리]
    C2[이메일 발송]
    C3[이미지 리사이징]
  end
  P1 --> Q1
  P2 --> Q2
  P3 --> Q1
  Q1 --> C1
  Q1 --> C3
  Q2 --> C2

RabbitMQ vs Kafka 비교

항목RabbitMQKafka
모델큐, Exchange토픽, 파티션
메시지 보존소비 후 삭제 (기본)보존 기간 동안 유지
재처리별도 구현오프셋 이동으로 가능
처리량높음 (메시지 크기·확인 방식에 따라 다름)매우 높음 (배치·순차 디스크 쓰기에 최적화)
C++ 클라이언트rabbitmq-c, SimpleAmqpClientlibrdkafka

RabbitMQ 브로커와 SimpleAmqpClient·rabbitmq-c 설치

필수 의존성

항목버전비고
C++C++14 이상C++17 권장
RabbitMQ3.8+브로커 (Docker 권장)
SimpleAmqpClient2.xrabbitmq-c C++ 래퍼
rabbitmq-c0.14+AMQP C 라이브러리
CMake3.16+find_package 지원

RabbitMQ 브로커 실행 (Docker)

# RabbitMQ + 관리 UI (포트 5672 AMQP, 15672 웹)
docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672 \
  rabbitmq:3.12-management
# 기본 계정: guest / guest (로컬만)
# 관리 UI: http://localhost:15672

SimpleAmqpClient 설치 (vcpkg 권장)

# vcpkg로 설치 (rabbitmq-c 자동 의존)
vcpkg install simpleamqpclient
# CMake 사용 시
cmake -B build -DCMAKE_TOOLCHAIN_FILE=[vcpkg]/scripts/buildsystems/vcpkg.cmake

rabbitmq-c 직접 설치 (SimpleAmqpClient 없이)

# Ubuntu/Debian
sudo apt-get install librabbitmq-dev
# macOS (Homebrew)
brew install rabbitmq-c
# 소스 빌드
git clone https://github.com/alanxz/rabbitmq-c.git
cd rabbitmq-c
mkdir build && cd build
cmake .. -DCMAKE_BUILD_TYPE=Release
cmake --build .
sudo cmake --build . --target install

CMakeLists.txt 기본 설정

cmake_minimum_required(VERSION 3.16)
project(rabbitmq_demo LANGUAGES CXX)
set(CMAKE_CXX_STANDARD 17)
# 설치 방법에 따라 CMake 패키지 이름·타깃이 다를 수 있음 (없으면 pkg-config의 libSimpleAmqpClient 사용)
find_package(SimpleAmqpClient REQUIRED)
add_executable(rabbitmq_producer producer.cpp)
target_link_libraries(rabbitmq_producer PRIVATE SimpleAmqpClient::SimpleAmqpClient)
add_executable(rabbitmq_consumer consumer.cpp)
target_link_libraries(rabbitmq_consumer PRIVATE SimpleAmqpClient::SimpleAmqpClient)

rabbitmq-c만 사용 시 CMakeLists.txt

find_package(rabbitmq-c REQUIRED)
target_link_libraries(rabbitmq_producer PRIVATE rabbitmq::rabbitmq)

첫 연결과 Hello World 프로듀서·컨슈머

SimpleAmqpClient로 최소 연결

AmqpClient::Channel::Create()로 호스트·포트·인증 정보를 넘겨 채널을 생성합니다. SimpleAmqpClient의 Channel 객체 하나가 TCP 연결 하나와 그 위의 채널을 함께 관리합니다. 최근 버전에서는 이 생성 함수가 OpenOpts를 받는 형태로 바뀌고 기존 형태는 deprecated 되었으므로, 사용하는 버전의 헤더를 확인합니다.

#include <SimpleAmqpClient/SimpleAmqpClient.h>
#include <iostream>
int main() {
    try {
        // 1. 채널 생성 (연결 + 로그인 자동)
        auto channel = AmqpClient::Channel::Create("localhost", 5672, "guest", "guest");
        // 2. 큐 선언 (없으면 생성)
        channel->DeclareQueue("hello", /*passive=*/false,
                             /*durable=*/false, /*exclusive=*/false,
                             /*auto_delete=*/false);
        // 3. 메시지 발행
        auto msg = AmqpClient::BasicMessage::Create("Hello, RabbitMQ!");
        channel->BasicPublish("", "hello", msg);
        std::cout << "메시지 발행 완료" << std::endl;
    } catch (const std::exception& e) {
        std::cerr << "에러: " << e.what() << std::endl;
        return 1;
    }
    return 0;
}

DeclareQueue("hello", ...)는 큐가 없으면 만들고 있으면 그대로 쓰며, passive=true면 존재 여부만 확인합니다. 빈 exchange("")로 발행하면 기본 exchange가 routing key와 이름이 같은 큐로 메시지를 보냅니다.

Hello World: 프로듀서 + 컨슈머

// producer_hello.cpp
#include <SimpleAmqpClient/SimpleAmqpClient.h>
#include <iostream>
int main() {
    auto channel = AmqpClient::Channel::Create("localhost", 5672, "guest", "guest");
    channel->DeclareQueue("hello", false, false, false, false);
    for (int i = 0; i < 5; ++i) {
        std::string body = "Hello " + std::to_string(i);
        auto msg = AmqpClient::BasicMessage::Create(body);
        channel->BasicPublish("", "hello", msg);
        std::cout << "Sent: " << body << std::endl;
    }
    return 0;
}
// consumer_hello.cpp
#include <SimpleAmqpClient/SimpleAmqpClient.h>
#include <iostream>
#include <csignal>
static volatile sig_atomic_t g_running = 1;
void sig_handler(int) { g_running = 0; }
int main() {
    signal(SIGINT, sig_handler);
    auto channel = AmqpClient::Channel::Create("localhost", 5672, "guest", "guest");
    channel->DeclareQueue("hello", false, false, false, false);
    // 인자 순서: (queue, consumer_tag, no_local, no_ack, exclusive)
    std::string consumer_tag = channel->BasicConsume("hello", "", /*no_local=*/true,
                                                     /*no_ack=*/true, /*exclusive=*/false);
    std::cout << "메시지 수신 대기 중... (Ctrl+C로 종료)" << std::endl;
    while (g_running) {
        AmqpClient::Envelope::ptr_t envelope;
        if (channel->BasicConsumeMessage(consumer_tag, envelope, 1000)) {
            std::cout << "Received: " << envelope->Message()->Body() << std::endl;
        }
    }
    return 0;
}

주의: no_ack=true면 서버가 메시지를 보내자마자 삭제합니다. 처리 실패 시 유실되므로, 실전에서는 no_ack=false로 ACK를 수동 처리해야 합니다.

rabbitmq-c로 직접 연결 (C API)

SimpleAmqpClient 없이 rabbitmq-c만 사용하는 예시입니다.

// rabbitmq_c_hello.cpp - rabbitmq-c C API
#include <amqp.h>
#include <amqp_tcp_socket.h>
#include <amqp_framing.h>
#include <iostream>
#include <cstring>
int main() {
    amqp_connection_state_t conn = amqp_new_connection();
    amqp_socket_t* socket = amqp_tcp_socket_new(conn);
    if (!socket) {
        std::cerr << "소켓 생성 실패" << std::endl;
        return 1;
    }
    if (amqp_socket_open(socket, "localhost", 5672) != AMQP_STATUS_OK) {
        std::cerr << "연결 실패" << std::endl;
        amqp_destroy_connection(conn);
        return 1;
    }
    amqp_rpc_reply_t reply = amqp_login(conn, "/", 0, 131072, 0,
                                        AMQP_SASL_METHOD_PLAIN, "guest", "guest");
    if (reply.reply_type != AMQP_RESPONSE_NORMAL) {
        std::cerr << "로그인 실패" << std::endl;
        amqp_destroy_connection(conn);
        return 1;
    }
    amqp_channel_open(conn, 1);
    reply = amqp_get_rpc_reply(conn);
    if (reply.reply_type != AMQP_RESPONSE_NORMAL) {
        std::cerr << "채널 열기 실패" << std::endl;
        amqp_connection_close(conn, AMQP_REPLY_SUCCESS);
        amqp_destroy_connection(conn);
        return 1;
    }
    // 큐 선언
    amqp_queue_declare_ok_t* r = amqp_queue_declare(conn, 1, amqp_cstring_bytes("hello"),
                                                    0, 0, 0, 0, amqp_empty_table);
    reply = amqp_get_rpc_reply(conn);
    if (reply.reply_type != AMQP_RESPONSE_NORMAL) {
        std::cerr << "큐 선언 실패" << std::endl;
        amqp_connection_close(conn, AMQP_REPLY_SUCCESS);
        amqp_destroy_connection(conn);
        return 1;
    }
    // 메시지 발행
    std::string body = "Hello from rabbitmq-c!";
    amqp_basic_properties_t props;
    props._flags = AMQP_BASIC_CONTENT_TYPE_FLAG | AMQP_BASIC_DELIVERY_MODE_FLAG;
    props.content_type = amqp_cstring_bytes("text/plain");
    props.delivery_mode = 1;  // non-persistent
    amqp_basic_publish(conn, 1, amqp_cstring_bytes(""), amqp_cstring_bytes("hello"),
                       0, 0, &props, amqp_cstring_bytes(body.c_str()));
    std::cout << "메시지 발행 완료" << std::endl;
    amqp_channel_close(conn, 1, AMQP_REPLY_SUCCESS);
    amqp_connection_close(conn, AMQP_REPLY_SUCCESS);
    amqp_destroy_connection(conn);
    return 0;
}

컴파일:

g++ -std=c++17 rabbitmq_c_hello.cpp -o rabbitmq_c_hello -lrabbitmq

작업 큐·prefetch·NAK·RAII 래퍼 예제

작업 큐 (Task Queue) — durable + persistent

실전에서는 큐와 메시지를 영구 저장해 브로커 재시작 시에도 유지합니다.

// task_producer.cpp
#include <SimpleAmqpClient/SimpleAmqpClient.h>
#include <iostream>
#include <string>
int main(int argc, char* argv[]) {
    auto channel = AmqpClient::Channel::Create("localhost", 5672, "guest", "guest");
    // durable=true: 큐를 디스크에 저장
    channel->DeclareQueue("task_queue", false, true, false, false);
    std::string message = (argc > 1) ? argv[1] : "Hello World!";
    auto msg = AmqpClient::BasicMessage::Create(message);
    msg->DeliveryMode(2);  // 2 = persistent (디스크에 저장)
    channel->BasicPublish("", "task_queue", msg);
    std::cout << " [x] Sent '" << message << "'" << std::endl;
    return 0;
}

작업 큐 컨슈머 — prefetch + ACK

// task_consumer.cpp
#include <SimpleAmqpClient/SimpleAmqpClient.h>
#include <algorithm>
#include <chrono>
#include <csignal>
#include <iostream>
#include <thread>
static volatile sig_atomic_t g_running = 1;
void sig_handler(int) { g_running = 0; }
int main() {
    signal(SIGINT, sig_handler);
    auto channel = AmqpClient::Channel::Create("localhost", 5672, "guest", "guest");
    channel->DeclareQueue("task_queue", false, true, false, false);
    // prefetch_count=1: 한 번에 하나씩만 가져옴 → 워커 간 공정 분배
    channel->BasicQos(0, 1, false);
    std::string consumer_tag = channel->BasicConsume("task_queue", "",
                                                     /*no_local=*/true,
                                                     /*no_ack=*/false,  // 수동 ACK
                                                     /*exclusive=*/false);
    std::cout << " [*] Waiting for messages. Ctrl+C to exit." << std::endl;
    while (g_running) {
        AmqpClient::Envelope::ptr_t envelope;
        if (channel->BasicConsumeMessage(consumer_tag, envelope, 1000)) {
            std::string body = envelope->Message()->Body();
            std::cout << " [x] Received '" << body << "'" << std::endl;
            // 작업 시뮬레이션 (점 하나당 1초)
            size_t count = std::count(body.begin(), body.end(), '.');
            std::this_thread::sleep_for(std::chrono::seconds(count > 0 ? count : 1));
            std::cout << " [x] Done" << std::endl;
            channel->BasicAck(envelope);  // 처리 완료 후 ACK
        }
    }
    return 0;
}

BasicQos(0, 1, false)로 prefetch를 1로 두면 워커는 ACK하지 않은 메시지를 하나만 가질 수 있어, 바쁜 워커에 메시지가 쌓이지 않고 비어 있는 워커로 갑니다. no_ack=false이므로 처리를 마친 뒤 BasicAck를 보내야 브로커가 메시지를 지웁니다. ACK 전에 연결이 끊기면 브로커가 그 메시지를 다른 컨슈머에게 다시 전달합니다. SimpleAmqpClient의 BasicConsume에도 prefetch 인자(기본값 1)가 있어 같은 역할을 합니다.

실패 시 NAK (재큐잉)

// 처리 실패 시 BasicReject로 재큐잉
void process_message(AmqpClient::Channel::ptr_t channel,
                    AmqpClient::Envelope::ptr_t envelope) {
    std::string body = envelope->Message()->Body();
    try {
        do_work(body);
        channel->BasicAck(envelope);  // 성공
    } catch (const std::exception& e) {
        std::cerr << "처리 실패: " << e.what() << std::endl;
        // requeue=true: 다시 큐에 넣음. 계속 실패하는 메시지는 무한히 재전달되므로
        // 재시도 횟수를 세거나 requeue=false + DLX로 격리하는 편이 안전함
        channel->BasicReject(envelope, /*requeue=*/true);
    }
}

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

// rabbitmq_producer.hpp
#pragma once
#include <SimpleAmqpClient/SimpleAmqpClient.h>
#include <memory>
#include <stdexcept>
#include <string>
class RabbitMQProducer {
public:
    RabbitMQProducer(const std::string& host = "localhost",
                    int port = 5672,
                    const std::string& user = "guest",
                    const std::string& password = "guest")
        : host_(host), user_(user), password_(password), port_(port) {}
    void connect() {
        channel_ = AmqpClient::Channel::Create(host_, port_, user_, password_);
    }
    void declareQueue(const std::string& queue, bool durable = true) {
        if (!channel_) throw std::runtime_error("Not connected");
        channel_->DeclareQueue(queue, false, durable, false, false);
    }
    void publish(const std::string& queue, const std::string& message,
                bool persistent = true) {
        if (!channel_) throw std::runtime_error("Not connected");
        auto msg = AmqpClient::BasicMessage::Create(message);
        msg->DeliveryMode(persistent ? 2 : 1);
        channel_->BasicPublish("", queue, msg);
    }
private:
    std::string host_, user_, password_;
    int port_;
    AmqpClient::Channel::ptr_t channel_;
};
// rabbitmq_consumer.hpp
#pragma once
#include <SimpleAmqpClient/SimpleAmqpClient.h>
#include <atomic>
#include <functional>
#include <stdexcept>
#include <string>
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), user_(user), password_(password), port_(port) {}
    void connect() {
        channel_ = AmqpClient::Channel::Create(host_, port_, user_, password_);
    }
    void declareQueue(const std::string& queue, bool durable = true) {
        if (!channel_) throw std::runtime_error("Not connected");
        channel_->DeclareQueue(queue, false, durable, false, false);
    }
    void setPrefetch(int count) {
        if (!channel_) throw std::runtime_error("Not connected");
        channel_->BasicQos(0, count, false);
    }
    void consume(const std::string& queue, MessageHandler handler) {
        if (!channel_) throw std::runtime_error("Not connected");
        std::string consumer_tag = channel_->BasicConsume(queue, "", /*no_local=*/true,
                                                          /*no_ack=*/false, /*exclusive=*/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};
};

실전 예제: 주문 이벤트 파이프라인

// order_pipeline.cpp
#include "rabbitmq_producer.hpp"
#include "rabbitmq_consumer.hpp"
#include <iostream>
void run_producer() {
    RabbitMQProducer producer("localhost", 5672, "guest", "guest");
    producer.connect();
    producer.declareQueue("orders");
    for (int i = 0; i < 10; ++i) {
        std::string order = R"({"order_id":)" + std::to_string(i) +
                           R"(,"user_id":"u1","amount":100})";
        producer.publish("orders", order);
        std::cout << "Published order " << i << std::endl;
    }
}
void run_consumer() {
    RabbitMQConsumer consumer("localhost", 5672, "guest", "guest");
    consumer.connect();
    consumer.declareQueue("orders");
    consumer.setPrefetch(1);
    consumer.consume("orders", [](const std::string& body) -> bool {
        std::cout << "Processing: " << body << std::endl;
        // 주문 검증, 재고 차감 등
        return true;  // 성공 시 ACK
    });
}
int main(int argc, char* argv[]) {
    if (argc > 1 && std::string(argv[1]) == "producer") {
        run_producer();
    } else {
        run_consumer();
    }
    return 0;
}

토픽 라우팅·우선순위 큐·Dead Letter·메시지 TTL

토픽 라우팅 (Topic Exchange)

로그 레벨별·서비스별로 메시지를 라우팅할 때 topic exchange를 사용합니다. routing key는 점으로 구분한 단어 목록이고, 바인딩 키에서 *는 정확히 한 단어, #는 0개 이상의 단어와 일치합니다. 그래서 logs.error라는 바인딩은 logs.error.critical과 일치하지 않고, logs.error.*나 logs.#여야 일치합니다.

flowchart LR
  subgraph Producer[프로듀서]
    P[BasicPublish]
  end
  subgraph Exchange["logs (topic)"]
    E[Exchange]
  end
  subgraph Queues["큐 (바인딩 키)"]
    Q1["logs.error.*"]
    Q2["logs.*.critical"]
    Q3["logs.#"]
  end
  P -->|logs.error.critical| E
  E --> Q1
  E --> Q2
  E --> Q3

SimpleAmqpClient로 토픽 라우팅:

// topic_producer.cpp
#include <SimpleAmqpClient/SimpleAmqpClient.h>
#include <iostream>
int main(int argc, char* argv[]) {
    auto channel = AmqpClient::Channel::Create("localhost", 5672, "guest", "guest");
    // topic 타입 exchange 선언
    channel->DeclareExchange("logs_topic", "topic", false, false, false);
    std::string routing_key = (argc > 1) ? argv[1] : "logs.info.default";
    std::string message = (argc > 2) ? argv[2] : "Hello from topic!";
    auto msg = AmqpClient::BasicMessage::Create(message);
    msg->DeliveryMode(2);
    channel->BasicPublish("logs_topic", routing_key, msg);
    std::cout << " [x] Sent '" << routing_key << "':" << message << std::endl;
    return 0;
}
// topic_consumer.cpp - 특정 패턴만 구독
#include <SimpleAmqpClient/SimpleAmqpClient.h>
#include <iostream>
#include <csignal>
static volatile sig_atomic_t g_running = 1;
void sig_handler(int) { g_running = 0; }
int main(int argc, char* argv[]) {
    std::string binding_key = (argc > 1) ? argv[1] : "logs.#";  // 모든 로그
    signal(SIGINT, sig_handler);
    auto channel = AmqpClient::Channel::Create("localhost", 5672, "guest", "guest");
    channel->DeclareExchange("logs_topic", "topic", false, false, false);
    // 익명 큐 생성 (exclusive, 연결 끊기면 삭제)
    std::string queue = channel->DeclareQueue("", false, false, true, true);
    channel->BindQueue(queue, "logs_topic", binding_key);
    std::string consumer_tag = channel->BasicConsume(queue, "", /*no_local=*/true,
                                                     /*no_ack=*/false, /*exclusive=*/false);
    std::cout << " [*] Waiting for " << binding_key << ". Ctrl+C to exit." << std::endl;
    while (g_running) {
        AmqpClient::Envelope::ptr_t envelope;
        if (channel->BasicConsumeMessage(consumer_tag, envelope, 1000)) {
            std::cout << " [x] " << envelope->RoutingKey() << ": "
                      << envelope->Message()->Body() << std::endl;
            channel->BasicAck(envelope);
        }
    }
    return 0;
}

rabbitmq-c로 토픽 exchange:

// topic_rabbitmq_c.cpp (에러 검사 생략: 실제로는 각 호출 뒤 amqp_get_rpc_reply 확인)
#include <amqp.h>
#include <amqp_tcp_socket.h>
#include <amqp_framing.h>
#include <iostream>
int main() {
    amqp_connection_state_t conn = amqp_new_connection();
    amqp_socket_t* socket = amqp_tcp_socket_new(conn);
    amqp_socket_open(socket, "localhost", 5672);
    amqp_login(conn, "/", 0, 131072, 0, AMQP_SASL_METHOD_PLAIN, "guest", "guest");
    amqp_channel_open(conn, 1);
    // topic exchange 선언
    amqp_exchange_declare(conn, 1, amqp_cstring_bytes("logs_topic"),
                          amqp_cstring_bytes("topic"), 0, 0, 0, 0, amqp_empty_table);
    amqp_get_rpc_reply(conn);
    // 메시지 발행
    amqp_basic_properties_t props;
    props._flags = AMQP_BASIC_DELIVERY_MODE_FLAG;
    props.delivery_mode = 2;
    amqp_basic_publish(conn, 1, amqp_cstring_bytes("logs_topic"),
                       amqp_cstring_bytes("logs.error.critical"),
                       0, 0, &props, amqp_cstring_bytes("Critical error!"));
    amqp_channel_close(conn, 1, AMQP_REPLY_SUCCESS);
    amqp_connection_close(conn, AMQP_REPLY_SUCCESS);
    amqp_destroy_connection(conn);
    return 0;
}

우선순위 큐 (Priority Queue)

긴급 주문·VIP 고객 메시지를 먼저 처리하려면 x-max-priority 인자로 우선순위 큐를 선언합니다(1~10 정도가 권장 범위). 큐 인자는 선언 시에만 정할 수 있습니다. 아래는 rabbitmq-c 예시입니다(지정 초기화 .num_entries =는 C++20 문법).

// priority_queue_rabbitmq_c.cpp - x-max-priority (0~10)
amqp_table_entry_t entries[1];
entries[0].key = amqp_cstring_bytes("x-max-priority");
entries[0].value.kind = AMQP_FIELD_KIND_I32;
entries[0].value.value.i32 = 10;
amqp_table_t args = {.num_entries = 1, .entries = entries};
amqp_queue_declare(conn, 1, amqp_cstring_bytes("priority_queue"), 0, 1, 0, 0, args);
// 우선순위 9 메시지 발행 (높을수록 먼저 처리)
amqp_basic_properties_t props;
props._flags = AMQP_BASIC_DELIVERY_MODE_FLAG | AMQP_BASIC_PRIORITY_FLAG;
props.delivery_mode = 2;
props.priority = 9;
amqp_basic_publish(conn, 1, amqp_cstring_bytes(""), amqp_cstring_bytes("priority_queue"),
                   0, 0, &props, amqp_cstring_bytes("Urgent task!"));

SimpleAmqpClient도 DeclareQueue에 인자 테이블을 받는 오버로드가 있어 같은 설정을 할 수 있습니다.

AmqpClient::Table args;
args.insert(AmqpClient::TableEntry(AmqpClient::TableKey("x-max-priority"),
                                   AmqpClient::TableValue(static_cast<int32_t>(10))));
channel->DeclareQueue("priority_queue", /*passive=*/false, /*durable=*/true,
                      /*exclusive=*/false, /*auto_delete=*/false, args);
auto msg = AmqpClient::BasicMessage::Create("Urgent task!");
msg->Priority(9);
channel->BasicPublish("", "priority_queue", msg);

우선순위는 큐에 메시지가 쌓여 있을 때만 의미가 있습니다. prefetch가 크면 메시지가 이미 컨슈머 쪽으로 넘어가 있어 우선순위가 반영되지 않으므로, 우선순위 큐에는 작은 prefetch를 함께 씁니다.

Dead Letter Exchange (DLX)

처리 실패·만료·거부된 메시지를 별도 큐로 보내 재처리·분석할 때 사용합니다.

flowchart TB
  subgraph Main[메인 큐]
    Q1[orders]
  end
  subgraph DLX[Dead Letter]
    E[dlx_exchange]
    Q2[orders.dlq]
  end
  Q1 -->|reject/expire/maxlen| E
  E --> Q2
// dlx_rabbitmq_c.cpp - DLX용 exchange + 메인 큐 인자
amqp_exchange_declare(conn, 1, amqp_cstring_bytes("dlx_exchange"),
                      amqp_cstring_bytes("direct"), 0, 1, 0, 0, amqp_empty_table);
amqp_queue_declare(conn, 1, amqp_cstring_bytes("orders.dlq"), 0, 1, 0, 0, amqp_empty_table);
amqp_queue_bind(conn, 1, amqp_cstring_bytes("orders.dlq"),
                amqp_cstring_bytes("dlx_exchange"), amqp_cstring_bytes("dlq"), amqp_empty_table);
amqp_table_entry_t entries[2];
entries[0].key = amqp_cstring_bytes("x-dead-letter-exchange");
entries[0].value.kind = AMQP_FIELD_KIND_UTF8;
entries[0].value.value.bytes = amqp_cstring_bytes("dlx_exchange");
entries[1].key = amqp_cstring_bytes("x-dead-letter-routing-key");
entries[1].value.kind = AMQP_FIELD_KIND_UTF8;
entries[1].value.value.bytes = amqp_cstring_bytes("dlq");
amqp_table_t args = {.num_entries = 2, .entries = entries};
amqp_queue_declare(conn, 1, amqp_cstring_bytes("orders"), 0, 1, 0, 0, args);

DLX로 가는 경우: BasicReject(requeue=false), 메시지 TTL 만료, 큐 x-max-length 초과.

메시지 TTL (Time-To-Live)

큐 단위 또는 메시지 단위로 TTL을 설정해 오래된 메시지를 자동 삭제합니다.

// ttl_rabbitmq_c.cpp - x-message-ttl: 60000ms (60초)
amqp_table_entry_t entries[1];
entries[0].key = amqp_cstring_bytes("x-message-ttl");
entries[0].value.kind = AMQP_FIELD_KIND_I32;
entries[0].value.value.i32 = 60000;
amqp_table_t args = {.num_entries = 1, .entries = entries};
amqp_queue_declare(conn, 1, amqp_cstring_bytes("ttl_queue"), 0, 1, 0, 0, args);

큐 TTL은 만료된 메시지를 바로 제거하지만, 메시지별 TTL은 그 메시지가 큐의 맨 앞에 왔을 때 만료 여부를 확인합니다. 그래서 TTL이 긴 메시지 뒤에 있는 짧은 TTL 메시지는 제때 사라지지 않을 수 있습니다.

메시지별 TTL (SimpleAmqpClient):

// 메시지 단위 TTL - Expiration 헤더 (밀리초 문자열)
auto msg = AmqpClient::BasicMessage::Create("Short-lived");
msg->Expiration("5000");  // 5초 후 만료
channel->BasicPublish("", "queue", msg);

Connection refused, PRECONDITION_FAILED, DLX 무한 루프: 에러 해결

”Connection refused” / “Connection failed”

Channel::Create()가 예외를 던진다면 브로커가 실행 중이 아니거나, 호스트·포트가 틀렸거나, 방화벽에 막혔거나, Docker 환경에서 localhost가 다른 컨테이너를 가리키지 않는 경우가 대부분입니다.

# RabbitMQ 실행 확인
docker ps | grep rabbitmq
# Docker 컨테이너에서 호스트 RabbitMQ 접속 시
# host.docker.internal:5672 (macOS/Windows)
# 172.17.0.1:5672 (Linux)
// ✅ 연결 전 브로커 확인
try {
    auto channel = AmqpClient::Channel::Create("localhost", 5672, "guest", "guest");
} catch (const std::exception& e) {  // 버전에 따라 AmqpLibraryException 등 구체 타입이 다름
    std::cerr << "연결 실패: " << e.what() << std::endl;
    // 재시도 또는 폴백
}

”NOT_FOUND”와 조용히 사라지는 메시지

없는 큐에서 BasicConsume하면 NOT_FOUND로 채널이 닫힙니다. 반면 기본 exchange("")로 없는 큐 이름에 발행하면 에러가 나지 않고 메시지가 라우팅되지 못해 조용히 버려집니다(mandatory 플래그를 켜야 반송됨). 그래서 프로듀서와 컨슈머 양쪽에서 같은 옵션으로 큐를 선언해 두는 것이 안전합니다.

// ❌ 큐 선언 없이 발행: my_queue가 없으면 메시지가 조용히 버려짐
channel->BasicPublish("", "my_queue", msg);
// ✅ 올바른 코드: 선언 후 사용
channel->DeclareQueue("my_queue", false, true, false, false);
channel->BasicPublish("", "my_queue", msg);

메시지 유실 (브로커 재시작 시)

RabbitMQ 재시작 후 메시지가 사라졌다면 큐가 durable=false이거나 메시지가 DeliveryMode(1)(non-persistent)인 경우입니다. 둘 다 설정해야 재시작 후에도 남습니다. 다만 persistent 메시지도 디스크에 기록되기 전에 브로커가 죽으면 잃을 수 있으므로, 유실이 치명적이면 publisher confirm까지 써야 합니다.

// ✅ durable 큐 + persistent 메시지
channel->DeclareQueue("task_queue", false, true, false, false);  // durable=true
auto msg = AmqpClient::BasicMessage::Create(body);
msg->DeliveryMode(2);  // 2 = persistent
channel->BasicPublish("", "task_queue", msg);

메시지 중복 처리

처리 중 크래시하거나 ACK를 보내기 직전에 연결이 끊기면 브로커는 그 메시지를 다시 전달합니다. 처리 완료 후에만 ACK를 보내면 유실은 막을 수 있지만 중복은 막을 수 없으므로(at-least-once), 처리 로직을 멱등하게 만들어야 합니다. 반대로 처리 전에 ACK하면 중복은 없지만 크래시 시 유실됩니다.

// ✅ 처리 완료 후에만 ACK
void handle_message(AmqpClient::Envelope::ptr_t envelope) {
    std::string body = envelope->Message()->Body();
    process(body);           // 비즈니스 로직
    channel->BasicAck(envelope);  // 처리 성공 후 ACK
}
// 멱등성: 처리 로직이 중복 호출되어도 안전하게 (예: DB upsert)

수동 ACK와 관련해 SimpleAmqpClient에서 자주 나는 실수가 두 가지 있습니다. 첫째, BasicConsume의 no_ack 기본값이 true라서 BasicConsume(queue)처럼 인자를 생략하면 자동 ACK 모드가 됩니다. 이때는 코드에서 BasicAck를 불러도 의미가 없고 컨슈머가 죽는 순간 받아 둔 메시지가 사라집니다. 세 번째 인자(no_local)를 no_ack로 착각하기도 쉬우므로 위 예제처럼 인자 이름을 주석으로 적어 두는 편이 안전합니다. 둘째, 예외 처리 경로처럼 어떤 분기에서 ACK도 reject도 하지 않으면 그 메시지는 관리 UI에서 “Unacked”로 남아 prefetch 슬롯을 계속 차지합니다. prefetch=1이면 컨슈머가 더 이상 아무것도 받지 못해 멈춘 것처럼 보이고, 연결이 끊겨야 다른 컨슈머에게 재전달됩니다. 모든 경로에서 ACK나 reject가 정확히 한 번 호출되는지 확인해야 합니다.

멱등성 기록을 메모리 집합에 두면 컨슈머 재시작 직후(중복이 가장 많이 생기는 때)에 기록이 사라지고, 컨슈머 여러 대가 서로의 기록을 모릅니다. 비즈니스 쓰기와 같은 DB 트랜잭션 안에서 processed_messages(message_id PRIMARY KEY) 테이블에 함께 INSERT하고, 중복 키 에러가 나면 이미 처리된 것으로 보고 ACK하는 방식이 가장 단단합니다. 두 쓰기가 한 트랜잭션에 묶이므로 둘 중 하나만 반영되는 틈이 없습니다.

”PRECONDITION_FAILED” — 큐 선언 불일치

이미 존재하는 큐와 다른 옵션(durable, exclusive, 인자 등)으로 선언하면 PRECONDITION_FAILED로 채널이 닫힙니다.

// 기존 큐와 동일한 옵션으로 선언해야 함
// 또는 passive=true로 존재 확인만
channel->DeclareQueue("task_queue", true, false, false, false);  // passive: 존재만 확인

”ACCESS_REFUSED” — 인증 실패

잘못된 사용자명·비밀번호이거나, 원격 호스트에서 guest 계정으로 접속한 경우입니다. guest는 기본 설정에서 localhost 접속만 허용됩니다.

# RabbitMQ에서 새 사용자 생성
rabbitmqctl add_user myuser mypassword
rabbitmqctl set_permissions -p / myuser ".*" ".*" ".*"
// ✅ 올바른 인증 정보
auto channel = AmqpClient::Channel::Create("localhost", 5672, "myuser", "mypassword");

워커 간 공정 분배 실패

prefetch 제한이 없으면 브로커는 큐의 메시지를 컨슈머들에게 라운드 로빈으로 미리 밀어 넣으므로, 처리 속도와 관계없이 느린 워커에도 같은 양이 쌓입니다. prefetch로 미확인 메시지 수를 제한합니다.

// ✅ prefetch=1로 한 번에 하나씩만 가져와 공정 분배
channel->BasicQos(0, 1, false);

SimpleAmqpClient 링크 에러 (dyld, .so)

dyld: Library not loaded나 공유 라이브러리를 찾지 못하는 에러는 실행 시 라이브러리 경로를 알려 주지 않았기 때문입니다.

# macOS에서 rpath 추가
c++ -std=c++17 main.cpp -Wl,-rpath,/usr/local/lib \
  -lSimpleAmqpClient -lrabbitmq -o app
# Linux: LD_LIBRARY_PATH
export LD_LIBRARY_PATH=/usr/local/lib:$LD_LIBRARY_PATH
./app

”Connection closed” — 장시간 유지 연결 끊김

한동안 쓰지 않은 연결로 발행할 때 실패한다면 heartbeat 타임아웃이나 중간 장비의 유휴 연결 정리로 연결이 끊긴 경우입니다. SimpleAmqpClient는 동기 라이브러리라 백그라운드에서 heartbeat를 보내지 않으므로, 연결이 끊겼다는 것을 다음 호출의 예외로 알게 됩니다. 예외가 나면 채널을 새로 만들고 한 번 재시도합니다.

// ✅ 실패 시 재연결 후 한 번 재시도 (단일 스레드 가정)
void publish_with_reconnect(AmqpClient::Channel::ptr_t& ch, const std::string& queue,
                            const AmqpClient::BasicMessage::ptr_t& msg) {
    try {
        ch->BasicPublish("", queue, msg);
    } catch (const std::exception&) {
        ch = AmqpClient::Channel::Create("localhost", 5672, "guest", "guest");
        ch->BasicPublish("", queue, msg);
    }
}

토픽 exchange “NOT_FOUND”

선언하지 않은 이름 있는 exchange에 발행하면 NOT_FOUND로 채널이 닫힙니다. 기본 exchange("")가 아닌 exchange는 DeclareExchange로 먼저 선언합니다.

// ✅ exchange 선언 후 발행
channel->DeclareExchange("logs_topic", "topic", false, false, false);
channel->BasicPublish("logs_topic", "app.error", msg);

우선순위 큐 “PRECONDITION_FAILED”

이미 존재하는 큐를 다른 x-max-priority 값(또는 없이)으로 다시 선언하면 PRECONDITION_FAILED가 납니다. 큐 인자는 바꿀 수 없으므로 큐를 지우고 다시 만들어야 합니다.

# 기존 큐 삭제 후 재선언 (개발 환경)
rabbitmqctl delete_queue priority_queue

DLX 무한 루프

DLQ에서 reject한 메시지가 다시 메인 큐로 돌아가는 경우는 DLQ에도 메인 큐로 향하는 DLX를 설정했기 때문입니다. DLQ에는 DLX를 설정하지 않거나, 별도 “final” DLQ로 이중 전달 구조를 사용합니다.


멱등성·Correlation ID·네이밍 규칙

멱등성 설계

메시지가 중복 전달되어도 안전하게 처리합니다.

// ✅ 멱등성: order_id로 중복 체크
bool process_order(const std::string& body) {
    auto order_id = parse_order_id(body);
    if (db->order_already_processed(order_id)) {
        return true;  // 이미 처리됨 → ACK
    }
    db->process_and_mark(order_id, body);
    return true;
}

메시지 크기 제한

너무 큰 메시지는 브로커 메모리와 네트워크에 부담을 줍니다. 프레임 크기(frame_max, 기본 128KB)는 메시지를 나눠 보내는 단위일 뿐이고, 메시지 자체의 최대 크기는 max_message_size 설정(최근 버전 기본 128MB)으로 제한됩니다. 그래도 수 MB 이상의 데이터는 외부 저장소에 두고 참조만 보내는 편이 낫습니다.

constexpr size_t MAX_MSG_SIZE = 1024 * 1024;  // 서비스 정책에 맞게 정하는 값 (예: 1MB)
if (message.size() > MAX_MSG_SIZE) {
    // 대용량은 외부 스토리지 URL만 전달
    message = R"({"ref":"s3://bucket/large-file.json"})";
}

Correlation ID로 추적

분산 환경에서 요청-응답을 연결할 때 correlation_id 헤더를 사용합니다.

auto msg = AmqpClient::BasicMessage::Create(body);
msg->CorrelationId("req-" + std::to_string(request_id));
channel->BasicPublish("", "orders", msg);

큐·Exchange 네이밍 규칙

권장: {서비스}.{도메인}.{용도}
예: order.payment.events, notification.email.tasks

연결 실패 시 지수 백오프

int retry_delay_ms = 100;
for (int attempt = 0; attempt < 5; ++attempt) {
    try {
        channel_ = AmqpClient::Channel::Create(host_, port_, user_, password_);
        break;
    } catch (const std::exception& e) {
        if (attempt == 4) throw;
        std::this_thread::sleep_for(std::chrono::milliseconds(retry_delay_ms));
        retry_delay_ms *= 2;
    }
}

채널 재사용·Prefetch·Publisher Confirms

채널 재사용

매 메시지마다 새 채널을 만들지 마세요. 채널은 스레드당 하나 또는 공유 풀로 재사용합니다.

// ❌ 나쁜 예
void send_message(const std::string& msg) {
    auto channel = AmqpClient::Channel::Create("localhost", 5672, "guest", "guest");
    channel->BasicPublish("", "queue", AmqpClient::BasicMessage::Create(msg));
}
// ✅ 좋은 예: 채널 재사용
class Producer {
    AmqpClient::Channel::ptr_t channel_;
public:
    Producer() {
        channel_ = AmqpClient::Channel::Create("localhost", 5672, "guest", "guest");
        channel_->DeclareQueue("queue", false, true, false, false);
    }
    void send(const std::string& msg) {
        channel_->BasicPublish("", "queue", AmqpClient::BasicMessage::Create(msg));
    }
};

Prefetch 조정

처리량이 높을 때 prefetch를 늘리면 왕복 횟수를 줄일 수 있습니다. 하지만 너무 크면 한 워커에 메시지가 몰립니다.

// 처리 속도가 빠른 작업: prefetch 증가
channel->BasicQos(0, 10, false);  // 한 번에 10개까지
// 처리 속도가 느린 작업: prefetch=1 유지
channel->BasicQos(0, 1, false);

Publisher Confirms (신뢰성)

persistent 메시지는 브로커가 메시지를 받은 뒤의 재시작에 대비하는 것이고, 브로커에 도착하기 전이나 디스크에 기록되기 전의 유실은 막지 못합니다. 이를 확인하는 장치가 publisher confirm으로, 브로커가 메시지를 안전하게 받았을 때 프로듀서에 ack를 보냅니다. rabbitmq-c에서는 amqp_confirm_select로 채널을 confirm 모드로 바꾸고 basic.ack를 기다립니다. SimpleAmqpClient는 버전에 따라 BasicPublish가 내부적으로 확인을 기다리는 방식으로 동작하므로, 사용하는 버전의 문서를 확인합니다. 비동기로 대량의 confirm을 다뤄야 한다면 AMQP-CPP 같은 비동기 클라이언트가 편합니다.

메시지 크기 최적화

큰 메시지는 압축 후 발행하는 것을 고려하세요. AMQP 메시지 본문은 바이너리이므로 압축 결과를 그대로 넣고, content_encoding 속성에 압축 방식을 적어 컨슈머가 알 수 있게 합니다.

연결 풀 (다중 스레드)

AMQP 채널과 SimpleAmqpClient의 Channel 객체는 스레드 안전하지 않으므로, 여러 스레드가 발행한다면 스레드마다 채널을 따로 두거나 락으로 보호합니다. 아래는 미리 만든 채널을 돌려 쓰는 단순한 풀이지만, 같은 채널을 두 스레드가 동시에 받지 않도록 하려면 반납 방식이 필요합니다.

class ChannelPool {
    std::vector<AmqpClient::Channel::ptr_t> channels_;
    std::mutex mtx_;
    size_t index_ = 0;
public:
    AmqpClient::Channel::ptr_t get() {
        std::lock_guard<std::mutex> lock(mtx_);
        auto& ch = channels_[index_ % channels_.size()];
        index_++;
        return ch;
    }
};

Graceful Shutdown·Health Check·Docker Compose

Graceful Shutdown

static std::atomic<bool> g_running{true};
void sig_handler(int) {
    g_running = false;
}
int main() {
    signal(SIGINT, sig_handler);
    signal(SIGTERM, sig_handler);
    auto channel = AmqpClient::Channel::Create("localhost", 5672, "guest", "guest");
    std::string consumer_tag = channel->BasicConsume("queue", "", /*no_local=*/true, /*no_ack=*/false, /*exclusive=*/false);
    while (g_running) {
        AmqpClient::Envelope::ptr_t envelope;
        if (channel->BasicConsumeMessage(consumer_tag, envelope, 1000)) {
            process(envelope);
            channel->BasicAck(envelope);
        }
    }
    channel->BasicCancel(consumer_tag);
    return 0;
}

설정 외부화

struct RabbitMQConfig {
    std::string host = "localhost";
    int port = 5672;
    std::string user = "guest";
    std::string password = "guest";
};
RabbitMQConfig loadFromEnv() {
    RabbitMQConfig c;
    if (const char* h = std::getenv("RABBITMQ_HOST")) c.host = h;
    if (const char* p = std::getenv("RABBITMQ_PORT")) c.port = std::stoi(p);
    if (const char* u = std::getenv("RABBITMQ_USER")) c.user = u;
    if (const char* pw = std::getenv("RABBITMQ_PASSWORD")) c.password = pw;
    return c;
}

Health Check

// 연결·로그인·채널 열기가 성공하면 정상으로 판단 (매 호출마다 새 연결을 만드므로 주기를 길게)
bool check_rabbitmq(const std::string& host, int port) {
    try {
        auto channel = AmqpClient::Channel::Create(host, port, "guest", "guest");
        return static_cast<bool>(channel);
    } catch (...) {
        return false;
    }
}

재시도 (지수 백오프)

template<typename Func>
auto withRetry(Func&& f, int maxRetries = 3) -> decltype(f()) {
    for (int i = 0; i < maxRetries; ++i) {
        try {
            return f();
        } catch (const std::exception& e) {
            if (i == maxRetries - 1) throw;
            auto delay = std::chrono::milliseconds(100 * (1 << i));
            std::this_thread::sleep_for(delay);
        }
    }
    throw std::runtime_error("Unreachable");
}
// 사용
auto channel = withRetry([&]() {
    return AmqpClient::Channel::Create("localhost", 5672, "guest", "guest");
});

로깅·메트릭

template<typename Func>
auto withTiming(const char* op, Func&& f) -> decltype(f()) {
    auto start = std::chrono::steady_clock::now();
    auto result = f();
    auto dur = std::chrono::duration_cast<std::chrono::milliseconds>(
        std::chrono::steady_clock::now() - start).count();
    LOG(INFO) << "RabbitMQ " << op << " took " << dur << "ms";
    return result;
}

Docker Compose 연동

# docker-compose.yml
version: '3'
services:
  app:
    build: .
    environment:
      RABBITMQ_HOST: rabbitmq
      RABBITMQ_PORT: 5672
    depends_on:
      - rabbitmq
  rabbitmq:
    image: rabbitmq:3.12-management
    ports:
      - "5672:5672"
      - "15672:15672"

메시지 순서와 다른 클라이언트

한 큐를 한 컨슈머가 소비하면 메시지는 발행 순서대로 전달됩니다. 하지만 컨슈머가 여럿이면 메시지가 나뉘어 처리 순서가 섞이고, 하나의 컨슈머라도 reject 후 requeue한 메시지는 큐의 원래 위치 근처로 돌아가거나 다시 전달되면서 순서가 바뀔 수 있습니다. 엄격한 순서가 필요하면 키별로 큐를 나누거나 Kafka의 파티션처럼 키 단위 순서를 보장하는 구조를 고려합니다.

C++ 클라이언트는 세 가지가 흔히 쓰입니다. SimpleAmqpClient는 rabbitmq-c 위의 동기 C++ 래퍼라 가장 쉽게 시작할 수 있고, rabbitmq-c는 C API로 세밀하게 제어할 수 있습니다. AMQP-CPP는 C++ 네이티브 비동기 라이브러리로, 이벤트 루프(libev, libuv, Boost.Asio 등)와 통합해 많은 메시지를 비동기로 다룰 때 적합합니다.

다음 글: RabbitMQ 고급: 토픽 라우팅·우선순위·Dead Letter(#52-8) 이전 글: C++에서 Kafka 쓰기


참고 자료


같이 보면 좋은 글