RabbitMQ의 메시징 모델 심층 분석

AMQP(Advanced Message Queuing Protocol)를 기반으로 구축된 RabbitMQ는 분산 시스템에서 효율적인 메시지 교환을 가능하게 하는 강력한 메시지 브로커입니다. Erlang 언어로 개발되어 뛰어난 안정성과 확장성을 자랑합니다. 다양한 애플리케이션 컴포넌트 간의 비동기 통신을 지원하며, 이는 시스템의 유연성과 견고성을 향상시키는 데 기여합니다.

RabbitMQ는 메시지 흐름을 관리하기 위한 여러 가지 패턴, 즉 '메시징 모델'을 제공합니다. 비록 일부 문서에서 6가지 모델을 언급하기도 하지만, 일반적으로 사용되는 순수한 메시지 큐 모델은 5가지로 압축됩니다. 특히, '발행-구독(Publish/Subscribe)' 계열에 속하는 세 가지 모델(Fanout, Direct, Topic)은 메시지 라우팅 방식에 따라 세분화됩니다. 본 글에서는 이 다섯 가지 핵심 메시징 모델의 개념과 활용 방법을 자세히 살펴봅니다.

1. 기본 메시징 모델 (Basic Messaging Model)

가장 기본적인 메시징 패턴으로, 하나의 생산자가 메시지를 특정 큐에 발행하고, 하나의 소비자가 해당 큐에서 메시지를 수신합니다. 이 모델은 다음과 같은 핵심 요소들로 구성됩니다.

  • 생산자(Producer): 메시지를 생성하여 RabbitMQ 브로커로 전송하는 애플리케이션입니다.
  • 소비자(Consumer): 큐에 도착한 메시지를 수신하고 처리하는 애플리케이션입니다.
  • 큐(Queue): 생산자가 보낸 메시지가 일시적으로 저장되는 공간입니다. 소비자는 이 큐에서 메시지를 가져갑니다.

이 모델은 우편 시스템에 비유할 수 있습니다. 메시지 생산자가 우편함(큐)에 편지(메시지)를 넣으면, 우체부(RabbitMQ)가 이를 관리하고, 수신자(소비자)가 해당 편지를 받아보는 것과 유사합니다.

메시지 수신 확인 (Acknowledge)

소비자가 큐에서 메시지를 성공적으로 처리했는지 여부는 RabbitMQ에 '수신 확인(ACK)' 메시지를 전송함으로써 알립니다. 이 수신 확인 방식에는 두 가지 종류가 있습니다.

  • 자동 수신 확인 (Auto-ACK): 소비자가 메시지를 받자마자 RabbitMQ에게 자동으로 수신 확인을 보냅니다. 메시지가 중요하지 않아 손실되어도 무방한 경우 유용하지만, 소비자가 메시지 처리 도중 오류로 인해 중단되면 해당 메시지는 손실될 수 있습니다.
  • 수동 수신 확인 (Manual-ACK): 소비자가 메시지를 완전히 처리한 후에 명시적으로 수신 확인을 보냅니다. 이 방식은 메시지 손실을 방지하는 데 필수적입니다. 소비자가 메시지를 수신한 후 처리 도중 실패하더라도, 수동 수신 확인을 보내지 않으면 RabbitMQ는 메시지를 큐에 유지하여 다른 소비자가 다시 처리하거나 재시도할 수 있도록 합니다.

다음은 수동 수신 확인을 설정하는 코드 예시입니다. basicConsume 메서드의 두 번째 인자를 false로 설정하면 수동 수신 확인 모드가 활성화됩니다.

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.DeliverCallback;

public class SimpleMessageConsumer {
    private final static String MY_QUEUE_NAME = "my_task_queue";

    public static void main(String[] args) throws Exception {
        // ... 채널 생성 로직 생략 ...
        // 실제 RabbitMQ 연결 및 채널 생성 로직을 여기에 구현
        Channel channel = getRabbitMQChannel(); // 가정: 채널 획득 메서드

        DeliverCallback deliveryCallback = (consumerTag, delivery) -> {
            String message = new String(delivery.getBody(), "UTF-8");
            System.out.println(" [x] 메시지 수신: '" + message + "'");
            try {
                // 메시지 처리 로직
                Thread.sleep(1000); // 작업 처리 시뮬레이션 (1초 소요)
                System.out.println(" [x] 메시지 처리 완료: '" + message + "'");
                channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); // 수동 ACK
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                System.err.println(" [x] 메시지 처리 실패 및 인터럽트: " + e.getMessage());
                // 필요시 nack 처리 (메시지를 다시 큐로 보내거나 폐기)
                // true: 메시지를 큐에 다시 넣음, false: 메시지를 폐기
                channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, true);
            }
        };

        // 수동 수신 확인 모드로 소비자 등록 (autoAck = false)
        channel.basicConsume(MY_QUEUE_NAME, false, deliveryCallback, consumerTag -> {});
        System.out.println(" [*] 메시지 대기 중. 종료하려면 Ctrl+C를 누르세요.");
    }

    // 실제 RabbitMQ 연결 및 채널 생성 로직을 반환하는 가상의 메서드
    private static Channel getRabbitMQChannel() {
        // 실제 구현 시 ConnectionFactory, Connection 등을 통해 채널 생성
        return null; // 예시를 위해 null 반환
    }
}

2. 워크 큐 모델 (Work Queues Model)

단일 소비자가 메시지를 처리하는 기본 모델과 달리, 워크 큐(Work Queues) 모델은 여러 소비자가 하나의 큐에 연결되어 메시지를 공동으로 처리합니다. 이 방식은 메시지 처리 시간이 길거나, 생산 속도가 소비 속도를 능가하여 메시지가 큐에 쌓이는 문제를 해결하는 데 매우 효과적입니다. 여러 소비자가 작업을 분산하여 처리함으로써 전체 처리량을 늘리고 지연 시간을 줄일 수 있습니다.

이 모델의 주요 특징은 다음과 같습니다.

  • 부하 분산: 큐에 있는 메시지들은 여러 소비자에게 분산되어 전달됩니다.
  • 메시지 중복 처리 방지: 하나의 메시지는 오직 하나의 소비자에게만 전달되며, 해당 소비자가 메시지 처리를 완료하면 큐에서 사라지므로 메시지 중복 처리는 발생하지 않습니다.

공정한 메시지 분배 (Fair Dispatch)

기본적으로 RabbitMQ는 메시지를 라운드 로빈(round-robin) 방식으로 소비자들에게 분배합니다. 즉, 각 소비자에게 동일한 수의 메시지를 순서대로 할당합니다. 그러나 만약 일부 소비자의 메시지 처리 속도가 현저히 느리다면, 느린 소비자에게 할당된 메시지는 처리되지 못하고 쌓여 전체 시스템의 효율을 떨어뜨릴 수 있습니다.

이러한 문제를 해결하기 위해 RabbitMQ는 basicQos 메서드를 통해 '공정한 분배(Fair Dispatch)'를 설정할 수 있습니다. basicQos(1)을 호출하면 RabbitMQ는 각 소비자가 동시에 처리할 수 있는 메시지의 수를 1개로 제한합니다. 이는 소비자가 현재 처리 중인 메시지에 대한 수신 확인(ACK)을 보내기 전까지는 새로운 메시지를 보내지 않도록 하여, 더 빠른 소비자가 더 많은 메시지를 처리하게 함으로써 효율적인 부하 분산을 가능하게 합니다.

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.DeliverCallback;

public class WorkerConsumer {
    private final static String TASK_QUEUE_NAME = "processing_tasks";

    public static void main(String[] args) throws Exception {
        // ... 채널 생성 로직 생략 ...
        // 실제 RabbitMQ 연결 및 채널 생성 로직을 여기에 구현
        Channel workerChannel = getRabbitMQChannel(); // 가정: 채널 획득 메서드

        // 워커 채널이 동시에 한 개의 메시지만 처리하도록 설정 (Fair Dispatch)
        // 각 소비자는 한 번에 1개의 메시지만 받음
        workerChannel.basicQos(1); 

        DeliverCallback taskDeliveryCallback = (consumerTag, delivery) -> {
            String taskMessage = new String(delivery.getBody(), "UTF-8");
            System.out.println(" [worker] 작업 메시지 수신: '" + taskMessage + "'");
            try {
                // 실제 작업 처리 로직 시뮬레이션
                Thread.sleep(5000); // 작업이 오래 걸리는 상황 가정 (5초 소요)
                System.out.println(" [worker] 작업 처리 완료: '" + taskMessage + "'");
                workerChannel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); // 작업 완료 후 수동 ACK
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                System.err.println(" [worker] 작업 처리 중단: " + e.getMessage());
                // 메시지 재큐 또는 폐기
                workerChannel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, true); 
            }
        };

        // 수동 ACK와 함께 소비자 등록
        workerChannel.basicConsume(TASK_QUEUE_NAME, false, taskDeliveryCallback, consumerTag -> {});
        System.out.println(" [*] 워커 프로세스 시작. 새로운 작업을 기다립니다.");
    }

    // 실제 RabbitMQ 연결 및 채널 생성 로직을 반환하는 가상의 메서드
    private static Channel getRabbitMQChannel() {
        return null; // 예시를 위해 null 반환
    }
}

3. 발행-구독 모델: Fanout (Publish/Subscribe - Fanout)

발행-구독 모델은 하나의 메시지가 여러 소비자에게 동시에 전달되어야 할 때 사용됩니다. Fanout 타입의 교환기(Exchange)는 이 모델의 핵심적인 역할을 합니다. 이 모델의 특징은 다음과 같습니다.

  1. 다수의 소비자: 여러 소비자가 동시에 메시지를 받을 수 있습니다.
  2. 개별 큐: 각 소비자는 메시지를 수신하기 위한 자신만의 독립적인 큐를 가집니다.
  3. 교환기에 큐 바인딩: 모든 소비자의 큐는 동일한 교환기에 바인딩(연결)되어야 합니다.
  4. 생산자는 교환기로만 전송: 생산자는 메시지를 큐가 아닌 교환기로만 전송합니다. 생산자는 특정 큐를 지정할 수 없으며, 메시지가 어떤 큐로 전달될지 관여하지 않습니다.
  5. 교환기의 브로드캐스트: 교환기는 자신에게 바인딩된 모든 큐로 수신된 메시지를 브로드캐스트(전파)합니다.
  6. 모든 소비자가 메시지 수신: 결과적으로, 교환기에 바인딩된 모든 큐의 소비자들은 동일한 메시지를 받게 됩니다.

이 모델은 공지사항 시스템이나 실시간 알림 시스템처럼 모든 구독자에게 동일한 정보를 동시에 전달해야 할 때 매우 유용합니다. 생산자는 단순히 메시지를 교환기에 발행하고, RabbitMQ는 이 메시지를 모든 구독 큐로 복제하여 전달하는 역할을 수행합니다. 이 과정에서 큐와 교환기 간의 '바인딩' 설정이 필수적입니다.

4. 발행-구독 모델: Direct (Publish/Subscribe - Direct)

Fanout 모델이 모든 구독자에게 메시지를 브로드캐스트하는 반면, Direct 모델은 특정 조건에 따라 메시지를 선택적으로 라우팅해야 할 때 사용됩니다. 이 모델에서는 '라우팅 키(Routing Key)'라는 개념이 도입되어 메시지 필터링이 가능해집니다.

Direct 모델의 작동 방식은 다음과 같습니다.

  • 라우팅 키 지정 바인딩: 큐가 교환기에 바인딩될 때, 단순히 연결되는 것이 아니라 특정 '라우팅 키'를 지정하여 바인딩됩니다.
  • 메시지에 라우팅 키 포함: 생산자는 메시지를 교환기로 전송할 때, 해당 메시지에 특정 '라우팅 키'를 함께 포함해야 합니다.
  • 정확한 일치 기반 라우팅: 교환기는 수신된 메시지의 라우팅 키와, 자신에게 바인딩된 큐들의 라우팅 키를 비교합니다. 오직 메시지의 라우팅 키와 정확히 일치하는 라우팅 키로 바인딩된 큐로만 메시지를 전달합니다.

이 모델은 시스템 로그(예: info, warning, error)를 분류하거나, 특정 작업 유형에 따라 메시지를 다른 처리기로 보내야 할 때 적합합니다. 예를 들어, 'error' 라우팅 키를 가진 메시지는 오류 처리 큐로, 'info' 라우팅 키를 가진 메시지는 정보성 로그 큐로만 전달되도록 설정할 수 있습니다. 한 큐는 여러 라우팅 키로 교환기에 바인딩될 수 있습니다.

5. 발행-구독 모델: Topic (Publish/Subscribe - Topic)

Topic 모델은 Direct 모델과 유사하게 라우팅 키를 사용하여 메시지를 필터링하지만, 훨씬 더 유연한 패턴 매칭 기능을 제공합니다. 이 모델에서는 라우팅 키에 와일드카드(wildcard)를 사용하여 다수의 라우팅 키 패턴에 메시지를 일치시킬 수 있습니다. 라우팅 키는 일반적으로 마침표(.)로 구분된 단어들로 구성됩니다. 예를 들어, stock.us.ny.aapl과 같습니다.

Topic 모델에서 사용되는 와일드카드 규칙은 다음과 같습니다.

  • * (별표): 정확히 하나의 단어를 대체합니다.
  • # (샵): 하나 또는 그 이상의 단어들을 대체합니다 (0개도 포함).

와일드카드 사용 예시:

  • *.news: sport.news, weather.news는 일치하지만, korea.sport.news는 일치하지 않습니다. (하나의 단어만 일치)
  • sport.#: sport.football, sport.team.players, sport (단어 0개)는 모두 일치합니다. (하나 또는 그 이상의 단어 일치)
  • audit.#: audit.irs.corporate 또는 audit.irs와 일치합니다.
  • usa.*: usa.eastusa.west와 일치하지만, usa.east.newyork와는 일치하지 않습니다.

Topic 모델은 복잡한 이벤트 라우팅, 예를 들어 특정 지역의 특정 유형 주식 거래 알림(stock.us.ny.aapl)과 같이 계층적이고 다양한 조건에 따라 메시지를 분배해야 할 때 매우 강력합니다. 큐는 필요한 라우팅 패턴으로 교환기에 바인딩되어, 해당 패턴에 맞는 메시지만 수신하게 됩니다.

지속성 (Persistence)

앞서 설명한 수신 확인(ACK) 메커니즘은 소비자가 메시지를 처리하는 도중 발생할 수 있는 메시지 손실을 방지합니다. 그러나 메시지 브로커인 RabbitMQ 자체가 예기치 않게 종료되거나 재시작될 경우, 아직 소비자에게 전달되지 않은 메시지나 처리 대기 중인 큐의 메시지들은 손실될 수 있습니다. 이러한 상황을 방지하고 메시지의 안정적인 전달을 보장하기 위해 '지속성(Persistence)' 기능이 필수적입니다.

메시지를 지속시키기 위해서는 다음 세 가지 요소가 모두 지속 가능한(Durable) 상태로 선언되어야 합니다.

1. 교환기(Exchange) 지속성

교환기를 내구성이 있는(durable) 상태로 선언하면, RabbitMQ 서버가 재시작되어도 해당 교환기가 사라지지 않고 유지됩니다.

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;

public class DurablePublisher {
    private final static String DURABLE_EXCHANGE_NAME = "my_durable_topic_exchange";

    public static void main(String[] args) throws Exception {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost"); // RabbitMQ 호스트 설정 (필요에 따라 변경)
        
        try (Connection connection = factory.newConnection();
             Channel publisherChannel = connection.createChannel()) {

            // 교환기 선언: "topic" 타입, durable = true (지속성 활성화)
            publisherChannel.exchangeDeclare(DURABLE_EXCHANGE_NAME, "topic", true);
            System.out.println(" [*] 내구성 있는 교환기 '" + DURABLE_EXCHANGE_NAME + "'가 선언되었습니다.");

            // 이후 메시지 발행 로직...
        }
    }
}

2. 큐(Queue) 지속성

큐를 내구성이 있는 상태로 선언하면, RabbitMQ 서버가 재시작되어도 큐가 삭제되지 않고 그 안에 저장된 메시지들이 보존됩니다.

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;

public class DurableConsumer {
    private final static String DURABLE_QUEUE_NAME = "my_durable_queue";

    public static void main(String[] args) throws Exception {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost"); // RabbitMQ 호스트 설정 (필요에 따라 변경)

        try (Connection connection = factory.newConnection();
             Channel consumerChannel = connection.createChannel()) {

            // 큐 선언: durable = true (지속성 활성화)
            // queueDeclare(queue, durable, exclusive, autoDelete, arguments)
            consumerChannel.queueDeclare(DURABLE_QUEUE_NAME, true, false, false, null);
            System.out.println(" [*] 내구성 있는 큐 '" + DURABLE_QUEUE_NAME + "'가 선언되었습니다.");

            // 이후 메시지 소비 로직...
        }
    }
}

3. 메시지 지속성

교환기와 큐가 내구성이 있더라도, 메시지 자체가 지속성을 갖도록 설정하지 않으면 RabbitMQ 서버 재시작 시 메시지는 손실될 수 있습니다. 메시지를 발행할 때 MessageProperties.PERSISTENT_TEXT_PLAIN 또는 MessageProperties.PERSISTENT_BASIC과 같은 속성을 사용하여 메시지를 지속 가능한 상태로 만듭니다.

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;

public class PersistentMessageProducer {
    private final static String PERSISTENT_EXCHANGE_NAME = "my_persistent_exchange";
    private final static String ROUTING_KEY = "event.critical";

    public static void main(String[] args) throws Exception {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost"); // RabbitMQ 호스트 설정 (필요에 따라 변경)

        try (Connection connection = factory.newConnection();
             Channel producerChannel = connection.createChannel()) {

            // 교환기도 durable로 선언 (메시지 지속성을 위한 필수 조건)
            producerChannel.exchangeDeclare(PERSISTENT_EXCHANGE_NAME, "direct", true); 

            String criticalMessage = "시스템에 중요한 이벤트가 발생했습니다!";
            // 메시지 발행: MessageProperties.PERSISTENT_TEXT_PLAIN을 사용하여 메시지 지속성 설정
            producerChannel.basicPublish(
                PERSISTENT_EXCHANGE_NAME,
                ROUTING_KEY,
                MessageProperties.PERSISTENT_TEXT_PLAIN, // 메시지 지속성 속성
                criticalMessage.getBytes("UTF-8")
            );
            System.out.println(" [x] 지속 가능한 메시지 전송됨: '" + criticalMessage + "'");

        }
    }
}

태그: RabbitMQ MessageQueue AMQP MessagingModels PublishSubscribe

8월 2일 19:03에 게시됨