当前位置: 首页 > article >正文

高度细化的SAGA模式实现:基于Spring Boot与RabbitMQ的跨服务事务

场景与技术栈

场景:电商系统中的订单创建流程,涉及订单服务(Order Service)、库存服务(Inventory Service)、支付服务(Payment Service)。

技术栈:

  • Java 11

  • Spring Boot 2.7.0

  • MySQL 8.0

  • RabbitMQ

项目结构与组件
  • order-service

  • inventory-service

  • payment-service

  • saga-coordinator

代码实现

1. Saga Coordinator

@SpringBootApplicationpublic class SagaCoordinatorApplication {    public static void main(String[] args) {        SpringApplication.run(SagaCoordinatorApplication.class, args);    }}@Service@RequiredArgsConstructorpublic class SagaService {    private final SagaOrderService sagaOrderService;    private final SagaInventoryService sagaInventoryService;    private final SagaPaymentService sagaPaymentService;    private final SagaTransactionRepository sagaTransactionRepository;    public void placeOrder(Long orderId, Long productId, int quantity, double amount) throws ExecutionException, InterruptedException {        SagaTransaction transaction = new SagaTransaction();        transaction.setOrderId(orderId);        transaction.setProductId(productId);        transaction.setQuantity(quantity);        transaction.setAmount(amount);        transaction.setStatus(SagaTransaction.Status.INITIATED);        sagaTransactionRepository.save(transaction);        CompletableFuture.runAsync(() -> sagaOrderService.createOrder(transaction))            .thenRunAsync(() -> sagaInventoryService.decreaseInventory(transaction))            .thenRunAsync(() -> sagaPaymentService.charge(transaction))            .exceptionally(ex -> {                rollback(transaction);                throw new RuntimeException("Saga failed", ex);            });    }    private void rollback(SagaTransaction transaction) {        CompletableFuture.runAsync(() -> sagaPaymentService.refund(transaction))            .thenRunAsync(() -> sagaInventoryService.increaseInventory(transaction))            .thenRunAsync(() -> sagaOrderService.cancelOrder(transaction))            .exceptionally(ex -> {                transaction.setStatus(SagaTransaction.Status.FAILED);                sagaTransactionRepository.save(transaction);                throw new RuntimeException("Rollback failed", ex);            });    }}

2. SagaOrderService

@Servicepublic class SagaOrderService {    private final OrderService orderService;    private final RabbitTemplate rabbitTemplate;    @Autowired    public SagaOrderService(OrderService orderService, RabbitTemplate rabbitTemplate) {        this.orderService = orderService;        this.rabbitTemplate = rabbitTemplate;    }    public void createOrder(SagaTransaction transaction) {        // 创建订单逻辑        // ...        // 发送创建订单完成的消息        rabbitTemplate.convertAndSend("order.create", transaction);    }    public void cancelOrder(SagaTransaction transaction) {        // 取消订单逻辑        // ...        // 发送取消订单完成的消息        rabbitTemplate.convertAndSend("order.cancel", transaction);    }}

3. SagaInventoryService

@Servicepublic class SagaInventoryService {    private final InventoryService inventoryService;    private final RabbitTemplate rabbitTemplate;    @Autowired    public SagaInventoryService(InventoryService inventoryService, RabbitTemplate rabbitTemplate) {        this.inventoryService = inventoryService;        this.rabbitTemplate = rabbitTemplate;    }    public void decreaseInventory(SagaTransaction transaction) {        // 扣减库存逻辑        // ...        // 发送扣减库存完成的消息        rabbitTemplate.convertAndSend("inventory.decrease", transaction);    }    public void increaseInventory(SagaTransaction transaction) {        // 增加库存逻辑        // ...        // 发送增加库存完成的消息        rabbitTemplate.convertAndSend("inventory.increase", transaction);    }}

4. SagaPaymentService

@Servicepublic class SagaPaymentService {    private final PaymentService paymentService;    private final RabbitTemplate rabbitTemplate;    @Autowired    public SagaPaymentService(PaymentService paymentService, RabbitTemplate rabbitTemplate) {        this.paymentService = paymentService;        this.rabbitTemplate = rabbitTemplate;    }    public void charge(SagaTransaction transaction) {        // 收款逻辑        // ...        // 发送收款完成的消息        rabbitTemplate.convertAndSend("payment.charge", transaction);    }    public void refund(SagaTransaction transaction) {        // 退款逻辑        // ...        // 发送退款完成的消息        rabbitTemplate.convertAndSend("payment.refund", transaction);    }}
SagaTransaction Entity
@Entitypublic class SagaTransaction {    @Id    @GeneratedValue(strategy = GenerationType.IDENTITY)    private Long id;    private Long orderId;    private Long productId;    private int quantity;    private double amount;    private Status status; // INITIATED, COMPLETED, ROLLING_BACK, ROLLED_BACK, FAILED    // Getters & Setters}
消息队列监听

在每个微服务中,需要添加消息队列的监听器,以便在接收到消息时执行相应的操作。例如,在order-service中:

@Componentpublic class OrderSagaListener {    private final SagaOrderService sagaOrderService;    @Autowired    public OrderSagaListener(SagaOrderService sagaOrderService) {        this.sagaOrderService = sagaOrderService;    }    @RabbitListener(queues = "order.create")    public void handleOrderCreate(SagaTransaction transaction) {        sagaOrderService.createOrder(transaction);    }    @RabbitListener(queues = "order.cancel")    public void handleOrderCancel(SagaTransaction transaction) {        sagaOrderService.cancelOrder(transaction);    }}
额外细节
  • 为确保事务的一致性,可以使用RabbitMQ的发布确认(Publisher Confirms)机制。

  • 每个微服务的数据库事务应该使用@Transactional注解来保证ACID属性。

  • 需要设计失败重试和事务状态检查机制,确保在故障恢复时能够正确地执行补偿操作。

通过上述设计,SAGA模式与RabbitMQ的结合,不仅能够处理跨服务的事务,还能够通过消息队列实现服务解耦和消息的异步处理,提高系统的稳定性和可扩展性。


http://www.kler.cn/news/329307.html

相关文章:

  • 甄选范文“论软件的可靠性设计”,软考高级论文,系统架构设计师论文
  • Vue页面,基础配置
  • 机器学习模型评估
  • Web APIs 4:日期对象、时间戳、节点操作、swiper插件
  • VS code user setting 与 workspace setting 的区别
  • 前端规范工程-2:JS代码规范(Prettier + ESLint)
  • consul 介绍与使用,以及spring boot 项目的集成
  • Servlet——springMvc底层原理
  • 苏州 数字化科技展厅展馆-「世岩科技」一站式服务商
  • RD-Agent Windows安装教程
  • 第一节- C++入门
  • 图论(dfs系列) 9/27
  • telnet发送邮件教程:安全配置与操作指南?
  • 09_OpenCV彩色图片直方图
  • git cherry-pick作用
  • ClientlocaleController
  • Dify:一个简化大模型应用的开源平台
  • python中统计奇数和偶数之和
  • 如何在 Kubernetes 集群中安装和配置 OpenEBS 持久化块存储?
  • 卸载apt-get 安装的PostgreSQL版本
  • TI DSP TMS320F280025 Note14:模数转换器ADC原理分析与应用
  • gd32jlink第一次下载可以用,重新上电后不行了
  • 第十四周:机器学习笔记
  • 10 个最佳 Golang 库
  • 常见的 C++ 库介绍
  • C++学习笔记----8、掌握类与对象(二)---- 成员函数的更多知识(1)
  • 昇思MindSpore进阶教程--下沉模式
  • C++和OpenGL实现3D游戏编程【连载12】——游戏中音效的使用
  • DTH11温湿度传感器
  • python学习记录5