RabbitMQ란?
RabbitMQ는 서비스나 애플리케이션 사이에서 메시지를 전달하는 메시지 브로커(Message Broker) 이다.
서비스가 다른 서비스를 직접 호출하는 대신 RabbitMQ에 메시지를 전달하면, RabbitMQ가 메시지를 Queue에 보관하고 해당 Queue를 구독하는 Consumer에게 전달한다.
예를 들어 주문이 생성되었을 때 Order Service가 Product Service를 직접 호출하지 않고 RabbitMQ를 사용한다면 다음과 같은 구조가 된다.
Order Service
↓
RabbitMQ
↓
Product Queue
↓
Product Service
Order Service는 메시지를 생성해 RabbitMQ로 보내는 Producer 역할을 한다.
RabbitMQ는 전달받은 메시지를 Queue에 보관하고, Product Service는 Queue에서 메시지를 가져와 처리하는 Consumer 역할을 한다.
이러한 구조를 사용하면 Order Service와 Product Service가 서로를 직접 호출하지 않아도 되기 때문에 서비스 사이의 결합도를 낮출 수 있다.
RabbitMQ의 장점
1. 신뢰성
RabbitMQ는 메시지 승인 방식인 ACK(Acknowledgement) 를 지원한다.
Consumer가 메시지를 정상적으로 처리하면 ACK를 RabbitMQ에 전달하고, RabbitMQ는 해당 메시지가 성공적으로 처리되었다고 판단한다.
반대로 메시지 처리 중 오류가 발생하여 ACK가 전달되지 않으면 메시지를 다시 Queue에 넣거나 다른 Consumer에게 전달하도록 구성할 수 있다.
또한 Queue를 Durable로 설정하고 메시지를 Persistent 형태로 전송하면 RabbitMQ가 재시작되더라도 메시지가 유지되도록 구성할 수 있다.
단, RabbitMQ에 메시지를 보낸다고 해서 모든 메시지가 자동으로 디스크에 영구 저장되는 것은 아니다. Queue와 메시지의 영속성 설정이 함께 필요하다.
2. 유연성
RabbitMQ는 여러 메시지 전달 방식을 지원한다.
Exchange의 종류와 Routing Key를 이용하면 하나의 메시지를 특정 Queue로 전달하거나 여러 Queue에 동시에 전달하는 등의 구조를 만들 수 있다.
RabbitMQ는 기본적으로 AMQP를 사용하며, 플러그인을 통해 STOMP와 MQTT 등의 프로토콜도 지원한다.
3. 확장성
RabbitMQ는 클러스터링을 통해 여러 노드로 구성할 수 있다.
이를 통해 장애에 대비한 고가용성 환경을 구성하거나 여러 Consumer를 실행하여 메시지 처리 부하를 분산할 수 있다.
4. 관리 및 모니터링
RabbitMQ는 웹 기반 Management UI를 제공한다.
관리 화면을 통해 다음과 같은 정보를 확인할 수 있다.
- Exchange 목록
- Queue 목록
- Binding 관계
- Queue에 쌓인 메시지 수
- Producer와 Consumer 연결 상태
- 메시지 처리 속도
5. 비동기 처리
Producer는 Consumer의 작업이 끝날 때까지 기다리지 않고 RabbitMQ에 메시지를 전송한 뒤 자신의 작업을 계속할 수 있다.
따라서 처리 시간이 오래 걸리는 작업을 비동기로 분리할 수 있다.
RabbitMQ의 단점
1. 설정 및 운영 복잡성
Exchange, Queue, Routing Key, Binding, ACK, 재시도 정책 등을 올바르게 설정해야 한다.
서비스 규모가 커질수록 Queue와 메시지 흐름을 관리하는 작업도 복잡해질 수 있다.
2. 메시지 중복 처리 가능성
ACK가 전달되지 않거나 네트워크 오류가 발생하면 동일한 메시지가 다시 전달될 수 있다.
따라서 Consumer는 같은 메시지를 여러 번 받더라도 문제가 발생하지 않도록 멱등성을 고려해야 한다.
3. 성능 관리 필요
Queue에 메시지가 지나치게 많이 쌓이거나 Consumer의 처리 속도가 느리면 메시지 지연이 발생할 수 있다.
메시지 크기, Consumer 개수, Prefetch 설정 등을 적절하게 조절해야 한다.
4. 운영 비용
RabbitMQ 서버 구축뿐만 아니라 장애 대응, 모니터링, 백업, 클러스터 관리 등에 추가적인 운영 비용이 발생한다.
RabbitMQ의 구성 요소
메시지(Message)
RabbitMQ를 통해 전달되는 데이터 단위이다.
현재 프로젝트에서는 주문 ID, 상품 ID, 주문 수량 등의 정보를 가진 DeliveryMessage 객체가 메시지로 사용된다.
RabbitMQ는 Java 객체 자체를 이해하는 것이 아니라 바이트 데이터를 전달한다.
따라서 Jackson 기반의 MessageConverter를 사용하면 Java 객체를 JSON 형태로 변환하여 메시지로 전송하고, Consumer에서는 다시 Java 객체로 변환해 받을 수 있다.
DeliveryMessage 객체
↓
Jackson JSON 변환
↓
RabbitMQ 전송
↓
Jackson 객체 변환
↓
DeliveryMessage 객체
Producer
메시지를 생성하여 RabbitMQ로 전송하는 주체이다.
현재 프로젝트에서는 Order Service가 주문 생성 후 상품 재고 차감 메시지를 전송하므로 Producer 역할을 한다.
Spring에서는 RabbitTemplate.convertAndSend();을 이용해 메시지를 전송할 수 있다.
Queue
메시지가 Consumer에 의해 처리되기 전까지 대기하는 공간이다.
Producer가 전송한 메시지는 Queue에 들어가고, 해당 Queue를 구독하는 Consumer가 메시지를 가져가 처리한다.
현재 프로젝트에서는 다음과 같은 Queue를 사용한다.
productQueue
paymentQueue
errorOrderQueue
errorProductQueue
Consumer
Queue의 메시지를 가져와 처리하는 주체이다.
Spring에서는 @RabbitListener를 사용하여 특정 Queue를 구독할 수 있다.
메시지가 Queue에 들어오면 Spring이 @RabbitListener가 선언된 메서드를 자동으로 호출한다.
현재 프로젝트에서는 Product Service의 ProductEndpoint가 productQueue를 구독하는 Consumer 역할을 한다.
Exchange
Producer가 전송한 메시지를 어떤 Queue로 전달할지 결정하는 역할을 한다.
Exchange는 메시지의 Routing Key와 Binding 정보를 확인하여 적절한 Queue로 메시지를 라우팅한다.
Order Service
↓
Exchange
↓
Product Queue
↓
Product Service
Binding
Exchange와 Queue를 연결하는 설정이다.
Binding에는 Routing Key가 포함될 수 있으며, Exchange는 메시지의 Routing Key와 Binding Key를 비교하여 메시지를 전달할 Queue를 결정한다.
그림으로 표현하면 Exchange와 Queue 사이를 연결하는 화살표라고 볼 수 있다.
Exchange ── Binding ──> Product Queue
RabbitMQ에서 AMQP란?
AMQP는 Advanced Message Queuing Protocol의 약자로, 메시지 브로커를 통해 메시지를 전달하기 위한 표준 프로토콜이다.
메시지 생성, 전송, Queue 저장, 라우팅, ACK 처리 등의 방식을 표준화한다.
RabbitMQ는 AMQP 모델을 기반으로 Exchange, Queue, Binding, Routing Key를 사용하여 메시지를 전달한다.
AMQP의 주요 개념은 다음과 같다.
- Message
- Producer
- Consumer
- Exchange
- Queue
- Routing Key
- Binding
- ACK
전체 메시지 처리 구조
현재 프로젝트의 정상적인 주문 처리 흐름은 다음과 같다.
클라이언트
↓ HTTP 요청
OrderEndpoint
↓
OrderService
↓ RabbitTemplate
RabbitMQ
↓
productQueue
↓ @RabbitListener
ProductEndpoint
↓
ProductService
↓
상품 재고 차감
각 계층의 책임은 다음과 같이 구분할 수 있다.
OrderEndpoint
- 클라이언트의 HTTP 요청 수신
- OrderService 호출
OrderService
- 주문 생성
- 주문 저장
- RabbitTemplate을 이용한 메시지 전송
ProductEndpoint
- @RabbitListener를 이용한 Queue 구독
- 메시지 수신
- ProductService 호출
ProductService
- 상품 재고 차감
- 재고 복구
- 실제 비즈니스 로직 수행
RabbitMQ 구조 설정
다음 코드는 Exchange, Queue, Binding을 Spring Bean으로 등록하는 설정이다.
@Configuration
public class OrderApplicationQueueConfig {
@Bean
public MessageConverter messageConverter() {
return new JacksonJsonMessageConverter();
}
@Value("${message.exchange}")
private String exchange;
@Value("${message.queue.product}")
private String queueProduct;
@Value("${message.queue.payment}")
private String queuePayment;
@Value("${message.err.exchange}")
private String exchangeErr;
@Value("${message.queue.err.order}")
private String queueErrOrder;
@Value("${message.queue.err.product}")
private String queueErrProduct;
@Bean public TopicExchange exchange() { return new TopicExchange(exchange);}
@Bean public Queue queueProduct(){ return new Queue(queueProduct);}
@Bean public Queue queuePayment(){ return new Queue(queuePayment);}
// 아래 두개는 order -> exchange에서 productQueue로 가는 화살표를 바인딩하는 것과
// product -> exchange에서 paymentQueue로 가는 화살표를 바인딩 한 것
@Bean public Binding bindingProduct(){ return BindingBuilder.bind(queueProduct()).to(exchange()).with(queueProduct); }
@Bean public Binding bindingPayment(){ return BindingBuilder.bind(queuePayment()).to(exchange()).with(queuePayment); }
@Bean public TopicExchange exchangeErr() { return new TopicExchange(exchangeErr);}
@Bean public Queue queueErrOrder() { return new Queue(queueErrOrder); }
@Bean public Queue queueErrProduct() { return new Queue(queueErrProduct);}
@Bean public Binding bindingErrOrder(){ return BindingBuilder.bind(queueErrOrder()).to(exchangeErr()).with(queueErrOrder);}
@Bean public Binding bindingErrProduct(){ return BindingBuilder.bind(queueErrProduct()).to(exchangeErr()).with(queueErrProduct);}
이 설정은 메시지를 실제로 보내거나 처리하는 코드가 아니다.
RabbitMQ 내부에 다음과 같은 구조를 등록하는 역할을 한다.
Normal Exchange
├── productQueue
└── paymentQueue
Error Exchange
├── queueErrOrder
└── queueErrProduct
Binding은 Exchange와 Queue를 연결하고, 어떤 Routing Key를 가진 메시지가 어떤 Queue로 전달될지를 정의한다.
OrderEndpoint
@RestController
@RequiredArgsConstructor
public class OrderEndpoint {
private final OrderService orderService;
@GetMapping("/order/{orderId}")
public ResponseEntity<Order> getOrder(
@PathVariable("orderId") UUID orderId
) {
Order order = orderService.getOrder(orderId);
return ResponseEntity.ok(order);
}
@PostMapping("/order")
public ResponseEntity<Order> order(
@RequestBody OrderRequestDto orderRequestDto
) {
Order order = orderService.createOrder(orderRequestDto);
return ResponseEntity.ok(order);
}
}
OrderEndpoint는 RabbitMQ로 메시지를 직접 전송하는 역할이 아니다.
클라이언트의 HTTP 요청을 받아 OrderService에 전달하는 진입점 역할을 한다.
POST /order 요청이 들어오면 요청 데이터를 OrderService의 createOrder() 메서드로 전달한다.
Client
↓ HTTP POST
OrderEndpoint
↓
OrderService
즉, OrderEndpoint의 주요 책임은 다음과 같다.
- HTTP 요청 수신
- 요청 데이터 전달
- Service 호출
- 처리 결과 반환
실제 주문 생성과 RabbitMQ 메시지 전송은 OrderService에서 처리한다.
OrderService
@Slf4j
@Service
@RequiredArgsConstructor
public class OrderService {
@Value("${message.queue.product}")
private String productQueue;
private final RabbitTemplate rabbitTemplate;
/**
* 데이터베이스 대신 사용하는 임시 메모리 저장소
*/
private final Map<UUID, Order> orderStore = new HashMap<>();
public Order createOrder(
OrderEndpoint.OrderRequestDto orderRequestDto
) {
Order order = orderRequestDto.toOrder();
DeliveryMessage deliveryMessage =
orderRequestDto.toDeliveryMessage(order.getOrderId());
orderStore.put(order.getOrderId(), order);
log.info("SEND MESSAGE: {}", deliveryMessage);
rabbitTemplate.convertAndSend(
productQueue,
deliveryMessage
);
return order;
}
public Order getOrder(UUID orderId) {
return orderStore.get(orderId);
}
}
OrderService는 주문 생성과 메시지 전송을 담당한다.
주문 요청이 들어오면 먼저 Order 객체를 생성하고, Product Service로 전달할 DeliveryMessage 객체를 생성한다.
생성한 주문은 현재 데이터베이스 대신 Map에 임시로 저장한다.
그 후 RabbitTemplate의 convertAndSend()를 사용하여 메시지를 productQueue로 전송한다.
rabbitTemplate.convertAndSend(
productQueue,
deliveryMessage
);
RabbitTemplate은 Spring에서 RabbitMQ로 메시지를 전송하기 위해 사용하는 객체이다.
등록된 MessageConverter가 DeliveryMessage 객체를 JSON 등의 메시지 형태로 변환한 뒤 RabbitMQ로 전달한다.
따라서 OrderService의 역할은 Queue에서 메시지를 가져오는 것이 아니다.
현재 코드에서 OrderService는 메시지를 생성하고 전송하는 Producer 역할을 한다.
OrderService
↓ DeliveryMessage 생성
RabbitTemplate
↓
productQueue
현재 전송 방식에서 주의할 점
현재 코드는 다음과 같이 인자 두 개를 사용한다.
rabbitTemplate.convertAndSend(
productQueue,
deliveryMessage
);
이 방식은 직접 만든 TopicExchange의 이름을 지정하지 않는다.
Spring AMQP에서는 첫 번째 인자를 Routing Key로 사용하여 RabbitMQ의 기본 Exchange를 통해 이름이 같은 Queue로 메시지를 전달한다.
따라서 현재 흐름은 다음과 가깝다.
OrderService
↓
RabbitMQ 기본 Exchange
↓ Queue 이름과 Routing Key 일치
productQueue
앞에서 생성한 TopicExchange를 명시적으로 사용하려면 Exchange 이름과 Routing Key를 모두 전달해야 한다.
rabbitTemplate.convertAndSend(
exchangeName,
productQueue,
deliveryMessage
);
이 경우 메시지 흐름은 다음과 같다.
OrderService
↓
TopicExchange
↓ Binding 및 Routing Key 확인
productQueue
즉, 현재 코드처럼 Queue 이름만 전달하면 직접 만든 Topic Exchange를 거치지 않을 수 있다는 점을 구분해야 한다.
ProductEndpoint와 ProductService
Product Service에서는 ProductEndpoint와 ProductService가 서로 다른 역할을 담당한다.
ProductEndpoint는 RabbitMQ Queue를 구독하고 메시지를 수신하는 Consumer 역할을 한다.
ProductService는 수신한 메시지의 상품 정보를 검증하고, 처리 결과에 따라 다음 Queue로 메시지를 전달하는 비즈니스 로직 역할을 한다.
전체적인 역할은 다음과 같이 구분할 수 있다.
RabbitMQ Queue
↓
ProductEndpoint
↓
ProductService
↓
다음 Queue로 메시지 전달
ProductEndpoint
@Slf4j
@Component
@RequiredArgsConstructor
public class ProductEndpoint {
private final ProductService productService;
@RabbitListener(queues = "${message.queue.product}")
public void receiveMessage(
DeliveryMessage deliveryMessage
) {
log.info("PRODUCT RECEIVE: {}", deliveryMessage);
productService.reduceProductAmount(deliveryMessage);
}
@RabbitListener(queues = "${message.queue.err.product}")
public void receiveErrorMessage(
DeliveryMessage deliveryMessage
) {
log.info("PRODUCT ERROR RECEIVE: {}", deliveryMessage);
productService.rollbackProduct(deliveryMessage);
}
}
ProductEndpoint는 두 개의 Queue를 구독한다.
하나는 정상적인 상품 처리 메시지를 받는 productQueue이고, 다른 하나는 이후 단계에서 발생한 오류를 전달받는 errorProductQueue이다.
정상 메시지 수신
@RabbitListener(queues = "${message.queue.product}")
public void receiveMessage(DeliveryMessage deliveryMessage) {
log.info("PRODUCT RECEIVE: {}", deliveryMessage);
productService.reduceProductAmount(deliveryMessage);
}
다음 설정을 통해 productQueue를 지속적으로 구독한다.
@RabbitListener(queues = "${message.queue.product}")
productQueue에 메시지가 들어오면 Spring AMQP가 receiveMessage()를 자동으로 호출한다.
RabbitMQ에 저장된 메시지는 설정된 MessageConverter를 통해 DeliveryMessage 객체로 변환되어 메서드의 매개변수로 전달된다.
productService.reduceProductAmount(deliveryMessage);
ProductEndpoint는 상품 검증을 직접 수행하지 않고, 메시지를 ProductService의 reduceProductAmount()에 전달한다.
오류 메시지 수신
@RabbitListener(queues = "${message.queue.err.product}")
public void receiveErrorMessage(DeliveryMessage deliveryMessage) {
log.info("PRODUCT ERROR RECEIVE: {}", deliveryMessage);
productService.rollbackProduct(deliveryMessage);
}
다음 설정을 통해 errorProductQueue를 구독한다.
@RabbitListener(queues = "${message.queue.err.product}")
이 Queue는 Product Service 자체에서 처음 발생한 상품 검증 오류를 받기 위한 Queue가 아니다.
Product Service가 정상적으로 메시지를 다음 단계인 Payment Service로 전달한 후, Payment Service 등 이후 처리 단계에서 오류가 발생했을 때 Product Service로 오류 메시지를 되돌려 보내기 위한 Queue이다.
오류 메시지가 errorProductQueue에 들어오면 receiveErrorMessage()가 호출되고, 해당 메시지는 ProductService의 rollbackProduct()로 전달된다.
productService.rollbackProduct(deliveryMessage);
실제 시스템이라면 이 시점에 이전에 수행했던 상품 재고 차감 작업을 취소하는 보상 처리를 구현할 수 있다.
하지만 현재 코드에는 데이터베이스의 재고를 복구하는 로직이 구현되어 있지 않다.
ProductService
@Slf4j
@Service
@RequiredArgsConstructor
public class ProductService {
@Value("${message.queue.payment}")
private String paymentQueue;
@Value("${message.queue.err.order}")
private String errOrderQueue;
private final RabbitTemplate rabbitTemplate;
public void reduceProductAmount(DeliveryMessage deliveryMessage) {
Integer productId = deliveryMessage.getProductId();
Integer productQuantity = deliveryMessage.getProductQuantity();
if (productId != 1 || productQuantity > 1) {
this.rollbackProduct(deliveryMessage);
return;
}
rabbitTemplate.convertAndSend(
paymentQueue,
deliveryMessage
);
}
public void rollbackProduct(DeliveryMessage deliveryMessage) {
log.info("PRODUCT ROLLBACK!!!");
if (StringUtils.hasText(deliveryMessage.getErrorType())) {
deliveryMessage.setErrorType("PRODUCT ERROR");
}
rabbitTemplate.convertAndSend(
errOrderQueue,
deliveryMessage
);
}
}
ProductService는 ProductEndpoint에서 전달받은 메시지를 처리하고, 처리 결과에 따라 다음 Queue를 결정한다.
현재 코드의 주요 역할은 다음과 같다.
- 상품 ID와 상품 수량을 확인한다.
- 상품 정보가 유효한지 검증한다.
- 검증에 성공하면 메시지를 paymentQueue로 전달한다.
- 검증에 실패하면 오류 메시지를 errOrderQueue로 전달한다.
- 이후 단계에서 오류 메시지가 돌아오면 Order Service 쪽으로 오류를 전달한다.
상품 정보 검증
public void reduceProductAmount(DeliveryMessage deliveryMessage) {
Integer productId = deliveryMessage.getProductId();
Integer productQuantity = deliveryMessage.getProductQuantity();
reduceProductAmount()는 DeliveryMessage에서 상품 ID와 상품 수량을 가져온다.
Integer productId = deliveryMessage.getProductId();
Integer productQuantity = deliveryMessage.getProductQuantity();
이후 다음 조건을 사용하여 상품 정보를 검증한다.
if (productId != 1 || productQuantity > 1) {
this.rollbackProduct(deliveryMessage);
return;
}
현재 예제에서는 다음 중 하나라도 해당하면 상품 처리를 실패한 것으로 판단한다.
- 상품 ID가 1이 아닌 경우
- 요청한 상품 수량이 1보다 큰 경우
즉, 현재 코드는 실제 데이터베이스에서 상품이나 재고를 조회하는 것이 아니라, 고정된 값을 기준으로 상품 처리 성공과 실패를 구분하는 간단한 예제이다.
검증에 실패하면 rollbackProduct()를 호출한다.
this.rollbackProduct(deliveryMessage);
return;
return을 사용했기 때문에 이후의 paymentQueue 전송 코드는 실행되지 않는다.
상품 검증 성공
상품 ID와 수량이 조건에 맞으면 다음 코드가 실행된다.
rabbitTemplate.convertAndSend(
paymentQueue,
deliveryMessage
);
RabbitTemplate은 DeliveryMessage를 RabbitMQ의 paymentQueue로 전송한다.
이후 Payment Service가 paymentQueue를 구독하고 있다면 해당 메시지를 수신하여 결제 처리를 이어갈 수 있다.
ProductService
↓
RabbitTemplate
↓
paymentQueue
↓
Payment Service
현재 코드에서는 실제 상품 재고를 데이터베이스에서 차감하는 로직이 없다.
따라서 reduceProductAmount()라는 메서드 이름과 달리, 실제 동작은 다음과 같다.
상품 ID와 수량 확인
↓
상품 정보 검증
↓
paymentQueue로 메시지 전달
현재 예제 기준으로는 “상품 재고를 차감한다”라고 설명하기보다 상품 정보를 검증하고 다음 결제 단계로 메시지를 전달한다라고 설명하는 것이 정확하다.
상품 검증 실패
상품 ID나 상품 수량이 조건에 맞지 않으면 rollbackProduct()가 호출된다.
if (productId != 1 || productQuantity > 1) {
this.rollbackProduct(deliveryMessage);
return;
}
이 경우 메시지는 errorProductQueue로 이동하지 않는다.
reduceProductAmount()가 rollbackProduct()를 직접 호출하고, rollbackProduct()가 곧바로 메시지를 errOrderQueue로 전송한다.
따라서 Product Service 자체에서 발생한 상품 검증 실패 흐름은 다음과 같다.
productQueue
↓
ProductEndpoint.receiveMessage()
↓
ProductService.reduceProductAmount()
↓
상품 검증 실패
↓
ProductService.rollbackProduct()
↓
errOrderQueue
↓
Order Service
rollbackProduct()
public void rollbackProduct(DeliveryMessage deliveryMessage) {
log.info("PRODUCT ROLLBACK!!!");
if (StringUtils.hasText(deliveryMessage.getErrorType())) {
deliveryMessage.setErrorType("PRODUCT ERROR");
}
rabbitTemplate.convertAndSend(
errOrderQueue,
deliveryMessage
);
}
rollbackProduct()는 상품 처리 실패 또는 이후 단계의 오류 메시지를 처리한다.
먼저 다음 로그를 출력한다.
log.info("PRODUCT ROLLBACK!!!");
그다음 메시지의 errorType을 확인한다.
if (StringUtils.hasText(deliveryMessage.getErrorType())) {
deliveryMessage.setErrorType("PRODUCT ERROR");
}
마지막으로 메시지를 errOrderQueue로 전송한다.
rabbitTemplate.convertAndSend(
errOrderQueue,
deliveryMessage
);
이를 통해 Product Service보다 이전 단계인 Order Service에 오류 사실을 전달한다.
현재 rollbackProduct()에는 실제 상품 재고를 복구하는 코드가 없다.
따라서 현재 메서드가 실제로 수행하는 역할은 다음과 같다.
오류 로그 출력
↓
오류 정보 확인 또는 변경
↓
errOrderQueue로 메시지 전달
즉, 현재 예제에서 rollbackProduct()는 실제 재고 복구 메서드라기보다 오류 메시지를 이전 단계로 전달하는 메서드에 가깝다.
ProductEndpoint와 ProductService의 책임 구분
구성 요소역할
| ProductEndpoint | ProductQueue를 구독한다. |
| ProductEndpoint | 정상 메시지와 오류 메시지를 수신한다. |
| ProductEndpoint | 수신한 메시지를 ProductService에 전달한다. |
| ProductService | 상품 ID와 상품 수량을 검증한다. |
| ProductService | 정상 처리 시 메시지를 paymentQueue로 전달한다. |
| ProductService | 오류 처리 시 메시지를 errOrderQueue로 전달한다. |
ProductEndpoint는 RabbitMQ와 Product Service 사이를 연결하는 메시지 수신 계층이다.
ProductService는 메시지의 내용을 확인하고 정상 흐름으로 보낼지, 오류 흐름으로 보낼지를 결정하는 비즈니스 로직 계층이다.
정상 처리 흐름
정상적인 주문 처리 과정은 다음과 같다.
- 클라이언트가 주문 생성 요청을 보낸다.
- OrderEndpoint가 HTTP 요청을 받는다.
- OrderEndpoint가 OrderService.createOrder()를 호출한다.
- OrderService가 Order 객체를 생성한다.
- OrderService가 DeliveryMessage를 생성한다.
- OrderService가 RabbitTemplate을 사용하여 메시지를 전송한다.
- 메시지가 productQueue에 저장된다.
- ProductEndpoint.receiveMessage()가 메시지를 수신한다.
- ProductEndpoint가 ProductService.reduceProductAmount()를 호출한다.
- ProductService가 상품 ID와 상품 수량을 검증한다.
- 검증에 성공하면 메시지를 paymentQueue로 전송한다.
- Payment Service가 메시지를 받아 다음 처리를 진행한다.
Client
↓ HTTP 요청
OrderEndpoint
↓
OrderService
↓ RabbitTemplate
RabbitMQ 기본 Exchange
↓
productQueue
↓ @RabbitListener
ProductEndpoint.receiveMessage()
↓
ProductService.reduceProductAmount()
↓ 상품 검증 성공
RabbitTemplate
↓
paymentQueue
↓
Payment Service
현재 convertAndSend(queueName, message) 방식은 직접 생성한 TopicExchange 이름을 지정하지 않는다.
따라서 이 메시지는 RabbitMQ의 기본 Exchange를 통해 Queue 이름과 같은 Routing Key를 사용하여 전달된다.
Product Service 자체에서 검증이 실패한 흐름
Product Service가 상품 정보를 검증하는 과정에서 오류를 발견한 경우이다.
예를 들어 상품 ID가 1이 아니거나 상품 수량이 1보다 크면 이 흐름이 실행된다.
productQueue
↓
ProductEndpoint.receiveMessage()
↓
ProductService.reduceProductAmount()
↓
상품 ID 또는 수량 검증 실패
↓
ProductService.rollbackProduct()
↓ RabbitTemplate
errOrderQueue
↓
Order Service
이 경우에는 errorProductQueue나 Error Exchange를 거치지 않는다.
ProductService가 직접 rollbackProduct()를 호출하여 errOrderQueue로 메시지를 전달한다.
정리
현재 코드에서 각 구성 요소의 역할은 다음과 같다.
ProductEndpoint
- productQueue를 구독한다.
- errorProductQueue를 구독한다.
- RabbitMQ 메시지를 DeliveryMessage로 전달받는다.
- 정상 메시지는 reduceProductAmount()로 전달한다.
- 오류 메시지는 rollbackProduct()로 전달한다.
- RabbitMQ 메시지를 소비하는 Consumer 역할을 한다.
ProductService
- 상품 ID와 상품 수량을 검증한다.
- 상품 검증에 성공하면 메시지를 paymentQueue로 전달한다.
- 상품 검증에 실패하면 메시지를 errOrderQueue로 전달한다.
- 이후 단계에서 오류가 돌아오면 해당 오류 메시지를 errOrderQueue로 전달한다.
- 메시지의 다음 처리 방향을 결정하는 비즈니스 로직 역할을 한다.
현재 예제에는 실제 데이터베이스의 재고 차감과 재고 복구 로직이 없다.
따라서 전체 흐름은 다음과 같이 이해하는 것이 정확하다.
Order Service
↓
productQueue
↓
ProductEndpoint
↓
ProductService
↓
상품 정보 검증
↙ ↘
성공 실패
↓ ↓
paymentQueue errOrderQueue
Kafka란?
Kafka는 분산 이벤트 스트리밍 플랫폼(Distributed Event Streaming Platform) 으로, 대용량의 데이터를 실시간으로 수집하고 저장하며 처리하기 위해 사용된다.
RabbitMQ와 마찬가지로 Producer와 Consumer 구조를 사용하지만, Kafka는 이벤트 로그를 장기간 저장하고 여러 Consumer가 동일한 데이터를 반복해서 읽을 수 있는 구조라는 점에서 차이가 있다.
예를 들어 주문이 생성되면 Producer가 주문 이벤트를 Kafka의 Topic에 저장하고, 배송 서비스, 결제 서비스, 알림 서비스 등 여러 Consumer가 동일한 이벤트를 각각 읽어 자신의 작업을 수행할 수 있다.
Producer
↓
Kafka Topic
├── Delivery Service
├── Payment Service
└── Notification Service
이처럼 Kafka는 여러 서비스가 하나의 이벤트를 공유해야 하는 이벤트 기반(Event-Driven) 시스템에서 많이 사용된다.
Kafka의 구성 요소
메시지(Message)
Kafka를 통해 전달되는 데이터 단위이다.
현재 프로젝트의 DeliveryMessage와 같은 객체가 하나의 메시지가 될 수 있으며, Producer가 생성하여 Topic에 저장한다.
Producer
메시지를 생성하여 Kafka의 Topic으로 전송하는 역할을 수행한다.
Producer는 Topic 이름을 지정하여 메시지를 전송하며, 필요에 따라 Key를 함께 전달할 수 있다.
Topic
Kafka에서 메시지를 저장하는 논리적인 공간이다.
RabbitMQ의 Queue와 비슷한 역할을 하지만, 여러 개의 Partition으로 구성될 수 있다는 차이가 있다.
Producer가 전송한 메시지는 Topic에 저장되고, Consumer는 Topic으로부터 메시지를 읽어 처리한다.
Partition
Topic을 물리적으로 분할한 단위이다.
각 Partition은 독립적으로 메시지를 저장하며, 메시지마다 Offset이라는 고유 번호를 가진다.
Partition을 여러 개로 구성하면 여러 Consumer가 동시에 메시지를 처리할 수 있어 높은 처리량을 제공한다.
Order Topic
Partition 0
Offset 0
Offset 1
Offset 2
Partition 1
Offset 0
Offset 1
Offset 2
Offset
Partition 안에서 메시지를 구분하는 번호이다.
Consumer는 자신이 어디까지 읽었는지 Offset을 기록하여 이후에는 다음 메시지부터 이어서 읽을 수 있다.
Key
메시지를 어느 Partition에 저장할지 결정하는 값이다.
같은 Key를 가진 메시지는 동일한 Partition으로 저장되므로 메시지의 순서를 유지할 수 있다.
예를 들어 주문 번호를 Key로 사용하면 같은 주문과 관련된 이벤트는 항상 동일한 Partition에 저장된다.
Consumer
Topic에 저장된 메시지를 읽어 처리하는 역할을 한다.
Kafka에서는 Consumer가 Offset을 기준으로 데이터를 읽으며, Consumer Group을 사용하여 여러 Consumer가 Partition을 나누어 처리할 수 있다.
이를 통해 하나의 Topic을 여러 Consumer가 병렬로 처리할 수 있다.
Broker
Kafka 서버를 의미한다.
Broker는 Topic과 Partition을 저장하고 Producer와 Consumer 사이에서 메시지를 전달하는 역할을 수행한다.
여러 Broker를 클러스터로 구성하면 높은 처리량과 장애 대응이 가능하다.
ZooKeeper
기존 Kafka에서는 Broker의 메타데이터 관리와 클러스터 관리를 위해 ZooKeeper를 사용하였다.
Broker의 상태를 관리하고 리더 선출 등의 역할을 담당하였다.
다만 최근 Kafka(KRaft 모드)에서는 ZooKeeper 없이도 클러스터를 구성할 수 있으며, 최신 버전에서는 ZooKeeper 사용이 점차 줄어들고 있다.
Kafka와 RabbitMQ의 차이점
| 항목 | RabbitMQ | Kafka |
| 목적 | 메시지 전달(Message Broker)에 중점 | 대용량 이벤트 스트리밍(Event Streaming)에 중점 |
| 메시지 저장 장소 | Queue | Topic(Partition) |
| 저장 방식 | Consumer가 읽으면 일반적으로 Queue에서 제거 | 일정 기간 Topic에 저장되며 여러 Consumer가 반복해서 읽을 수 있음 |
| 메시지 지속성 | 주로 단기 저장(메모리 또는 디스크) | 장기 저장(디스크 기반 로그) |
| 메시지 순서 | Queue 단위로 순서 보장 | Partition 내부에서 순서 보장 |
| 병렬 처리 | 여러 Queue 또는 Consumer를 이용 | 여러 Partition을 이용한 병렬 처리 |
| 소비 방식 | Consumer가 메시지를 가져가면 일반적으로 제거 | Offset을 이용하여 원하는 위치부터 다시 읽을 수 있음 |
| 주요 사용 사례 | 작업 큐, 비동기 처리, 요청/응답, 마이크로서비스 간 메시지 전달 | 실시간 데이터 스트리밍, 로그 수집, 이벤트 소싱, 빅데이터 처리 |
RabbitMQ와 Kafka는 언제 사용할까?
RabbitMQ는 메시지를 안정적으로 전달하고 작업을 비동기로 처리하는 것에 강점이 있다.
예를 들어 주문 생성 후 재고 차감, 결제 요청, 이메일 발송과 같이 서비스 간 작업을 연결하는 경우 적합하다.
반면 Kafka는 대용량 이벤트를 저장하고 여러 서비스가 동시에 소비하는 구조에 적합하다.
예를 들어 사용자 행동 로그 분석, 실시간 모니터링, 이벤트 소싱(Event Sourcing), 데이터 파이프라인 구축과 같은 환경에서 많이 사용된다.
즉,
- RabbitMQ는 "메시지를 안전하게 전달하는 것"에 초점을 둔 메시지 브로커이고,
- Kafka는 "이벤트를 저장하고 실시간으로 처리하는 것"에 초점을 둔 분산 스트리밍 플랫폼이라고 이해하면 된다.
'심화_AI를 활용한 백엔드 아키텍처 심화 과정' 카테고리의 다른 글
| [내일배움캠프 사전캠프] 23회차 TIL(7/27 월) - 대규모 시스템 설계 (0) | 2026.07.27 |
|---|---|
| [내일배움캠프 사전캠프] 22회차 TIL(7/24 금) - Spring Session, Redis Cache (0) | 2026.07.24 |
| [내일배움캠프 사전캠프] 21회차 TIL(7/23 목) - Redis (0) | 2026.07.23 |
| [내일배움캠프 사전캠프] 20회차 TIL(7/22 수) - Docker, Docker-Compose (0) | 2026.07.22 |
| [내일배움캠프 사전캠프] 19회차 TIL(7/21 화) - Docker Image, Container (0) | 2026.07.21 |