RabbitMQ 아키텍처에서 메시지 교환기(Exchange)는 생산자가 전송한 메시지를 수신하고, 정의된 라우팅 규칙에 따라 특정 메시지 큐(Queue)로 전달하는 핵심 구성 요소입니다. RabbitMQ는 다양한 메시지 전달 패턴을 지원하기 위해 fanout, direct, topic, headers 네 가지 주요 Exchange 유형을 제공합니다. 이 글에서는 실제 활용 빈도가 높은 fanout, direct, topic 세 가지 Exchange 유형의 특징, 동작 방식, 그리고 상대적인 성능을 비교 분석합니다.
일반적으로 메시지 처리 속도 측면에서는 fanout이 가장 빠르며, 그 다음으로 direct, 마지막으로 topic 순으로 나타납니다. 대략적인 성능 비율은 fanout : direct : topic이 11 : 10 : 6 정도입니다.
Direct Exchange
Direct Exchange는 메시지의 라우팅 키(Routing Key)를 기반으로 정확히 일치하는 큐로 메시지를 전송합니다. Exchange에 바인딩된 큐는 특정 라우팅 키를 지정해야 하며, 메시지가 이 키와 완벽하게 일치할 때만 해당 큐로 전달됩니다. 예를 들어, 라우팅 키 'payment.success'로 바인딩된 큐는 'payment.success' 라우팅 키를 가진 메시지만 수신하며, 'payment.fail' 또는 'payment.success.user' 같은 메시지는 수신하지 않습니다.
- 메시지 라우팅은 라우팅 키의 정확한 일치 여부에 따라 결정됩니다.
- 큐는 특정 라우팅 키를 지정하여 Exchange에 바인딩되어야 합니다.
- RabbitMQ는 기본적으로 이름이 없는 "default Exchange"를 제공하며, 이는 Direct Exchange로 동작합니다. 이 경우 큐 이름이 라우팅 키로 사용됩니다.
- 지정된 라우팅 키와 일치하는 큐가 없으면 메시지는 버려집니다.
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;
public class DirectPublisher {
private static final String DIRECT_EXCHANGE = "direct_exchange_example";
private static final String INFO_QUEUE = "info_message_queue";
private static final String ERROR_QUEUE = "error_message_queue";
private static final String INFO_KEY = "info";
private static final String ERROR_KEY = "error";
public static void main(String[] argv) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost"); // RabbitMQ 서버 호스트 설정
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
// Direct Exchange 선언
channel.exchangeDeclare(DIRECT_EXCHANGE, "direct");
// 큐 선언 (durable, non-exclusive, non-auto-delete)
channel.queueDeclare(INFO_QUEUE, true, false, false, null);
channel.queueDeclare(ERROR_QUEUE, true, false, false, null);
// 큐와 Exchange를 라우팅 키로 바인딩
channel.queueBind(INFO_QUEUE, DIRECT_EXCHANGE, INFO_KEY);
channel.queueBind(ERROR_QUEUE, DIRECT_EXCHANGE, ERROR_KEY);
// 메시지 발행 (info 메시지)
String infoMessage = "시스템 정보: 정상 작동 중입니다.";
channel.basicPublish(DIRECT_EXCHANGE, INFO_KEY, MessageProperties.PERSISTENT_TEXT_PLAIN, infoMessage.getBytes("UTF-8"));
System.out.println(" [x] Sent '" + INFO_KEY + "':'" + infoMessage + "'");
// 메시지 발행 (error 메시지)
String errorMessage = "시스템 오류: 데이터베이스 연결 실패!";
channel.basicPublish(DIRECT_EXCHANGE, ERROR_KEY, MessageProperties.PERSISTENT_TEXT_PLAIN, errorMessage.getBytes("UTF-8"));
System.out.println(" [x] Sent '" + ERROR_KEY + "':'" + errorMessage + "'");
}
}
}
Fanout Exchange
Fanout Exchange는 가장 간단한 라우팅 방식으로, 수신된 모든 메시지를 자신에게 바인딩된 모든 큐로 브로드캐스트합니다. 메시지에 포함된 라우팅 키는 Fanout Exchange에서 무시됩니다. 이는 마치 서브넷 브로드캐스트와 유사하게 동작하여, 연결된 모든 구독자에게 메시지의 복사본을 전달합니다.
- 라우팅 키를 사용하지 않으며, 모든 바인딩된 큐로 메시지를 전달합니다.
- 메시지 복사본이 여러 큐로 동시에 전달되어 브로드캐스트 시나리오에 적합합니다.
- Exchange에 바인딩된 큐가 없다면 메시지는 소실됩니다.
- 가장 빠른 메시지 전달 성능을 제공합니다.
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;
public class FanoutPublisher {
private static final String FANOUT_EXCHANGE = "log_fanout_exchange";
public static void main(String[] argv) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost"); // RabbitMQ 서버 호스트 설정
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
// Fanout Exchange 선언
channel.exchangeDeclare(FANOUT_EXCHANGE, "fanout");
// 여러 임시 큐를 생성하고 동일한 Fanout Exchange에 바인딩
// 임시 큐는 소비자가 연결을 끊으면 자동으로 삭제됩니다.
String auditLogQueue = channel.queueDeclare().getQueue();
String errorLogQueue = channel.queueDeclare().getQueue();
String debugLogQueue = channel.queueDeclare().getQueue();
// Fanout은 라우팅 키를 무시하므로 바인딩 키는 빈 문자열("")을 사용합니다.
channel.queueBind(auditLogQueue, FANOUT_EXCHANGE, "");
channel.queueBind(errorLogQueue, FANOUT_EXCHANGE, "");
channel.queueBind(debugLogQueue, FANOUT_EXCHANGE, "");
String logMessage = "시스템 로그: 일반적인 이벤트 발생! (모든 큐로 전송)";
channel.basicPublish(FANOUT_EXCHANGE, "", MessageProperties.PERSISTENT_TEXT_PLAIN, logMessage.getBytes("UTF-8"));
System.out.println(" [x] Sent (Fanout) '" + logMessage + "' to all bound queues.");
}
}
}
Topic Exchange
Topic Exchange는 라우팅 키와 패턴 매칭을 통해 메시지를 라우팅하는 유연한 방식입니다. 큐는 하나 이상의 패턴(바인딩 키)을 사용하여 Exchange에 바인딩됩니다. 메시지의 라우팅 키가 이 패턴과 일치하면 메시지가 해당 큐로 전달됩니다. Topic Exchange는 복잡한 이벤트 필터링 및 멀티캐스팅 시나리오에 유용합니다.
- 라우팅 키와 바인딩 키 간의 패턴 매칭을 사용합니다.
- 바인딩 키에 두 가지 와일드카드 문자를 사용할 수 있습니다:
*(별표): 정확히 하나의 단어를 대체합니다. 예를 들어,log.*는log.info나log.error와 일치하지만log.info.app과는 일치하지 않습니다.#(샵): 0개 이상의 단어를 대체합니다. 예를 들어,log.#는log.info,log.error,log.info.app모두와 일치합니다.
- 바인딩 키는 점(.)으로 구분된 단어들로 구성됩니다.
- 라우팅 키와 일치하는 패턴을 가진 큐가 없으면 메시지는 소실됩니다.
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;
public class TopicPublisher {
private static final String TOPIC_EXCHANGE = "event_topic_exchange";
public static void main(String[] argv) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost"); // RabbitMQ 서버 호스트 설정
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
// Topic Exchange 선언
channel.exchangeDeclare(TOPIC_EXCHANGE, "topic");
// 큐 선언 및 패턴 바인딩
// Durable 큐로 선언
String criticalAlertsQueue = "critical_alerts_queue";
String generalLogsQueue = "general_logs_queue";
String debugInfoQueue = "debug_info_queue";
channel.queueDeclare(criticalAlertsQueue, true, false, false, null);
channel.queueDeclare(generalLogsQueue, true, false, false, null);
channel.queueDeclare(debugInfoQueue, true, false, false, null);
// 바인딩 키 설정
channel.queueBind(criticalAlertsQueue, TOPIC_EXCHANGE, "system.alert.critical"); // 정확한 경고
channel.queueBind(generalLogsQueue, TOPIC_EXCHANGE, "system.log.*"); // system.log.info, system.log.warn 등
channel.queueBind(debugInfoQueue, TOPIC_EXCHANGE, "system.log.#"); // system.log.debug.module1 등 (모든 하위 로그)
channel.queueBind(debugInfoQueue, TOPIC_EXCHANGE, "application.#"); // application의 모든 이벤트
// 메시지 발행
String message1 = "데이터베이스 오류: 연결 실패!";
channel.basicPublish(TOPIC_EXCHANGE, "system.alert.critical", MessageProperties.PERSISTENT_TEXT_PLAIN, message1.getBytes("UTF-8"));
System.out.println(" [x] Sent 'system.alert.critical':'" + message1 + "'");
String message2 = "새 사용자 로그인: admin";
channel.basicPublish(TOPIC_EXCHANGE, "system.log.info", MessageProperties.PERSISTENT_TEXT_PLAIN, message2.getBytes("UTF-8"));
System.out.println(" [x] Sent 'system.log.info':'" + message2 + "'");
String message3 = "모듈 초기화 성공.";
channel.basicPublish(TOPIC_EXCHANGE, "application.startup.success", MessageProperties.PERSISTENT_TEXT_PLAIN, message3.getBytes("UTF-8"));
System.out.println(" [x] Sent 'application.startup.success':'" + message3 + "'");
}
}
}