C++ RabbitMQ 고급 기능 | Topic·Priority Queue·Dead Letter·TTL

들어가며: “단순 큐만으로는 부족해요”

RabbitMQ 기본(#52-7)에서 작업 큐를 구현했습니다. 하지만 큐 하나만으로는 금방 한계에 부딪힙니다. error, warning, info 로그를 한 큐에 넣으면 에러 알림 서비스와 로그 수집기가 같은 메시지를 나눠 가져가 버립니다. VIP 주문이 앞에 쌓인 일반 주문 뒤에서 기다리고, 형식이 잘못된 메시지를 requeue하면 같은 메시지가 계속 돌아와 컨슈머가 그것만 붙잡고 있게 됩니다. “세일 종료 1시간 전” 같은 알림은 늦게 처리되면 의미가 없어 버려야 하고, “주문 생성됨” 이벤트는 재고·결제·알림 서비스가 각자 받아야 합니다.

이 글은 이 문제들을 Exchange 라우팅, 우선순위 큐, Dead Letter Exchange, TTL로 해결하는 방법을 rabbitmq-c와 SimpleAmqpClient 코드로 다룹니다.

요구 환경: C++17 이상, RabbitMQ 3.8+, SimpleAmqpClient 또는 rabbitmq-c


실전 문제를 고급 기능에 매핑하기

문제별 기능 대응표

문제해결 기능핵심 개념
로그 레벨별 분기 (error/warning/info)토픽 Exchangerouting key 패턴 매칭
VIP·긴급 주문 우선 처리우선순위 큐x-max-priority
실패 메시지 무한 재시도Dead Letter Exchangex-dead-letter-exchange
오래된 메시지 자동 폐기메시지 TTLexpiration, x-message-ttl
한 이벤트를 여러 서비스가 구독Fanout Exchange브로드캐스트
특정 키만 받고 싶음Direct Exchange정확한 routing key 매칭

아키텍처 개요

flowchart TB
    subgraph Producer[프로듀서]
        P1[로그 error]
        P2[로그 warning]
        P3[주문 VIP]
        P4[주문 일반]
    end
    subgraph Exchange["Exchange (토픽/다이렉트)"]
        EX[logs.orders]
    end
    subgraph Queues[큐]
        Q1[error_queue]
        Q2[all_logs_queue]
        Q3[priority_queue]
        Q4[dlq_dead_letter]
    end
    subgraph Consumers[컨슈머]
        C1[에러 알림]
        C2[로그 수집]
        C3[주문 처리]
        C4[실패 분석]
    end
    P1 -->|routing: logs.error| EX
    P2 -->|routing: logs.warning| EX
    P3 -->|priority: 10| EX
    P4 -->|priority: 1| EX
    EX --> Q1
    EX --> Q2
    EX --> Q3
    Q3 -.->|실패 시| Q4
    Q1 --> C1
    Q2 --> C2
    Q3 --> C3
    Q4 --> C4

Exchange와 라우팅 개념

Exchange 타입 비교

타입동작사용 사례
directrouting_key가 정확히 일치하는 큐로 전달작업 유형별 분기 (email, sms)
topic패턴 매칭 (* 한 단어, # 0개 이상)로그 레벨·카테고리 (logs.error, order.vip)
fanout모든 바인딩된 큐로 브로드캐스트이벤트 복제 (주문 생성 → N개 서비스)
headers헤더 키-값 매칭복잡한 라우팅 (선택적 사용)

기본 Exchange vs 명시적 Exchange

// 기본 Exchange ("") 사용 시: routing_key = 큐 이름
// - 큐가 미리 존재해야 함
channel->BasicPublish("", "my_queue", msg);
// 명시적 Exchange 사용: Exchange → 바인딩된 큐로 라우팅
channel->BasicPublish("logs_topic", "logs.error", msg);

시퀀스: 토픽 라우팅 흐름

sequenceDiagram
    participant P as 프로듀서
    participant EX as logs_topic (Exchange)
    participant Q1 as error_queue
    participant Q2 as all_logs_queue
    participant C1 as 에러 알림
    participant C2 as 로그 수집
    P->>EX: BasicPublish("logs.error", body)
    EX->>EX: routing_key 매칭
    EX->>Q1: logs.error 바인딩 매칭
    EX->>Q2: logs.# 바인딩 매칭
    Q1->>C1: 메시지 전달
    Q2->>C2: 메시지 전달

토픽·Fanout Exchange 구현

토픽 Exchange 선언 (rabbitmq-c)

먼저 C 클라이언트인 rabbitmq-c로 Exchange와 큐, 바인딩을 선언합니다. rabbitmq-c 0.12부터는 헤더가 <rabbitmq-c/amqp.h>, <rabbitmq-c/tcp_socket.h>로 옮겨졌고 아래의 예전 경로는 deprecated 경고와 함께 계속 동작합니다.

// topic_exchange_setup.cpp - Exchange + 큐 바인딩
#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 (amqp_socket_open(socket, "localhost", 5672) != AMQP_STATUS_OK) {
        std::cerr << "연결 실패" << std::endl;
        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_destroy_connection(conn);
        return 1;
    }
    // 1. 토픽 Exchange 선언
    amqp_exchange_declare(conn, 1, amqp_cstring_bytes("logs_topic"),
                         amqp_cstring_bytes("topic"), 0, 1, 0, 0, amqp_empty_table);
    reply = amqp_get_rpc_reply(conn);
    if (reply.reply_type != AMQP_RESPONSE_NORMAL) {
        std::cerr << "Exchange 선언 실패" << std::endl;
        amqp_destroy_connection(conn);
        return 1;
    }
    // 2. 큐 선언
    amqp_queue_declare(conn, 1, amqp_cstring_bytes("error_logs"), 0, 1, 0, 0, amqp_empty_table);
    amqp_queue_declare(conn, 1, amqp_cstring_bytes("all_logs"), 0, 1, 0, 0, amqp_empty_table);
    reply = amqp_get_rpc_reply(conn);
    // 3. 바인딩: error_logs <- logs.error
    amqp_queue_bind(conn, 1, amqp_cstring_bytes("error_logs"),
                    amqp_cstring_bytes("logs_topic"),
                    amqp_cstring_bytes("logs.error"), amqp_empty_table);
    // 4. 바인딩: all_logs <- logs.# (모든 로그)
    amqp_queue_bind(conn, 1, amqp_cstring_bytes("all_logs"),
                    amqp_cstring_bytes("logs_topic"),
                    amqp_cstring_bytes("logs.#"), amqp_empty_table);
    reply = amqp_get_rpc_reply(conn);
    std::cout << "Exchange + 바인딩 설정 완료" << 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 topic_exchange_setup.cpp -o topic_setup -lrabbitmq

SimpleAmqpClient로 토픽 프로듀서

SimpleAmqpClient(rabbitmq-c 위의 C++ 래퍼)도 DeclareExchange, DeclareQueue, BindQueue를 제공하므로 같은 선언을 C++ 쪽에서 할 수 있습니다.

channel->DeclareExchange("logs_topic", AmqpClient::Channel::EXCHANGE_TYPE_TOPIC,
                         false /*passive*/, true /*durable*/, false /*auto_delete*/);
channel->DeclareQueue("error_logs", false, true, false, false);
channel->BindQueue("error_logs", "logs_topic", "logs.error");

아래 프로듀서는 Exchange가 이미 선언되어 있다고 가정하고 발행만 합니다.

// topic_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");
    // routing_key: logs.error, logs.warning, logs.info 등
    std::string level = (argc > 1) ? argv[1] : "info";
    std::string routing_key = "logs." + level;
    std::string body = (argc > 2) ? argv[2] : "Sample log message";
    auto msg = AmqpClient::BasicMessage::Create(body);
    msg->DeliveryMode(2);
    // Exchange "logs_topic"으로 발행, routing_key로 라우팅
    channel->BasicPublish("logs_topic", routing_key, msg);
    std::cout << " [x] Sent [" << routing_key << "] " << body << std::endl;
    return 0;
}

토픽 컨슈머 (에러 전용)

// topic_consumer_error.cpp - error_logs 큐만 소비
#include <SimpleAmqpClient/SimpleAmqpClient.h>
#include <csignal>
#include <iostream>
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");
    // 인자: queue, consumer_tag, no_local, no_ack, exclusive, prefetch_count
    std::string consumer_tag = channel->BasicConsume("error_logs", "", false, false, false, 1);
    std::cout << " [*] error_logs 큐 대기 중... Ctrl+C로 종료" << 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] ERROR: " << body << std::endl;
            // 에러 알림 (Slack, PagerDuty 등) 발송
            channel->BasicAck(envelope);
        }
    }
    channel->BasicCancel(consumer_tag);
    return 0;
}

Fanout Exchange: 브로드캐스트

한 이벤트를 여러 서비스가 각자 받을 때 사용합니다.

// fanout_setup.cpp - rabbitmq-c
// Exchange 선언 (fanout)
amqp_exchange_declare(conn, 1, amqp_cstring_bytes("order_events"),
                     amqp_cstring_bytes("fanout"), 0, 1, 0, 0, amqp_empty_table);
// 각 서비스별 큐 생성 및 바인딩 (routing_key 불필요)
amqp_queue_declare(conn, 1, amqp_cstring_bytes("inventory_queue"), 0, 1, 0, 0, amqp_empty_table);
amqp_queue_bind(conn, 1, amqp_cstring_bytes("inventory_queue"),
                amqp_cstring_bytes("order_events"), amqp_empty_bytes, amqp_empty_table);
amqp_queue_declare(conn, 1, amqp_cstring_bytes("notification_queue"), 0, 1, 0, 0, amqp_empty_table);
amqp_queue_bind(conn, 1, amqp_cstring_bytes("notification_queue"),
                amqp_cstring_bytes("order_events"), amqp_empty_bytes, amqp_empty_table);
// fanout_producer.cpp - SimpleAmqpClient
// order_events Exchange로 발행 (routing_key 무시됨)
channel->BasicPublish("order_events", "", msg);

우선순위 큐

우선순위 큐 선언

우선순위 큐는 x-max-priority 인자로 최대 우선순위 값을 설정합니다. 1255까지 지원되지만, 클래식 큐는 우선순위 값마다 내부 서브 큐를 따로 두므로 값이 클수록 메모리와 CPU를 더 씁니다. 공식 문서도 15 정도의 작은 값을 권장하며, 아래 예제는 10을 씁니다.

// priority_queue_setup.cpp - rabbitmq-c
#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);
    // x-max-priority: 1~255 지원, 작은 값 권장
    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;
    args.num_entries = 1;
    args.entries = entries;
    amqp_queue_declare(conn, 1, amqp_cstring_bytes("priority_orders"),
                       0, 1, 0, 0, args);
    amqp_rpc_reply_t reply = amqp_get_rpc_reply(conn);
    if (reply.reply_type != AMQP_RESPONSE_NORMAL) {
        std::cerr << "우선순위 큐 선언 실패" << std::endl;
        return 1;
    }
    std::cout << "우선순위 큐 선언 완료 (max-priority=10)" << std::endl;
    amqp_channel_close(conn, 1, AMQP_REPLY_SUCCESS);
    amqp_connection_close(conn, AMQP_REPLY_SUCCESS);
    amqp_destroy_connection(conn);
    return 0;
}

우선순위 메시지 발행

메시지에 priority 속성을 설정합니다. 숫자가 클수록 높은 우선순위.

// priority_producer.cpp - SimpleAmqpClient
#include <SimpleAmqpClient/SimpleAmqpClient.h>
#include <iostream>
int main() {
    auto channel = AmqpClient::Channel::Create("localhost", 5672, "guest", "guest");
    // VIP 주문 (우선순위 10)
    auto vip_msg = AmqpClient::BasicMessage::Create(R"({"order_id":1,"type":"vip"})");
    vip_msg->DeliveryMode(2);
    vip_msg->Priority(10);  // 최고 우선순위
    channel->BasicPublish("", "priority_orders", vip_msg);
    std::cout << " [x] VIP 주문 발행 (priority=10)" << std::endl;
    // 일반 주문 (우선순위 1)
    auto normal_msg = AmqpClient::BasicMessage::Create(R"({"order_id":2,"type":"normal"})");
    normal_msg->DeliveryMode(2);
    normal_msg->Priority(1);
    channel->BasicPublish("", "priority_orders", normal_msg);
    std::cout << " [x] 일반 주문 발행 (priority=1)" << std::endl;
    return 0;
}

x-max-priority보다 큰 우선순위로 발행한 메시지는 최댓값으로 취급되고, 우선순위를 주지 않은 메시지는 0으로 취급됩니다. rabbitmq-c를 쓴다면 amqp_basic_properties_t의 priority 필드와 플래그를 직접 설정합니다.

rabbitmq-c로 우선순위 발행

// priority_publish_rabbitmq_c.cpp
amqp_basic_properties_t props;
props._flags = AMQP_BASIC_CONTENT_TYPE_FLAG | AMQP_BASIC_DELIVERY_MODE_FLAG | AMQP_BASIC_PRIORITY_FLAG;
props.content_type = amqp_cstring_bytes("application/json");
props.delivery_mode = 2;
props.priority = 10;  // VIP
amqp_basic_publish(conn, 1, amqp_cstring_bytes(""), amqp_cstring_bytes("priority_orders"),
                   0, 0, &props, amqp_cstring_bytes(body.c_str()));

Dead Letter Exchange (DLX)

DLX 개념

메시지는 컨슈머가 basic.nack이나 basic.reject를 requeue=false로 호출했을 때, 메시지 TTL이 만료됐을 때, x-max-length 같은 길이 제한을 넘어 큐 앞쪽 메시지가 밀려났을 때(기본 overflow 동작인 drop-head), 쿼럼 큐에서 재전달 횟수 제한(delivery-limit)을 넘었을 때 Dead Letter Exchange로 보내집니다.

flowchart LR
    subgraph Main[메인 큐]
        Q1[order_queue]
    end
    subgraph DLX[Dead Letter]
        EX[dlx_exchange]
        Q2[dead_letter_queue]
    end
    Q1 -->|실패/만료/초과| EX
    EX --> Q2

DLX 설정이 있는 큐 선언

// dlx_setup.cpp - rabbitmq-c
#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);
    // 1. DLX용 Exchange 선언
    amqp_exchange_declare(conn, 1, amqp_cstring_bytes("dlx_exchange"),
                         amqp_cstring_bytes("direct"), 0, 1, 0, 0, amqp_empty_table);
    // 2. Dead Letter 큐 선언 및 바인딩
    amqp_queue_declare(conn, 1, amqp_cstring_bytes("dead_letter_queue"), 0, 1, 0, 0, amqp_empty_table);
    amqp_queue_bind(conn, 1, amqp_cstring_bytes("dead_letter_queue"),
                    amqp_cstring_bytes("dlx_exchange"),
                    amqp_cstring_bytes("dead"), amqp_empty_table);
    // 3. 메인 큐: x-dead-letter-exchange, x-dead-letter-routing-key 설정
    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("dead");
    amqp_table_t args;
    args.num_entries = 2;
    args.entries = entries;
    amqp_queue_declare(conn, 1, amqp_cstring_bytes("order_queue"),
                       0, 1, 0, 0, args);
    std::cout << "DLX 설정 완료: order_queue -> dlx_exchange -> dead_letter_queue" << std::endl;
    amqp_channel_close(conn, 1, AMQP_REPLY_SUCCESS);
    amqp_connection_close(conn, AMQP_REPLY_SUCCESS);
    amqp_destroy_connection(conn);
    return 0;
}

실패 시 DLQ로 보내기 (requeue=false)

// dlx_consumer.cpp - 처리 실패 시 DLQ로
#include <SimpleAmqpClient/SimpleAmqpClient.h>
#include <csignal>
#include <iostream>
#include <stdexcept>
static volatile sig_atomic_t g_running = 1;
void sig_handler(int) { g_running = 0; }
bool process_order(const std::string& body) {
    // JSON 파싱, 비즈니스 검증 등
    if (body.find("invalid") != std::string::npos) {
        throw std::runtime_error("Invalid order format");
    }
    return true;
}
int main() {
    signal(SIGINT, sig_handler);
    auto channel = AmqpClient::Channel::Create("localhost", 5672, "guest", "guest");
    std::string consumer_tag = channel->BasicConsume("order_queue", "", false, false, false, 1);
    std::cout << " [*] order_queue 대기 중..." << std::endl;
    while (g_running) {
        AmqpClient::Envelope::ptr_t envelope;
        if (channel->BasicConsumeMessage(consumer_tag, envelope, 1000)) {
            std::string body = envelope->Message()->Body();
            try {
                if (process_order(body)) {
                    channel->BasicAck(envelope);
                    std::cout << " [x] 처리 완료: " << body << std::endl;
                }
            } catch (const std::exception& e) {
                std::cerr << " [!!] 처리 실패 (DLQ로 전달): " << e.what() << std::endl;
                // requeue=false → Dead Letter Exchange로 전달
                channel->BasicReject(envelope, false);
            }
        }
    }
    channel->BasicCancel(consumer_tag);
    return 0;
}

재시도 횟수 제한 후 DLQ (헤더 활용)

무한 재시도를 막으려면 재시도 횟수를 어딘가에 기록해야 합니다. 흔한 실수는 헤더의 카운터를 읽기만 하고 BasicReject(requeue=true)로 되돌리는 것입니다. requeue는 메시지를 원본 그대로 큐에 돌려놓으므로 헤더 값이 영원히 0이고, 결국 무한 재시도가 됩니다. 카운터를 늘리려면 헤더를 바꾼 복사본을 다시 발행하고 원본은 ack해야 합니다.

// retry_with_header.cpp
void handle_with_retry_limit(AmqpClient::Channel::ptr_t channel,
                             AmqpClient::Envelope::ptr_t envelope) {
    auto msg = envelope->Message();
    int retry_count = 0;
    if (msg->HeaderTableIsSet()) {
        const AmqpClient::Table& headers = msg->HeaderTable();
        auto it = headers.find("x-retry-count");
        if (it != headers.end()) retry_count = static_cast<int>(it->second.GetInteger());
    }
    const int max_retries = 3;
    try {
        process_order(msg->Body());
        channel->BasicAck(envelope);
    } catch (const std::exception& e) {
        if (retry_count >= max_retries) {
            std::cerr << "최대 재시도 초과, DLQ로 전달" << std::endl;
            channel->BasicReject(envelope, false);  // DLX로
        } else {
            AmqpClient::Table headers = msg->HeaderTableIsSet() ? msg->HeaderTable()
                                                                : AmqpClient::Table{};
            headers["x-retry-count"] = AmqpClient::TableValue(retry_count + 1);
            msg->HeaderTable(headers);
            // 재시도 전용 큐(TTL로 지연)로 보내거나, 단순하게는 원래 큐로 다시 발행
            channel->BasicPublish("", "order_queue", msg);  // 기본 Exchange: 라우팅 키 = 큐 이름
            channel->BasicAck(envelope);  // 재발행 후 원본 ack
        }
    }
}

재발행과 ack 사이에 프로세스가 죽으면 메시지가 두 번 존재할 수 있으므로, 처리 로직은 중복에 견디도록 만들어야 합니다. RabbitMQ가 자동으로 붙이는 x-death 헤더(어느 큐에서 몇 번 dead-letter됐는지 기록)를 재시도 큐 패턴과 함께 쓰거나, 쿼럼 큐의 x-delivery-limit을 쓰면 직접 카운터를 관리하지 않아도 됩니다.


메시지 TTL과 큐 TTL

메시지별 TTL (per-message expiration)

발행 시 각 메시지에 expiration 속성을 설정합니다. 밀리초 문자열.

// message_ttl_producer.cpp - SimpleAmqpClient
#include <SimpleAmqpClient/SimpleAmqpClient.h>
#include <iostream>
int main() {
    auto channel = AmqpClient::Channel::Create("localhost", 5672, "guest", "guest");
    channel->DeclareQueue("ttl_queue", false, true, false, false);
    auto msg = AmqpClient::BasicMessage::Create("1시간 후 만료되는 메시지");
    msg->DeliveryMode(2);
    msg->Expiration("3600000");  // 3600000 ms = 1시간
    channel->BasicPublish("", "ttl_queue", msg);
    std::cout << " [x] TTL 1시간 메시지 발행" << std::endl;
    return 0;
}

큐 TTL (x-message-ttl)

큐에 들어오는 모든 메시지에 기본 TTL을 적용합니다.

// queue_ttl_setup.cpp - rabbitmq-c
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;  // 60초
amqp_table_t args;
args.num_entries = 1;
args.entries = entries;
amqp_queue_declare(conn, 1, amqp_cstring_bytes("expiring_queue"),
                   0, 1, 0, 0, args);

메시지 TTL과 큐의 x-message-ttl이 둘 다 있으면 더 작은 값이 적용됩니다. 큐 자체를 일정 시간 미사용 시 삭제하는 것은 x-expires로, 메시지 TTL과는 다른 설정입니다.

TTL + DLX: 만료 메시지를 DLQ로

만료된 메시지를 버리지 않고 DLQ로 보내 분석할 수 있습니다.

// ttl_dlx_setup.cpp - 큐 인자
// x-message-ttl: 60000 (60초)
// x-dead-letter-exchange: dlx_exchange
// x-dead-letter-routing-key: expired
amqp_table_entry_t entries[3];
entries[0].key = amqp_cstring_bytes("x-message-ttl");
entries[0].value.kind = AMQP_FIELD_KIND_I32;
entries[0].value.value.i32 = 60000;
entries[1].key = amqp_cstring_bytes("x-dead-letter-exchange");
entries[1].value.kind = AMQP_FIELD_KIND_UTF8;
entries[1].value.value.bytes = amqp_cstring_bytes("dlx_exchange");
entries[2].key = amqp_cstring_bytes("x-dead-letter-routing-key");
entries[2].value.kind = AMQP_FIELD_KIND_UTF8;
entries[2].value.value.bytes = amqp_cstring_bytes("expired");

exchange NOT_FOUND, x-max-priority 불일치, DLX 미동작: 에러 해결

”NOT_FOUND - no exchange ‘logs_topic’”

선언하지 않은 Exchange로 발행하면 브로커가 NOT_FOUND로 채널을 닫습니다. 기본 Exchange("")와 amq.* 내장 Exchange 외에는 미리 선언해야 합니다. Exchange 이름의 오타(대소문자 포함)도 같은 에러로 나타납니다.

// ✅ Exchange 선언 후 사용
amqp_exchange_declare(conn, 1, amqp_cstring_bytes("logs_topic"),
                     amqp_cstring_bytes("topic"), 0, 1, 0, 0, amqp_empty_table);

”PRECONDITION_FAILED - inequivalent arg ‘x-max-priority’”

같은 이름의 큐가 이미 다른 인자(여기서는 x-max-priority 없이)로 선언되어 있으면 재선언이 PRECONDITION_FAILED로 실패합니다. RabbitMQ는 선언 인자를 바꾸는 재선언을 허용하지 않으므로, 큐를 지우고 새 인자로 선언하거나 새 이름의 큐로 옮겨야 합니다. 큐를 지우면 쌓여 있던 메시지도 사라지므로 운영 중이라면 새 큐를 만들어 소비자를 옮기는 편이 안전합니다.

rabbitmqctl delete_queue order_queue

우선순위가 동작하지 않음

Priority(10)으로 발행했는데 순서가 그대로라면 먼저 큐에 x-max-priority가 설정되어 있는지 확인합니다. 설정이 없는 큐는 우선순위 속성을 무시합니다. 설정이 있어도 컨슈머의 prefetch가 크면 효과가 거의 없습니다. 우선순위는 큐 안에서 기다리는 메시지에만 적용되는데, prefetch가 크면 메시지가 큐에 쌓이기도 전에 컨슈머 버퍼로 먼저 넘어가 버리기 때문입니다. 우선순위가 중요한 큐는 prefetch를 1처럼 작게 둡니다. rabbitmq-c로 발행한다면 props._flags에 AMQP_BASIC_PRIORITY_FLAG를 넣었는지도 확인합니다.

DLX로 메시지가 안 감

requeue=false로 거부했는데 dead_letter_queue가 비어 있다면 세 곳을 봅니다. 원래 큐에 x-dead-letter-exchange가 설정되어 있는지, DLX Exchange가 실제로 존재하는지, 그리고 DLX에서 dead_letter_queue로 가는 바인딩이 dead-letter 라우팅 키(x-dead-letter-routing-key, 없으면 원래 라우팅 키)와 맞는지입니다. DLX Exchange가 없거나 바인딩이 맞지 않으면 메시지는 에러 없이 조용히 버려집니다.

”ACCESS_REFUSED - access to exchange ‘logs_topic’ in vhost ’/’ refused”

이름 문제가 아니라 권한 문제입니다. 접속한 사용자에게 해당 vhost의 write(발행) 또는 configure(선언) 권한이 없을 때 납니다. rabbitmqctl list_permissions -p /로 사용자 권한을 확인하고, rabbitmqctl set_permissions -p / app ".*" ".*" ".*"처럼 필요한 범위만큼 부여합니다.

토픽 바인딩 후 메시지가 큐에 안 옴

대개 라우팅 키 패턴이 생각과 다르게 매칭되는 경우입니다. 매칭 규칙은 다음과 같습니다.

# 패턴 규칙
logs.error  → logs.error 바인딩에 매칭
logs.error  → logs.# 바인딩에 매칭
logs.a.b    → logs.# 에만 매칭 (logs.* 는 한 단계만)

TTL이 적용되지 않음

TTL은 메시지가 큐에서 기다리는 동안에만 적용되고, 컨슈머에게 전달된 뒤에는 만료되지 않습니다. 또 하나 자주 놓치는 점은 메시지별 expiration이 만료된 메시지를 즉시 지우지 않는다는 것입니다. 클래식 큐는 만료된 메시지를 큐의 맨 앞에 도달했을 때 버리거나 dead-letter하므로, 앞에 TTL이 긴 메시지가 있으면 그 뒤의 짧은 TTL 메시지는 만료 시각이 지나도 남아 있습니다. 큐 전체에 같은 TTL을 주는 x-message-ttl은 메시지가 들어온 순서대로 만료되므로 이런 문제가 없습니다.


바인딩 수·우선순위 범위·Exchange 타입별 처리량

토픽 바인딩

토픽 Exchange는 발행할 때마다 라우팅 키를 바인딩 패턴들과 비교하므로, 바인딩이 수천 개 수준으로 많아지면 라우팅 비용이 늘어납니다. 소비자별로 꼭 필요한 패턴만 바인딩하고, 정확히 일치하는 키만 쓴다면 direct Exchange가 더 단순합니다.

우선순위 범위

앞에서 설명했듯 클래식 큐의 우선순위는 값마다 서브 큐를 두는 방식이라, x-max-priority를 크게 잡을수록 자원을 더 씁니다. 실제로 구분해야 하는 단계 수만큼만 둡니다.

DLQ 모니터링

DLQ에 메시지가 쌓이면 비정상입니다. 알림·대시보드로 모니터링하며, 주기적으로 원인 분석 후 재처리 또는 폐기합니다.

// DLQ 메시지 수 확인 (rabbitmqctl 또는 관리 API)
// rabbitmqctl list_queues name messages

TTL 큐와 영구 큐 분리

TTL이 짧은 메시지(알림)와 영구 메시지(주문)를 같은 큐에 넣지 마세요. 큐별로 TTL 정책이 다르므로 분리하는 것이 관리에 유리합니다.


선언 스크립트 분리·지연 재처리 큐·DLQ 컨슈머

Exchange·큐 선언 스크립트 분리

C++ 앱에서 매번 Exchange·큐를 선언하지 말고, 배포 스크립트나 설정 툴에서 한 번만 선언합니다.

#!/bin/bash
# setup_rabbitmq.sh: rabbitmqadmin 또는 HTTP API로 Exchange/Queue 선언
rabbitmqadmin declare exchange name=logs_topic type=topic durable=true
rabbitmqadmin declare queue name=error_logs durable=true
rabbitmqadmin declare binding source=logs_topic destination=error_logs routing_key=logs.error

재시도 큐 + TTL (지연 재처리)

실패 시 즉시 requeue하지 않으며, 재시도 전용 큐에 TTL을 두어 지연 후 재처리합니다.

flowchart LR
    Q1[main_queue] -->|실패| EX1[retry_exchange]
    EX1 --> Q2[retry_queue]
    Q2 -->|TTL 만료| EX2[main_exchange]
    EX2 --> Q1
// retry_queue: x-message-ttl=60000 (60초)
// x-dead-letter-exchange=main_exchange
// x-dead-letter-routing-key=retry
// → 60초 후 main_queue로 다시 들어감

설정 외부화

struct RabbitMQAdvancedConfig {
    std::string topic_exchange = "logs_topic";
    std::string priority_queue = "priority_orders";
    std::string dlx_exchange = "dlx_exchange";
    std::string dead_letter_queue = "dead_letter_queue";
    int max_priority = 10;
    int message_ttl_ms = 3600000;
};
RabbitMQAdvancedConfig loadFromEnv() {
    RabbitMQAdvancedConfig c;
    if (const char* x = std::getenv("RABBITMQ_TOPIC_EXCHANGE")) c.topic_exchange = x;
    if (const char* ttl = std::getenv("RABBITMQ_MSG_TTL_MS")) c.message_ttl_ms = std::stoi(ttl);
    return c;
}

DLQ 컨슈머 (분석·알림)

// dlq_consumer.cpp - 실패 메시지 수집 및 알림
void consume_dead_letters(AmqpClient::Channel::ptr_t channel) {
    std::string tag = channel->BasicConsume("dead_letter_queue", "", false, false, false);
    while (true) {
        AmqpClient::Envelope::ptr_t envelope;
        if (channel->BasicConsumeMessage(tag, envelope, 1000)) {
            std::string body = envelope->Message()->Body();
            // 로깅, 메트릭, PagerDuty 알림 등
            log_failure(body);
            notify_ops(body);
            channel->BasicAck(envelope);
        }
    }
}

Docker Compose (고급 설정)

# docker-compose.yml
services:
  rabbitmq:
    image: rabbitmq:3.13-management
    ports:
      - "5672:5672"
      - "15672:15672"
    environment:
      RABBITMQ_DEFAULT_USER: app
      RABBITMQ_DEFAULT_PASS: secret
    healthcheck:
      test: ["CMD", "rabbitmq-diagnostics", "check_running"]
      interval: 10s
      timeout: 5s
      retries: 3
  app:
    build: .
    environment:
      RABBITMQ_HOST: rabbitmq
      RABBITMQ_PORT: 5672
      RABBITMQ_USER: app
      RABBITMQ_PASSWORD: secret
    depends_on:
      rabbitmq:
        condition: service_healthy

다음 글: C++ Kafka 고급(#52-6) — 스트림 처리·트랜잭션 이전 글: C++에서 RabbitMQ 쓰기 — 프로듀서·컨슈머·작업 큐


참고 자료


같이 보면 좋은 글