在这里插入图片描述

肖哥弹架构 跟大家“弹弹” 分布式事务SAGA设计与实战应用,需要代码关注

欢迎 点赞,关注,转发。

关注公号Solomon肖哥弹架构获取更多精彩内容

历史热点文章

本文通过纯Java代码实现一个完整的电商订单履约系统,让读者初步对SAGA有理解,逐步将同步阻塞式流程改造为异步事件驱动的分布式架构。核心设计包括:

  1. Saga事务模式:通过事件队列拆解长事务,实现扣库存→支付→发货的最终一致性。
  2. 内嵌状态机:在订单类中硬编码状态流转规则,确保业务流程合法性。
  3. 异步持久化:使用文件存储订单状态,支持系统崩溃后恢复。
  4. 无框架设计:仅依赖Java基础API(如BlockingQueue),避免Spring等框架的隐含复杂性。

通过代码逐行解析与架构图,揭示分布式系统中状态管理事务补偿的核心设计思想。

一、系统架构图

在这里插入图片描述

架构说明

  1. 客户端:发起创建订单请求。
  2. OrderService
    • 核心逻辑(状态机+Saga协调器)
    • 持久化状态到OrderRepository
    • EventQueue推送事件
  3. EventQueue:事件总线(使用BlockingQueue实现)。
  4. 三方服务
    • InventoryService:同步扣减库存
    • PaymentService:异步支付(带回调)
    • ShippingService:同步发货

二、 时序图(正常流程)

在这里插入图片描述

三、 状态流转图

在这里插入图片描述

四、代码实现

1. 项目结构
src/
├── Main.java                  # 测试入口
├── model/
│   ├── Order.java             # 订单实体(含状态机逻辑)
│   ├── OrderEvent.java        # 事件定义
│   └── OrderItem.java         # 订单项
├── service/
│   ├── OrderService.java      # 核心逻辑(含Saga)
│   ├── InventoryService.java  # 库存服务
│   ├── PaymentService.java    # 支付服务
│   └── ShippingService.java   # 物流服务
└── repository/
    └── OrderRepository.java   # 订单持久化
2. 核心代码
2.1 订单实体(Order.java
/**
 * 订单实体类(内置状态机逻辑)
 */
public class Order {
    private String id;
    private OrderStatus status;
    private List<OrderItem> items;
    
    // 状态枚举
    public enum OrderStatus {
        CREATED, PAID, SHIPPED, COMPLETED, CANCELLED
    }

    /**
     * 状态变更(有限状态机核心逻辑)
     */
    public void updateStatus(OrderStatus newStatus) {
        if (!isValidTransition(this.status, newStatus)) {
            throw new IllegalStateException("非法状态流转: " + status + " -> " + newStatus);
        }
        this.status = newStatus;
    }

    /**
     * 校验状态流转是否合法
     */
    private boolean isValidTransition(OrderStatus current, OrderStatus target) {
        switch (current) {
            case CREATED:
                return target == OrderStatus.PAID || target == OrderStatus.CANCELLED;
            case PAID:
                return target == OrderStatus.SHIPPED || target == OrderStatus.CANCELLED;
            case SHIPPED:
                return target == OrderStatus.COMPLETED;
            default:
                return false; // 其他状态不可变
        }
    }
    // 其他getter/setter...
}
2.2 事件定义(OrderEvent.java
/**
 * 订单事件(用于驱动Saga)
 */
public class OrderEvent {
    public enum EventType {
        ORDER_CREATED,      // 订单创建
        PAYMENT_REQUESTED,  // 支付请求
        PAYMENT_COMPLETED,  // 支付完成
        SHIPPING_REQUESTED, // 发货请求
        ORDER_CANCELLED     // 订单取消
    }

    private String orderId;
    private EventType type;
    private Object data; // 附加数据

    public OrderEvent(String orderId, EventType type) {
        this.orderId = orderId;
        this.type = type;
    }
    // getter/setter...
}
2.3 订单服务(OrderService.java
/**
 * 订单服务(Saga协调器 + 状态机驱动)
 */
public class OrderService {
    private OrderRepository orderRepo = new OrderRepository();
    private InventoryService inventoryService = new InventoryService();
    private PaymentService paymentService = new PaymentService();
    private ShippingService shippingService = new ShippingService();
    private BlockingQueue<OrderEvent> eventQueue = new LinkedBlockingQueue<>(); // 事件队列

    /**
     * 创建订单(初始状态)
     */
    public Order createOrder(List<OrderItem> items) {
        Order order = new Order();
        order.setId(UUID.randomUUID().toString());
        order.setStatus(Order.OrderStatus.CREATED);
        order.setItems(items);
        orderRepo.save(order);

        // 触发订单创建事件(异步处理)
        eventQueue.add(new OrderEvent(order.getId(), OrderEvent.EventType.ORDER_CREATED));
        return order;
    }

    /**
     * 启动事件处理器(独立线程消费事件)
     */
    public void startEventProcessor() {
        new Thread(() -> {
            while (true) {
                try {
                    OrderEvent event = eventQueue.take(); // 阻塞获取事件
                    processEvent(event);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            }
        }).start();
    }

    /**
     * 处理事件(Saga的核心逻辑)
     */
    private void processEvent(OrderEvent event) {
        Order order = orderRepo.findById(event.getOrderId());
        if (order == null) return;

        try {
            switch (event.getType()) {
                case ORDER_CREATED:
                    // 扣减库存(Saga第一步)
                    inventoryService.deductStock(order.getItems());
                    // 触发支付请求(下一步)
                    eventQueue.add(new OrderEvent(order.getId(), OrderEvent.EventType.PAYMENT_REQUESTED));
                    break;

                case PAYMENT_REQUESTED:
                    // 异步支付(非阻塞)
                    paymentService.processPaymentAsync(order.getId(), (success) -> {
                        if (success) {
                            eventQueue.add(new OrderEvent(order.getId(), OrderEvent.EventType.PAYMENT_COMPLETED));
                        } else {
                            eventQueue.add(new OrderEvent(order.getId(), OrderEvent.EventType.ORDER_CANCELLED, "支付失败"));
                        }
                    });
                    break;

                case PAYMENT_COMPLETED:
                    order.updateStatus(Order.OrderStatus.PAID);
                    orderRepo.save(order);
                    // 触发发货(Saga下一步)
                    eventQueue.add(new OrderEvent(order.getId(), OrderEvent.EventType.SHIPPING_REQUESTED));
                    break;

                case SHIPPING_REQUESTED:
                    shippingService.shipOrder(order.getId());
                    order.updateStatus(Order.OrderStatus.SHIPPED);
                    orderRepo.save(order);
                    break;

                case ORDER_CANCELLED:
                    cancelOrder(order, (String) event.getData());
                    break;
            }
        } catch (Exception e) {
            // 任何失败触发取消
            cancelOrder(order, "处理事件失败: " + e.getMessage());
        }
    }

    /**
     * 取消订单(补偿逻辑)
     */
    private void cancelOrder(Order order, String reason) {
        if (order.getStatus() == Order.OrderStatus.PAID) {
            paymentService.refund(order.getId());
        }
        inventoryService.restoreStock(order.getItems());
        order.updateStatus(Order.OrderStatus.CANCELLED);
        orderRepo.save(order);
        System.out.println("订单取消: " + order.getId() + ", 原因: " + reason);
    }
}
2.4 支付服务(PaymentService.java
/**
 * 支付服务(模拟异步支付)
 */
public class PaymentService {
    /**
     * 异步支付(带回调)
     */
    public void processPaymentAsync(String orderId, Consumer<Boolean> callback) {
        new Thread(() -> {
            try {
                // 模拟支付处理耗时
                Thread.sleep(1000);
                boolean success = Math.random() > 0.3; // 70%成功
                callback.accept(success);
            } catch (InterruptedException e) {
                callback.accept(false);
            }
        }).start();
    }

    public void refund(String orderId) {
        System.out.println("[支付] 退款: " + orderId);
    }
}
2.5 订单存储(OrderRepository.java
/**
 * 订单持久化(模拟数据库)
 */
public class OrderRepository {
    private Map<String, Order> orders = new HashMap<>();
    private String dataFile = "orders.dat"; // 持久化文件

    public void save(Order order) {
        orders.put(order.getId(), order);
        serializeToFile(); // 每次更新后持久化
    }

    public Order findById(String id) {
        return orders.get(id);
    }

    // 序列化到文件
    private void serializeToFile() {
        try (ObjectOutputStream oos = new ObjectOutputStream(new FileOutputStream(dataFile))) {
            oos.writeObject(orders);
        } catch (IOException e) {
            e.printStackTrace();
        }
    }

    // 从文件加载
    @SuppressWarnings("unchecked")
    public void loadFromFile() {
        try (ObjectInputStream ois = new ObjectInputStream(new FileInputStream(dataFile))) {
            orders = (Map<String, Order>) ois.readObject();
        } catch (Exception e) {
            orders = new HashMap<>();
        }
    }
}
3. 测试运行(Main.java
public class Main {
    public static void main(String[] args) {
        // 初始化服务
        OrderService orderService = new OrderService();
        orderService.startEventProcessor(); // 启动事件处理器

        // 模拟下单
        Order order = orderService.createOrder(List.of(
            new OrderItem("product-1", 2),
            new OrderItem("product-2", 1)
        ));

        // 模拟系统重启(从文件恢复状态)
        try {
            Thread.sleep(1500); // 等待支付完成
            System.out.println("--- 模拟系统重启 ---");
            OrderService newOrderService = new OrderService();
            newOrderService.getOrderRepository().loadFromFile();
            newOrderService.startEventProcessor(); // 继续处理未完成订单
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
    }
}

输出

[库存] 扣减: [product-1 x 2, product-2 x 1]
[支付] 处理中: order-123
--- 模拟系统重启 ---
[支付] 完成: order-123 成功
[物流] 发货: order-123

五、关键设计解析

  1. Saga模式实现

    • 将长事务拆分为多个事件(ORDER_CREATEDPAYMENT_REQUESTEDSHIPPING_REQUESTED)。
    • 每个事件对应一个本地事务,失败时触发补偿(如取消订单)。
  2. 事件驱动异步化

    • 使用BlockingQueue作为事件总线,解耦服务调用。
    • 支付服务通过回调通知结果,避免阻塞主线程。
  3. 持久化与恢复

    • 订单状态变更后立即序列化到文件。
    • 系统重启后从文件加载未完成订单,继续处理。
  4. 状态机内嵌

    • Order类中的updateStatus方法硬编码状态流转规则,确保业务合法性。
Logo

电商企业物流数字化转型必备!快递鸟 API 接口,72 小时快速完成物流系统集成。全流程实战1V1指导,营造开放的API技术生态圈。

更多推荐