RabbitMQ로 메시지 큐 구축하기: Exchange 타입, AMQP 채널, DLX·Quorum Queue, Worker 패턴

이 글의 핵심

RabbitMQ의 Exchange 타입별 라우팅, 채널과 프리페치로 흐름을 제어하는 방법, 실패 메시지를 DLX로 보내는 패턴, Quorum Queue로 고가용성 확보하기를 이메일 전송 예제로 정리합니다.

이 글의 핵심

RabbitMQ로 메시지 큐 시스템을 구축하는 글입니다. Exchange, Queue, Binding, Routing, Worker Pattern, 고가용성까지 실전 예제로 정리했습니다.

실무에서 마주치는 문제들

API 호출이 실패하면 데이터가 유실돼요

동기 호출은 상대 서비스가 잠깐 죽어 있기만 해도 요청이 실패하고, 재시도 로직을 호출하는 쪽마다 따로 만들어야 합니다. 메시지 큐를 두면 요청은 브로커에 저장되고, 처리하는 쪽이 복구되면 이어서 가져가 처리합니다. 다만 “큐에 넣으면 절대 유실되지 않는다”는 자동 보장이 아니라, 뒤에서 설명할 durable 큐, persistent 메시지, publisher confirm, 수동 ack를 모두 올바르게 설정했을 때 얻는 성질입니다.

느린 작업이 전체를 막아요

이메일 전송처럼 외부 SMTP 서버를 거치는 작업은 수 초씩 걸릴 수 있습니다. 이를 HTTP 요청 안에서 기다리면 응답이 느려지고, 트래픽이 몰리면 요청 스레드나 커넥션이 고갈됩니다. API는 “이메일 전송 작업”을 큐에 넣고 바로 응답하고, 별도 워커가 큐에서 꺼내 처리하면 응답 시간이 작업 시간과 분리됩니다.

서비스 간 결합도가 높아요

주문 서비스가 결제·재고·알림 서비스를 직접 호출하면, 새 서비스가 추가될 때마다 주문 서비스 코드를 고쳐야 하고 하나가 느려지면 전체가 느려집니다. 주문 서비스는 “주문 생성됨” 메시지를 Exchange에 발행만 하고, 필요한 서비스가 각자 큐를 바인딩해 받게 하면 발행자는 누가 받는지 몰라도 됩니다.


메시지 브로커로서의 RabbitMQ

핵심 특징

RabbitMQ는 메시지 브로커입니다. 주요 개념:

  • Producer: 메시지 발행

  • Exchange: 메시지 라우팅

  • Queue: 메시지 저장

  • Consumer: 메시지 소비

  • Binding: Exchange와 Queue 연결 사용 사례:

  • 비동기 작업 처리

  • 마이크로서비스 통신

  • 이벤트 기반 아키텍처


Docker로 설치하기

Docker

docker run -d --name rabbitmq \
  -p 5672:5672 \
  -p 15672:15672 \
  rabbitmq:3-management

Management UI: http://localhost:15672 (guest/guest)


amqplib로 Producer와 Consumer 작성

설치

npm install amqplib

Producer (발행자)

// producer.ts
import amqp from 'amqplib';
async function sendMessage() {
  const connection = await amqp.connect('amqp://localhost');
  const channel = await connection.createChannel();
  const queue = 'tasks';
  const message = JSON.stringify({
    type: 'send_email',
    to: '[email protected]',
    subject: 'Hello',
  });
  await channel.assertQueue(queue, { durable: true });
  channel.sendToQueue(queue, Buffer.from(message), {
    persistent: true,
  });
  console.log('Sent:', message);
  setTimeout(() => {
    connection.close();
  }, 500);
}
sendMessage();

durable: true와 persistent: true는 짝으로 써야 의미가 있습니다. durable은 큐 선언이 브로커 재시작 후에도 남게 하고, persistent는 메시지를 디스크에 쓰도록 표시합니다. 둘 중 하나만 설정하면 재시작 후 큐는 있는데 메시지는 비어 있거나, 메시지를 저장하려 해도 큐 자체가 사라집니다. 또 이미 durable: false로 만든 큐를 같은 이름으로 durable: true로 다시 선언하면 “PRECONDITION_FAILED - inequivalent arg ‘durable’ for queue ‘tasks’” 에러와 함께 채널이 닫히므로, 큐 설정을 바꿀 때는 기존 큐를 지우거나 이름을 바꿔야 합니다.

마지막의 setTimeout(..., 500)은 예제를 단순하게 하려는 편법입니다. sendToQueue는 메시지를 클라이언트 버퍼에 넣고 바로 반환하므로, 곧바로 연결을 닫으면 아직 전송되지 않은 메시지가 사라질 수 있습니다. 실무에서는 createConfirmChannel()로 만든 채널에서 await channel.waitForConfirms()를 호출해 브로커가 메시지를 받았다고 확인한 뒤 닫는 것이 정석입니다. 또한 예제처럼 메시지마다 연결을 새로 여는 것은 비용이 커서, API 서버에서는 연결 하나를 앱 시작 시 만들어 재사용해야 합니다.

Consumer (소비자)

// consumer.ts
import amqp from 'amqplib';
async function consumeMessages() {
  const connection = await amqp.connect('amqp://localhost');
  const channel = await connection.createChannel();
  const queue = 'tasks';
  await channel.assertQueue(queue, { durable: true });
  channel.prefetch(1);  // 한 번에 하나씩 처리
  console.log('Waiting for messages...');
  channel.consume(queue, async (msg) => {
    if (msg) {
      const content = JSON.parse(msg.content.toString());
      console.log('Received:', content);
      // 작업 처리
      await processTask(content);
      // ACK (처리 완료)
      channel.ack(msg);
    }
  });
}
async function processTask(task: any) {
  console.log('Processing:', task);
  await new Promise(resolve => setTimeout(resolve, 1000));
  console.log('Done');
}
consumeMessages();

channel.consume의 기본값은 수동 ack입니다. 처리가 끝난 뒤 channel.ack(msg)를 호출해야 브로커가 메시지를 지우고, ack 전에 워커가 죽으면 브로커는 연결이 끊긴 것을 감지해 그 메시지를 다른 워커에게 다시 보냅니다. 이 덕분에 작업이 유실되지 않지만, 반대로 같은 메시지가 두 번 처리될 수 있다는 뜻이기도 합니다(at-least-once). 그래서 워커는 같은 메시지를 두 번 받아도 결과가 같도록 멱등하게 만들어야 합니다.

제가 이 코드 패턴에서 가장 자주 보는 문제는 ack를 호출하지 못하는 경로입니다. 위 코드에서 JSON.parse가 깨진 메시지로 예외를 던지거나 processTask가 reject되면, ack도 nack도 호출되지 않습니다. 그 메시지는 연결이 유지되는 동안 계속 “미확인” 상태로 남고, prefetch(1)이라면 워커는 다음 메시지를 받지 못한 채 멈춘 것처럼 보입니다. Management UI에서 “Unacked” 숫자가 줄지 않고 “Ready”만 쌓인다면 이 경우입니다. 모든 처리를 try/catch로 감싸 성공이면 ack, 실패면 nack을 반드시 호출하는 구조로 만들어야 합니다. 반대로 ack를 두 번 호출하면 “PRECONDITION_FAILED - unknown delivery tag” 에러로 채널이 닫힙니다.


Direct, Fanout, Topic Exchange

Direct Exchange

// Producer
await channel.assertExchange('logs', 'direct', { durable: true });
channel.publish('logs', 'error', Buffer.from('Error log'));
channel.publish('logs', 'info', Buffer.from('Info log'));
// Consumer (error만 받음)
await channel.assertExchange('logs', 'direct', { durable: true });
const queue = await channel.assertQueue('', { exclusive: true });
await channel.bindQueue(queue.queue, 'logs', 'error');
channel.consume(queue.queue, (msg) => {
  if (msg) {
    console.log('Error:', msg.content.toString());
    channel.ack(msg);
  }
});

Fanout Exchange (브로드캐스트)

// Producer
await channel.assertExchange('notifications', 'fanout', { durable: true });
channel.publish('notifications', '', Buffer.from('New notification'));
// Consumer 1
const queue1 = await channel.assertQueue('', { exclusive: true });
await channel.bindQueue(queue1.queue, 'notifications', '');
// Consumer 2
const queue2 = await channel.assertQueue('', { exclusive: true });
await channel.bindQueue(queue2.queue, 'notifications', '');

Topic Exchange (패턴 매칭)

// Producer
await channel.assertExchange('events', 'topic', { durable: true });
channel.publish('events', 'user.created', Buffer.from('User created'));
channel.publish('events', 'user.deleted', Buffer.from('User deleted'));
channel.publish('events', 'order.created', Buffer.from('Order created'));
// Consumer (user.* 패턴)
const queue = await channel.assertQueue('', { exclusive: true });
await channel.bindQueue(queue.queue, 'events', 'user.*');

세 타입의 차이는 라우팅 키를 해석하는 방식입니다. Direct는 라우팅 키가 바인딩 키와 정확히 같은 큐로만 보내고, Fanout은 라우팅 키를 무시하고 바인딩된 모든 큐에 복사합니다. Topic은 점(.)으로 구분된 단어 패턴을 쓰는데, *는 정확히 한 단어, #은 0개 이상의 단어와 일치합니다. 그래서 user.*는 user.created와는 맞지만 user.profile.updated와는 맞지 않고, 후자까지 받으려면 user.#을 써야 합니다. 이 차이를 놓쳐 일부 이벤트만 조용히 누락되는 경우가 흔합니다.

또 하나 알아 둘 점은 어떤 큐에도 라우팅되지 않은 메시지는 기본적으로 버려진다는 것입니다. 위 예제의 order.created는 바인딩된 큐가 없으므로 발행은 성공하지만 아무도 받지 못하고 사라집니다. 컨슈머보다 프로듀서를 먼저 띄웠을 때 초기 메시지가 없어지는 것도 같은 이유입니다. 중요한 메시지라면 발행 시 mandatory: true를 주고 return 이벤트로 라우팅 실패를 감지하거나, Exchange에 alternate exchange를 설정해 갈 곳 없는 메시지를 따로 모아야 합니다. 예제의 exclusive: true 익명 큐는 연결이 끊기면 삭제되므로, 컨슈머가 내려가 있는 동안의 메시지는 보관되지 않는다는 점도 기억해야 합니다.


AMQP 0-9-1 내부: 채널·프레임·흐름 제어

RabbitMQ 클라이언트가 amqp.connect 이후 createChannel을 호출하면, TCP 하나 위에 가상 연결인 채널이 다수 열립니다. 메시지는 프레임(frame) 단위로 오가며, publish는 기본적으로 비동기입니다. 따라서 publish 직후 바로 connection.close()를 호출하면 버퍼에 남은 프레임이 유실될 수 있어, 짧은 지연이나 Publisher Confirm 패턴이 필요합니다.

흐름 제어(prefetch) 는 QoS 프레임으로 브로커에 “이 컨슈머에게는 미확인 메시지를 최대 n개까지”라고 알립니다. prefetch(1)은 공정한 작업 분배에 유리하지만, 메시지 처리가 매우 짧으면 스루풋이 떨어질 수 있습니다. 반대로 prefetch를 크게 잡으면 한 워커에 긴 작업이 몰려 다른 큐 소비가 지연될 수 있으므로, 작업 길이 분포를 보고 조정합니다.


DLX, TTL, Quorum Queue로 실패 메시지 다루기

죽은 편지 큐(Dead Letter Exchange, DLX) 는 처리 실패·거부·TTL 초과 메시지를 별도 큐로 모아 운영자가 재처리·분석할 수 있게 합니다. basic.nack(requeue=false)만으로는 “어디로 갔는지”가 불명확할 수 있어, DLX + DLQ는 관측 가능성을 크게 높입니다.

RabbitMQ 3.8 이후 권장되는 Quorum Queue는 복제된 Raft 기반 큐로, 과반수 노드에 기록된 뒤에 확인을 돌려주므로 노드 하나가 죽어도 메시지가 유지됩니다. 예전의 클래식 미러링 큐는 오랫동안 폐기 예정 상태였다가 RabbitMQ 4.0에서 제거되었으므로, 새로 구성한다면 복제가 필요한 큐는 쿼럼 큐로 선언하는 것이 기본입니다(arguments: { 'x-queue-type': 'quorum' }). 쿼럼 큐는 반드시 durable이어야 하고, 복제 때문에 클래식 큐보다 디스크와 메모리를 더 쓰며, 전달 횟수 제한(x-delivery-limit)을 지원해 무한 재전달을 막을 수 있습니다. 트랜잭션 메일·결제 알림처럼 절대 유실되면 안 되는 워크로드는 이런 정책을 명시적으로 선택하고, 메시지 크기·TTL·최대 길이를 제한해 브로커 디스크를 보호합니다.

// DLX/DLQ 최소 예: 큐 선언 시 deadLetterExchange 지정 (브로커·정책에 맞게 조정)
await channel.assertExchange('tasks.dlx', 'direct', { durable: true });
await channel.assertQueue('tasks.dlq', { durable: true });
await channel.bindQueue('tasks.dlq', 'tasks.dlx', 'failed');

await channel.assertQueue('tasks', {
  durable: true,
  arguments: {
    'x-dead-letter-exchange': 'tasks.dlx',
    'x-dead-letter-routing-key': 'failed',
    'x-message-ttl': 3600000, // 1h — 운영 정책에 맞게
  },
});

여러 Worker로 작업 나누기

여러 Worker

// worker.ts
import amqp from 'amqplib';
async function startWorker(workerId: number) {
  const connection = await amqp.connect('amqp://localhost');
  const channel = await connection.createChannel();
  const queue = 'tasks';
  await channel.assertQueue(queue, { durable: true });
  channel.prefetch(1);
  console.log(`Worker ${workerId} started`);
  channel.consume(queue, async (msg) => {
    if (msg) {
      const task = JSON.parse(msg.content.toString());
      console.log(`Worker ${workerId} processing:`, task);
      await processTask(task);
      channel.ack(msg);
    }
  });
}
// 3개의 Worker 시작
for (let i = 1; i <= 3; i++) {
  startWorker(i);
}

이메일 전송 시스템 만들기

Producer (API 서버)

// api/send-email.ts
import express from 'express';
import amqp from 'amqplib';
const app = express();
app.use(express.json());
let channel: amqp.Channel;
async function setupRabbitMQ() {
  const connection = await amqp.connect('amqp://localhost');
  channel = await connection.createChannel();
  await channel.assertQueue('emails', { durable: true });
}
app.post('/api/send-email', async (req, res) => {
  const { to, subject, body } = req.body;
  channel.sendToQueue(
    'emails',
    Buffer.from(JSON.stringify({ to, subject, body })),
    { persistent: true }
  );
  res.json({ message: 'Email queued' });
});
setupRabbitMQ().then(() => {
  app.listen(3000, () => console.log('API server running'));
});

Consumer (Email Worker)

// workers/email-worker.ts
import amqp from 'amqplib';
import nodemailer from 'nodemailer';
const transporter = nodemailer.createTransport({
  host: 'smtp.gmail.com',
  port: 587,
  auth: {
    user: process.env.EMAIL_USER,
    pass: process.env.EMAIL_PASS,
  },
});
async function startEmailWorker() {
  const connection = await amqp.connect('amqp://localhost');
  const channel = await connection.createChannel();
  await channel.assertQueue('emails', { durable: true });
  channel.prefetch(1);
  console.log('Email worker started');
  channel.consume('emails', async (msg) => {
    if (msg) {
      const { to, subject, body } = JSON.parse(msg.content.toString());
      try {
        await transporter.sendMail({ to, subject, text: body });
        console.log('Email sent to:', to);
        channel.ack(msg);
      } catch (error) {
        console.error('Failed to send email:', error);
        channel.nack(msg, false, true);  // 재시도
      }
    }
  });
}
startEmailWorker();

nack(msg, false, true)의 세 번째 인자 true는 메시지를 큐 앞쪽으로 즉시 되돌리라는 뜻입니다. 일시적인 SMTP 장애라면 이것으로 충분하지만, 수신 주소가 잘못된 메시지처럼 몇 번을 재시도해도 실패하는 메시지(poison message)라면 워커는 그 메시지를 받자마자 실패하고 되돌리는 일을 끝없이 반복합니다. 로그가 초당 수백 줄씩 같은 에러로 채워지고 CPU를 소모하는데 다른 메시지는 거의 처리되지 않는, 운영에서 꽤 자주 보는 장애 형태입니다. 해결책은 재시도 횟수를 제한하는 것입니다. 쿼럼 큐의 x-delivery-limit을 쓰거나, 메시지 헤더에 재시도 횟수를 기록해 일정 횟수를 넘으면 nack(msg, false, false)로 앞에서 설정한 DLX로 보냅니다. 즉시 재시도 대신 TTL이 있는 대기 큐를 거쳐 몇 초 뒤 원래 큐로 돌아오게 하면 백오프 재시도도 구현할 수 있습니다.

API 쪽 코드에도 빈틈이 있습니다. sendToQueue의 반환값이 false면 클라이언트 쓰기 버퍼가 가득 찼다는 뜻이라 drain 이벤트를 기다려야 하고, RabbitMQ 연결이 끊기면 channel이 닫힌 채로 남아 이후 모든 요청이 “Channel closed” 에러로 실패합니다. connection.on('close')에서 재연결하는 로직을 두거나, 재연결을 처리해 주는 amqp-connection-manager 같은 래퍼를 쓰는 것이 일반적입니다.


클러스터로 고가용성 구성

클러스터

# 노드 1
docker run -d --hostname rabbit1 --name rabbit1 \
  -p 5672:5672 -p 15672:15672 \
  -e RABBITMQ_ERLANG_COOKIE='secret' \
  rabbitmq:3-management
# 노드 2
docker run -d --hostname rabbit2 --name rabbit2 \
  -p 5673:5672 -p 15673:15672 \
  -e RABBITMQ_ERLANG_COOKIE='secret' \
  --link rabbit1:rabbit1 \
  rabbitmq:3-management
# 클러스터 조인
docker exec -it rabbit2 rabbitmqctl stop_app
docker exec -it rabbit2 rabbitmqctl join_cluster rabbit@rabbit1
docker exec -it rabbit2 rabbitmqctl start_app

이 예제는 클러스터 조인 절차를 보여 주기 위한 최소 구성입니다. 노드들이 서로를 인증하는 Erlang 쿠키는 모든 노드에서 같아야 하며, 다르면 조인 시 “Authentication failed (rejected by the remote node), please check the Erlang cookie” 에러가 납니다. 최근 이미지에서는 RABBITMQ_ERLANG_COOKIE 환경 변수가 폐기 예정이라 쿠키 파일(/var/lib/rabbitmq/.erlang.cookie)을 마운트하는 방식이 권장되고, --link도 오래된 Docker 기능이라 실제로는 사용자 정의 네트워크나 Docker Compose로 묶습니다. 노드 수는 쿼럼 큐의 과반수 계산 때문에 3이나 5처럼 홀수로 두는 것이 일반적이며, 2노드 클러스터는 한 노드가 죽으면 과반을 잃어 쿼럼 큐가 쓰기를 멈춥니다. 클러스터는 같은 데이터센터 안의 낮은 지연 네트워크를 전제로 하므로, 리전 간 복제가 필요하면 Federation이나 Shovel 플러그인을 씁니다.


Management Plugin으로 모니터링

Management Plugin

# 플러그인 활성화
rabbitmq-plugins enable rabbitmq_management
# API로 확인
curl -u guest:guest http://localhost:15672/api/overview

Management UI와 API에서 가장 먼저 봐야 할 지표는 큐별 Ready(대기 중)와 Unacked(전달했지만 확인 전) 메시지 수, 그리고 발행·소비 속도입니다. Ready가 꾸준히 늘면 소비 속도가 생산을 따라가지 못하는 것이고, Unacked만 늘면 워커가 ack를 하지 못하고 있는 것입니다. 브로커 메모리가 기본 한도(가용 메모리의 40%)를 넘으면 RabbitMQ는 메모리 알람을 걸고 모든 발행 연결을 막습니다. 이때 프로듀서는 에러 없이 멈춘 것처럼 보이므로, 발행이 갑자기 응답하지 않는다면 Management UI 상단의 알람 표시를 먼저 확인해 볼 만합니다. 장기적으로는 rabbitmq_prometheus 플러그인으로 지표를 Prometheus·Grafana에 연결하는 구성이 흔합니다.


증상별 원인과 대응

증상흔한 원인대응
메시지가 소비되지 않음noAck: true로 consume했거나, prefetch로 막힘ack/nack 경로 확인, channel.checkQueue로 backlog
중복 처리네트워크 끊김 후 재전달, at-least-once 특성멱등 키·아웃박스 패턴
브로커 메모리 급증메시지가 쌓이고 소비 속도 < 생산 속도소비자 스케일아웃, TTL·max-length 정책
“RESOURCE_LOCKED” 등미러링·쿼럼 정책 불일치클러스터 버전·정책 문서 점검

RabbitMQ 요약

  • RabbitMQ: 메시지 브로커
  • Exchange: 메시지 라우팅
  • Queue: 메시지 저장
  • Worker Pattern: 병렬 처리
  • 고가용성: 클러스터 지원
  • AMQP: 표준 프로토콜

프로덕션 체크리스트

  • RabbitMQ 설치
  • Exchange 및 Queue 설계
  • Producer 구현
  • Consumer 구현
  • 에러 처리 및 재시도
  • 클러스터 구성
  • 모니터링 설정

같이 보면 좋은 글


자주 묻는 질문 (FAQ)

Q. RabbitMQ vs Kafka, 어떤 게 나은가요?

A. RabbitMQ는 메시지를 큐에 넣고 소비자가 ack하면 지우는 메시지 브로커이고, Kafka는 메시지를 로그에 보관하며 소비자가 각자 읽은 위치(offset)를 관리하는 이벤트 스트리밍 플랫폼입니다. 작업 분배, 복잡한 라우팅, 메시지별 재시도가 중요하면 RabbitMQ가, 대용량 이벤트를 오래 보관하고 여러 소비자가 과거 데이터를 다시 읽어야 하면 Kafka가 잘 맞습니다.

Q. Redis Pub/Sub vs RabbitMQ, 차이가 뭔가요?

A. Redis Pub/Sub은 메시지를 저장하지 않아서, 발행 순간에 연결되어 있지 않은 구독자는 그 메시지를 받지 못합니다. RabbitMQ는 큐에 메시지를 보관하고 ack로 처리 완료를 확인하므로, 소비자가 잠시 내려가 있어도 메시지를 이어서 처리할 수 있습니다. Redis에서도 영속성과 소비자 그룹이 필요하면 Pub/Sub 대신 Redis Streams를 쓸 수 있습니다.

Q. 메시지 순서가 보장되나요?

A. 하나의 큐에서 하나의 소비자가 받을 때는 발행 순서대로 전달됩니다. 소비자가 여럿이면 각자 병렬로 처리하므로 처리 완료 순서는 보장되지 않고, nack로 되돌아간 메시지가 다시 전달되면 순서가 바뀔 수도 있습니다. 순서가 중요하면 단일 활성 소비자(x-single-active-consumer)를 쓰거나, 키별로 큐를 나눠 같은 키의 메시지가 한 소비자에게만 가게 합니다.

Q. 프로덕션에서 사용해도 되나요?

A. 네. RabbitMQ는 오래 운영되어 온 오픈소스 브로커이고 클러스터·쿼럼 큐·모니터링 도구가 갖춰져 있습니다. 다만 이 글의 예제는 개념 설명용이므로, 운영에 쓰려면 publisher confirm, 재연결, 재시도 제한과 DLX, 기본 guest 계정 제거(guest는 기본적으로 localhost에서만 접속 가능)와 TLS 설정을 추가해야 합니다.