RabbitMQ는 메시지를 안전하고 유연하게 전달하는 데 최적화된 메시지 브로커입니다. RabbitMQ는 메시지를 큐(queue)에 저장하고, 필요할 때 적절한 수신자에게 전달합니다.
*, #)를 이용해 패턴이 일치하는 큐로 전달합니다.Kafka는 대규모 실시간 데이터 피드의 저장과 분석을 목적으로 하는 데이터 스트리밍 플랫폼입니다.
| 구분 | RabbitMQ | Kafka |
|---|---|---|
| 설계 철학 | 메시지의 안정적인 전달과 복잡한 라우팅에 집중 | 대규모 실시간 데이터의 저장과 스트림 처리에 집중 |
| 메시지 모델 | 큐(Queue) 중심 모델 | 토픽(Topic) 및 로그(Log) 중심 모델 |
| 데이터 보존 | 소비 완료 시 메시지 삭제 | 설정된 기간 동안 디스크에 로그 보존 |
| 적합한 사례 | 비동기 작업 큐, 요청/응답 패턴 | 로그 수집, 실시간 분석, 이벤트 소싱 |
RabbitMQ 실습의 핵심은 "주문(Order)이 발생했을 때, 결제(Payment)와 상품(Product) 서비스가 어떻게 각자의 메시지를 정확히 받아가는가"를 확인하는 것입니다.
설정(Configuration): TopicExchange를 사용하여 메시지가 큐로 흐르는 길을 정의합니다.
@Configuration
public class OrderApplicationQueueConfig {
@Value("${message.exchange}")
private String exchange;
@Value("${message.queue.product}")
private String queueProduct;
@Value("${message.queue.payment}")
private String queuePayment;
// Exchange 생성: 메시지를 분류하는 우체국 역할
@Bean
public TopicExchange exchange() { return new TopicExchange(exchange); }
// Queue 생성: 메시지가 저장되는 실제 공간
@Bean
public Queue queueProduct() { return new Queue(queueProduct); }
@Bean
public Queue queuePayment() { return new Queue(queuePayment); }
// Binding: 특정 키를 가진 메시지를 특정 큐로 연결
@Bean
public Binding bindingProduct() { return BindingBuilder.bind(queueProduct()).to(exchange()).with(queueProduct); }
@Bean
public Binding bindingPayment() { return BindingBuilder.bind(queuePayment()).to(exchange()).with(queuePayment); }
}
송신 및 수신(Producer & Consumer)
Producer: RabbitTemplate을 이용해 특정 큐(productQueue, paymentQueue)로 주문 ID를 발송합니다.
public void createOrder(String orderId) {
// productQueue로 orderId 발송
rabbitTemplate.convertAndSend(productQueue, orderId);
// paymentQueue로 orderId 발송
rabbitTemplate.convertAndSend(paymentQueue, orderId) ;
}
Consumer: @RabbitListener 어노테이션을 통해 지정된 큐를 실시간으로 감시하다가 메시지가 들어오면 즉시 처리합니다.
product Consumer
@RabbitListener(queues = "${message.queue.product}")
public void receiveMessage(String orderId) {
log.info("receive OrderId: {}, appName: {}", orderId, appName);
}
payment Consumer
@RabbitListener(queues = "${message.queue.payment}")
public void receiveMessage(String orderId) {
log.info("receive OrderId: {}, appName: {}", orderId, appName);
}
order/1 요청 시 하나의 요청이 product와 payment 두 개의 큐에 각각 쌓이는 것 확인productApplication을 두 개 실행하면 RabbitMQ가 메시지를 두 서버에 번갈아 전달합니다. 이를 통해 특정 서버에 부하가 쏠리지 않는 것 확인Kafka 실습은 "대량의 메시지를 얼마나 빠르게 처리하며, 컨슈머 그룹이 어떻게 독립적으로 데이터를 읽어가는가"에 집중했습니다.
프로듀서(Producer): 반복문을 통해 10개의 메시지를 순식간에 발행
public void sendMessage(String topic, String key, String message) {
for (int i = 0; i < 10; i++) {
kafkaTemplate.send(topic, key, message + " " + i);
}
}
컨슈머 그룹(Consumer Group): 동일한 데이터를 여러 목적의 그룹이 동시에 처리하는 구조
@Slf4j
@Component
public class ConsumerEndpoint {
@KafkaListener(groupId = "group_a", topics = "topic1")
public void consumerFromGroupA(String message) {
log.info("Group A consumed message form topic1 : " + message);
}
@KafkaListener(groupId = "group_b", topics = "topic1")
public void consumerFromGroupB(String message) {
log.info("Group B consumed message form topic1 : " + message);
}
@KafkaListener(groupId = "group_c", topics = "topic2")
public void consumerFromGroupC(String message) {
log.info("Group C consumed message form topic2 : " + message);
}
@KafkaListener(groupId = "group_c", topics = "topic3")
public void consumerFromGroupD(String message) {
log.info("Group C consumed message form topic3 : " + message);
}
@KafkaListener(groupId = "group_d", topics = "topic4")
public void consumerFromGroup0(String message) {
log.info("Group D consumed message form topic4 : " + message);
}
}
topic1 메시지를 group_a와 group_b가 각각 수신합니다. 이를 활용하면 하나의 이벤트로 여러 서비스(통계, 알림, 로그 등)를 동시에 구동할 수 있습니다.비즈니스 로직의 복잡한 연결은 RabbitMQ로, 대용량 로그 수집 및 실시간 분석은 Kafka로 구성하는 것이 가장 효율적이라고 볼 수 있습니다.