
어제에 이은 추가 실습으로 RabbitMQ를 사용해 SAGA 패턴을 구현해보았다.

Product Application은 Consumer의 입장으로서 Order Application으로부터 메시지 큐를 통해 요청을 받는다.
동시에 Producer의 입장으로 Payment Application에 메시지 큐를 통해 요청을 보낸다.
Product Application에서 에러가 발생하였을 경우 Order Application에 발생한 에러 메시지와 함께 에러 메시지 큐를 통해 요청을 보낸다.
동시에 Consumer의 입장으로서 Payment Application에서 에러가 발생한 탓에 에러 메시지 큐를 통해 받아온 요청이 있을 경우에도 Order Application에 발생한 에러 메시지와 함께 에러 메시지 큐를 통해 요청을 보낸다.
dependencies {
implementation 'org.springframework.boot:spring-boot-starter-amqp'
implementation 'org.springframework.boot:spring-boot-starter-web'
compileOnly 'org.projectlombok:lombok'
annotationProcessor 'org.projectlombok:lombok'
testImplementation 'org.springframework.boot:spring-boot-starter-test'
testImplementation 'org.springframework.amqp:spring-rabbit-test'
testRuntimeOnly 'org.junit.platform:junit-platform-launcher'
}
spring.application.name=order
message.exchange=market
message.queue.product=market.product
message.queue.payment=market.payment
message.err.exchange=market.err
message.queue.err.order=market.err.order
message.queue.err.product=market.err.product
spring.rabbitmq.host=localhost
spring.rabbitmq.port=5672
spring.rabbitmq.username=guest
spring.rabbitmq.password=guest
import lombok.*;
import java.util.UUID;
@Data
@Builder
@ToString
@AllArgsConstructor
@NoArgsConstructor
public class DeliveryMessage {
private UUID orderId;
private UUID paymentId;
private String userId;
private Integer productId;
private Integer productQuantity;
private Integer payAmount;
private String errorType;
}
import lombok.Builder;
import lombok.Data;
import lombok.ToString;
import java.util.UUID;
@Builder
@Data
@ToString
public class Order {
private UUID orderId;
private String userId;
private String orderStatus;
private String errorType;
public void cancelOrder(String receiveErrorType) {
orderStatus = "CANCEL";
errorType = receiveErrorType;
}
}
import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.core.TopicExchange;
import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class OrderApplicationQueueConfig {
@Bean
public Jackson2JsonMessageConverter producerJackson2MessageConverter() {
return new Jackson2JsonMessageConverter();
}
@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); }
@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); }
}
import lombok.Data;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.*;
import java.util.UUID;
@Slf4j
@RestController
@RequiredArgsConstructor
public class OrderEndpoint {
private final OrderService orderService;
private final RabbitTemplate rabbitTemplate;
@GetMapping("order/{orderId}")
public ResponseEntity<Order> getOrder(@PathVariable 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);
}
@RabbitListener(queues = "${message.queue.err.order}")
public void errOrder(DeliveryMessage message) {
log.info("ERROR RECEIVE !!!");
orderService.rollbackOrder(message);
}
@Data
public static class OrderRequestDto {
private String userId;
private Integer productId;
private Integer productQuantity;
private Integer payAmount;
public Order toOrder (){
return Order.builder()
.orderId(UUID.randomUUID())
.userId(userId)
.orderStatus("RECEIPT")
.build();
}
public DeliveryMessage toDeliveryMessage(UUID orderId){
return DeliveryMessage.builder()
.orderId(orderId)
.productId(productId)
.productQuantity(productQuantity)
.payAmount(payAmount)
.build();
}
}
}
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import java.util.HashMap;
import java.util.Map;
import java.util.UUID;
@Slf4j
@Service
@RequiredArgsConstructor
public class OrderService {
@Value("${message.queue.product}")
private String productQueue;
private final RabbitTemplate rabbitTemplate;
private 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.toString());
rabbitTemplate.convertAndSend(productQueue, deliveryMessage);
return order;
}
public void rollbackOrder(DeliveryMessage message) {
Order order = orderStore.get(message.getOrderId());
order.cancelOrder(message.getErrorType());
log.info(order.toString());
}
public Order getOrder(UUID orderId) {
return orderStore.get(orderId);
}
}
dependencies {
// Jackson 의존성
implementation 'com.fasterxml.jackson.core:jackson-databind'
implementation 'com.fasterxml.jackson.core:jackson-core'
implementation 'com.fasterxml.jackson.core:jackson-annotations'
implementation 'org.springframework.boot:spring-boot-starter-amqp'
compileOnly 'org.projectlombok:lombok'
annotationProcessor 'org.projectlombok:lombok'
testImplementation 'org.springframework.boot:spring-boot-starter-test'
testImplementation 'org.springframework.amqp:spring-rabbit-test'
testRuntimeOnly 'org.junit.platform:junit-platform-launcher'
}
spring.application.name=product
message.exchange=market
message.queue.product=market.product
message.queue.payment=market.payment
message.err.exchange=market.err
message.queue.err.order=market.err.order
message.queue.err.product=market.err.product
spring.rabbitmq.host=localhost
spring.rabbitmq.port=5672
spring.rabbitmq.username=guest
spring.rabbitmq.password=guest
import lombok.*;
import java.util.UUID;
@Data
@Builder
@ToString
@AllArgsConstructor
@NoArgsConstructor
public class DeliveryMessage {
private UUID orderId;
private UUID paymentId;
private String userId;
private Integer productId;
private Integer productQuantity;
private Integer payAmount;
private String errorType;
}
import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class ProductApplicationQueueConfig {
@Bean
public Jackson2JsonMessageConverter producerJackson2MessageConverter() {
return new Jackson2JsonMessageConverter();
}
}
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Slf4j
@Component
@RequiredArgsConstructor
public class ProductEndpoint {
private final ProductService productService;
@RabbitListener(queues = "${message.queue.product}")
public void receiveMessage(DeliveryMessage deliveryMessage) {
productService.reduceProductAmount(deliveryMessage);
log.info("PRODUCT RECEIVE:{}", deliveryMessage.toString());
}
@RabbitListener(queues="${message.queue.err.product}")
public void receiveErrorMessage(DeliveryMessage deliveryMessage) {
log.info("ERROR RECEIVE !!!");
productService.rollbackProduct(deliveryMessage);
}
}
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils;
@Slf4j
@Service
@RequiredArgsConstructor
public class ProductService {
private final RabbitTemplate rabbitTemplate;
@Value("${message.queue.payment}")
private String paymentQueue;
@Value("${message.queue.err.order}")
private String orderErrorQueue;
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(orderErrorQueue, deliveryMessage);
}
}
dependencies {
// Jackson 의존성
implementation 'com.fasterxml.jackson.core:jackson-databind'
implementation 'com.fasterxml.jackson.core:jackson-core'
implementation 'com.fasterxml.jackson.core:jackson-annotations'
implementation 'org.springframework.boot:spring-boot-starter-amqp'
compileOnly 'org.projectlombok:lombok'
annotationProcessor 'org.projectlombok:lombok'
testImplementation 'org.springframework.boot:spring-boot-starter-test'
testImplementation 'org.springframework.amqp:spring-rabbit-test'
testRuntimeOnly 'org.junit.platform:junit-platform-launcher'
}
spring.application.name=payment
message.exchange=market
message.queue.product=market.product
message.queue.payment=market.payment
message.err.exchange=market.err
message.queue.err.order=market.err.order
message.queue.err.product=market.err.product
spring.rabbitmq.host=localhost
spring.rabbitmq.port=5672
spring.rabbitmq.username=guest
spring.rabbitmq.password=guest
import lombok.*;
import java.util.UUID;
@Data
@Builder
@ToString
@AllArgsConstructor
@NoArgsConstructor
public class DeliveryMessage {
private UUID orderId;
private UUID paymentId;
private String userId;
private Integer productId;
private Integer productQuantity;
private Integer payAmount;
private String errorType;
}
import lombok.Builder;
import lombok.Data;
import java.util.UUID;
@Data
@Builder
public class Payment {
private UUID paymentId;
private String userId;
private String payAmount;
private String payStatus;
}
import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class PaymentApplicationQueueConfig {
@Bean
public Jackson2JsonMessageConverter producerJackson2MessageConverter() {
return new Jackson2JsonMessageConverter();
}
}
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Slf4j
@Component
@RequiredArgsConstructor
public class PaymentEndpoint {
private final PaymentService paymentService;
@RabbitListener(queues = "${message.queue.payment}")
public void receiveMessage(DeliveryMessage deliveryMessage) {
log.info("PAYMENT RECEIVE : {}", deliveryMessage.toString());
paymentService.createPayment(deliveryMessage);
}
}
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import java.util.UUID;
@Slf4j
@Service
@RequiredArgsConstructor
public class PaymentService {
private final RabbitTemplate rabbitTemplate;
@Value("${message.queue.err.product}")
private String productErrorQueue;
public void createPayment(DeliveryMessage deliveryMessage) {
Payment payment = Payment.builder()
.paymentId(UUID.randomUUID())
.userId(deliveryMessage.getUserId())
.payStatus("SUCCESS").build();
Integer payAmount = deliveryMessage.getPayAmount();
if (payAmount >= 10000) {
log.error("Payment amount exceeds limit: {}", payAmount);
payment.setPayStatus("CANCEL");
deliveryMessage.setErrorType("PAYMENT_LIMIT_EXCEEDED");
this.rollbackPayment(deliveryMessage);
}
}
public void rollbackPayment(DeliveryMessage deliveryMessage) {
log.info("PAYMENT ROLLBACK !!!");
rabbitTemplate.convertAndSend(productErrorQueue, deliveryMessage);
}
}