最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

SpringBoot+本地消息表實現(xiàn)分布式最終一致性

 更新時間:2026年02月26日 09:56:53   作者:三水不滴  
本文主要介紹了SpringBoot+本地消息表實現(xiàn)分布式最終一致性,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧

在分布式系統(tǒng)中,跨服務數(shù)據(jù)一致性是核心難題 —— 例如電商下單時,需同時完成訂單創(chuàng)建、庫存扣減、積分增加等操作,若某一步失敗,可能導致數(shù)據(jù)不一致。主流的分布式事務方案(如 Seata、RocketMQ 事務消息)需依賴額外中間件,運維成本高,不適用于小型團隊或資源有限的場景。本文詳細拆解 SpringBoot + 本地消息表 + 定時補償 的輕量級方案,無需額外中間件,通過 “本地事務保障 + 定時重試” 實現(xiàn)分布式最終一致性。

一、分布式一致性痛點與方案選型

1. 核心痛點

分布式場景下,跨服務調用面臨三大問題,導致數(shù)據(jù)不一致:

  • 網(wǎng)絡異常:服務間通信超時、斷連,部分操作執(zhí)行成功部分失??;
  • 服務宕機:某服務執(zhí)行中宕機,未完成后續(xù)操作;
  • 數(shù)據(jù)沖突:并發(fā)場景下,多服務同時操作同一數(shù)據(jù)導致沖突。

傳統(tǒng) “同步調用” 方案(下單→扣庫存→加積分)一旦中間環(huán)節(jié)失敗,需手動回滾所有已執(zhí)行操作,代碼復雜且可靠性低。

2. 方案設計理念

本地消息表方案的核心是 “本地事務原子性 + 消息異步補償”:

  1. 將跨服務操作轉化為 “業(yè)務操作 + 記錄消息” 的本地事務,確保兩者同時成功或同時回滾;
  2. 通過定時任務掃描未完成的消息,異步調用目標服務;
  3. 采用重試機制處理臨時失敗,達到最大重試次數(shù)后標記為死信,人工介入處理。

整體流程示意圖:

服務A(訂單):
1. 開啟本地事務 → 2. 創(chuàng)建訂單 → 3. 記錄消息到本地消息表 → 4. 提交事務
5. 定時任務掃描消息表 → 6. 發(fā)送消息給服務B(庫存)→ 7. 服務B執(zhí)行扣庫存 → 8. 回調更新消息狀態(tài)

3. 技術選型與優(yōu)勢

組件選型理由
開發(fā)框架SpringBoot(快速整合組件,簡化配置)
數(shù)據(jù)庫MySQL(支持事務、行鎖,滿足本地消息表存儲需求)
ORM 框架MyBatis-Plus(簡化 CRUD 操作,支持樂觀鎖 / 悲觀鎖)
定時任務Spring Scheduled(輕量級,無需額外部署,滿足定時掃描需求)
重試策略指數(shù)退避算法(避免頻繁重試導致服務壓力,適配臨時故障場景)
冪等保障消息唯一標識(msgId)+ 目標服務接口冪等設計

核心優(yōu)勢:

  • 無額外依賴:無需部署消息隊列、分布式事務中間件,降低運維成本;
  • 實現(xiàn)簡單:基于本地事務和定時任務,開發(fā)門檻低,易落地;
  • 可靠性高:消息持久化存儲,宕機后可恢復,通過重試保障最終一致性;
  • 適配場景廣:適用于訂單創(chuàng)建、支付回調、庫存同步等非實時強一致場景。

二、核心實現(xiàn)細節(jié)

1. 數(shù)據(jù)庫設計(本地消息表)

本地消息表與業(yè)務表在同一數(shù)據(jù)庫,確保業(yè)務操作與消息記錄原子性:

-- 本地消息表
CREATE TABLE `local_message` (
  `id` bigint NOT NULL AUTO_INCREMENT COMMENT '主鍵ID',
  `msg_id` varchar(64) NOT NULL COMMENT '消息唯一標識(UUID)',
  `msg_type` varchar(32) NOT NULL COMMENT '消息類型(如:ORDER_CREATE、STOCK_DEDUCT)',
  `msg_content` text NOT NULL COMMENT '消息內容(JSON格式)',
  `target_service` varchar(64) NOT NULL COMMENT '目標服務(如:stock-service)',
  `target_url` varchar(255) NOT NULL COMMENT '目標接口URL',
  `status` tinyint NOT NULL COMMENT '消息狀態(tài):0-待處理,1-發(fā)送中,2-已完成,3-失?。ㄋ佬牛?,
  `retry_count` int NOT NULL DEFAULT '0' COMMENT '重試次數(shù)',
  `next_retry_time` datetime NOT NULL COMMENT '下次重試時間',
  `create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '創(chuàng)建時間',
  `update_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新時間',
  PRIMARY KEY (`id`),
  UNIQUE KEY `uk_msg_id` (`msg_id`),
  KEY `idx_status_next_retry_time` (`status`,`next_retry_time`) COMMENT '查詢待重試消息索引'
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='本地消息表';

-- 訂單表(示例業(yè)務表)
CREATE TABLE `t_order` (
  `id` bigint NOT NULL AUTO_INCREMENT COMMENT '訂單ID',
  `order_no` varchar(64) NOT NULL COMMENT '訂單編號',
  `user_id` bigint NOT NULL COMMENT '用戶ID',
  `product_id` bigint NOT NULL COMMENT '商品ID',
  `quantity` int NOT NULL COMMENT '購買數(shù)量',
  `amount` decimal(10,2) NOT NULL COMMENT '訂單金額',
  `status` tinyint NOT NULL COMMENT '訂單狀態(tài):0-待支付,1-已支付,2-已取消',
  `create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
  `update_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  PRIMARY KEY (`id`),
  UNIQUE KEY `uk_order_no` (`order_no`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='訂單表';

2. 核心實體類

(1)本地消息實體

@Data
@TableName("local_message")
public class LocalMessage {
    @TableId(type = IdType.AUTO)
    private Long id;
    @TableField("msg_id")
    private String msgId;
    @TableField("msg_type")
    private String msgType;
    @TableField("msg_content")
    private String msgContent;
    @TableField("target_service")
    private String targetService;
    @TableField("target_url")
    private String targetUrl;
    @TableField("status")
    private Integer status;
    @TableField("retry_count")
    private Integer retryCount;
    @TableField("next_retry_time")
    private LocalDateTime nextRetryTime;
    @TableField("create_time")
    private LocalDateTime createTime;
    @TableField("update_time")
    private LocalDateTime updateTime;

    // 消息狀態(tài)枚舉
    public static final Integer STATUS_PENDING = 0;    // 待處理
    public static final Integer STATUS_SENDING = 1;   // 發(fā)送中
    public static final Integer STATUS_COMPLETED = 2; // 已完成
    public static final Integer STATUS_FAILED = 3;    // 失?。ㄋ佬牛?
}

(2)訂單實體

@Data
@TableName("t_order")
public class Order {
    @TableId(type = IdType.AUTO)
    private Long id;
    @TableField("order_no")
    private String orderNo;
    @TableField("user_id")
    private Long userId;
    @TableField("product_id")
    private Long productId;
    @TableField("quantity")
    private Integer quantity;
    @TableField("amount")
    private BigDecimal amount;
    @TableField("status")
    private Integer status;
    @TableField("create_time")
    private LocalDateTime createTime;
    @TableField("update_time")
    private LocalDateTime updateTime;

    // 訂單狀態(tài)枚舉
    public static final Integer STATUS_PENDING_PAY = 0; // 待支付
    public static final Integer STATUS_PAID = 1;        // 已支付
    public static final Integer STATUS_CANCELLED = 2;   // 已取消
}

3. 業(yè)務操作與消息記錄(本地事務)

核心邏輯:訂單創(chuàng)建與消息記錄在同一本地事務中執(zhí)行,確保原子性:

@Service
public class OrderService {

    @Autowired
    private OrderMapper orderMapper;
    @Autowired
    private LocalMessageMapper localMessageMapper;
    @Autowired
    private RestTemplate restTemplate;

    // 最大重試次數(shù)
    private static final int MAX_RETRY_COUNT = 5;
    // 初始重試間隔(5秒)
    private static final int INIT_RETRY_INTERVAL_SECONDS = 5;

    /**
     * 創(chuàng)建訂單 + 記錄本地消息(扣庫存)
     */
    @Transactional(rollbackFor = Exception.class)
    public Order createOrder(OrderDTO orderDTO) {
        // 1. 生成訂單編號
        String orderNo = IdUtil.fastSimpleUUID();

        // 2. 構建訂單對象
        Order order = new Order();
        order.setOrderNo(orderNo);
        order.setUserId(orderDTO.getUserId());
        order.setProductId(orderDTO.getProductId());
        order.setQuantity(orderDTO.getQuantity());
        order.setAmount(orderDTO.getAmount());
        order.setStatus(Order.STATUS_PENDING_PAY);
        order.setCreateTime(LocalDateTime.now());
        order.setUpdateTime(LocalDateTime.now());

        // 3. 保存訂單(本地事務第一步)
        orderMapper.insert(order);

        // 4. 構建本地消息(扣庫存消息)
        LocalMessage localMessage = buildDeductStockMessage(order);

        // 5. 保存本地消息(本地事務第二步)
        localMessageMapper.insert(localMessage);

        return order;
    }

    /**
     * 構建扣庫存本地消息
     */
    private LocalMessage buildDeductStockMessage(Order order) {
        LocalMessage message = new LocalMessage();
        // 消息唯一標識(UUID)
        message.setMsgId(IdUtil.fastUUID());
        // 消息類型
        message.setMsgType("STOCK_DEDUCT");
        // 消息內容(JSON格式,包含商品ID、數(shù)量、訂單號)
        StockDeductDTO deductDTO = new StockDeductDTO();
        deductDTO.setProductId(order.getProductId());
        deductDTO.setQuantity(order.getQuantity());
        deductDTO.setOrderNo(order.getOrderNo());
        message.setMsgContent(JSON.toJSONString(deductDTO));
        // 目標服務(庫存服務)
        message.setTargetService("stock-service");
        // 目標接口URL(庫存服務扣庫存接口)
        message.setTargetUrl("http://stock-service/api/stock/deduct");
        // 消息狀態(tài):待處理
        message.setStatus(LocalMessage.STATUS_PENDING);
        // 初始重試次數(shù):0
        message.setRetryCount(0);
        // 下次重試時間:當前時間 + 初始間隔
        message.setNextRetryTime(LocalDateTime.now().plusSeconds(INIT_RETRY_INTERVAL_SECONDS));
        message.setCreateTime(LocalDateTime.now());
        message.setUpdateTime(LocalDateTime.now());

        return message;
    }
}

4. 定時補償任務(掃描 + 發(fā)送消息)

通過 Spring Scheduled 定時掃描待處理消息,執(zhí)行發(fā)送邏輯,失敗則更新重試次數(shù)和下次重試時間:

@Component
@EnableScheduling
public class LocalMessageScheduledTask {

    @Autowired
    private LocalMessageMapper localMessageMapper;
    @Autowired
    private RestTemplate restTemplate;

    // 定時任務執(zhí)行間隔(30秒,可通過配置中心動態(tài)調整)
    @Scheduled(cron = "0/30 * * * * ?")
    public void processPendingMessages() {
        log.info("開始掃描待處理本地消息");

        // 1. 查詢待處理且已到重試時間的消息(加悲觀鎖,防止并發(fā)處理)
        List<LocalMessage> pendingMessages = localMessageMapper.selectPendingMessages(LocalDateTime.now());
        if (CollectionUtils.isEmpty(pendingMessages)) {
            log.info("無待處理本地消息");
            return;
        }

        // 2. 遍歷消息,執(zhí)行發(fā)送邏輯
        for (LocalMessage message : pendingMessages) {
            try {
                // 2.1 更新消息狀態(tài)為“發(fā)送中”
                updateMessageStatus(message.getId(), LocalMessage.STATUS_SENDING);

                // 2.2 發(fā)送消息到目標服務
                boolean sendSuccess = sendMessageToTargetService(message);
                if (sendSuccess) {
                    // 2.3 發(fā)送成功:更新狀態(tài)為“已完成”
                    updateMessageToCompleted(message.getId());
                    log.info("消息發(fā)送成功:msgId={}", message.getMsgId());
                } else {
                    // 2.4 發(fā)送失?。禾幚碇卦囘壿?
                    handleRetry(message);
                }
            } catch (Exception e) {
                log.error("處理消息失?。簃sgId={}, 異常={}", message.getMsgId(), e.getMessage(), e);
                // 異常時同樣處理重試邏輯
                handleRetry(message);
            }
        }
    }

    /**
     * 發(fā)送消息到目標服務
     */
    private boolean sendSuccess = sendMessageToTargetService(message);
                if (sendSuccess) {
                    // 2.3 發(fā)送成功:更新狀態(tài)為“已完成”
                    updateMessageToCompleted(message.getId());
                    log.info("消息發(fā)送成功:msgId={}", message.getMsgId());
                } else {
                    // 2.4 發(fā)送失?。禾幚碇卦囘壿?
                    handleRetry(message);
                }
            } catch (Exception e) {
                log.error("處理消息失?。簃sgId={}, 異常={}", message.getMsgId(), e.getMessage(), e);
                // 異常時同樣處理重試邏輯
                handleRetry(message);
            }
        }
    }

    /**
     * 發(fā)送消息到目標服務
     */
    private boolean sendMessageToTargetService(LocalMessage message) {
        try {
            // 構建請求頭
            HttpHeaders headers = new HttpHeaders();
            headers.setContentType(MediaType.APPLICATION_JSON);

            // 構建請求體
            HttpEntity<String> requestEntity = new HttpEntity<>(message.getMsgContent(), headers);

            // 發(fā)送POST請求
            ResponseEntity<String> response = restTemplate.postForEntity(
                    message.getTargetUrl(),
                    requestEntity,
                    String.class
            );

            // 響應狀態(tài)碼200且返回成功標識,視為發(fā)送成功
            return response.getStatusCode().is2xxSuccessful() 
                    && "success".equals(JSON.parseObject(response.getBody()).getString("code"));
        } catch (Exception e) {
            log.error("調用目標服務失?。簃sgId={}, targetUrl={}", message.getMsgId(), message.getTargetUrl(), e);
            return false;
        }
    }

    /**
     * 處理重試邏輯(指數(shù)退避)
     */
    private void handleRetry(LocalMessage message) {
        int currentRetryCount = message.getRetryCount() + 1;

        // 超過最大重試次數(shù):標記為死信
        if (currentRetryCount >= MAX_RETRY_COUNT) {
            updateMessageToFailed(message.getId(), currentRetryCount);
            log.warn("消息達到最大重試次數(shù),標記為死信:msgId={}, retryCount={}", message.getMsgId(), currentRetryCount);
            return;
        }

        // 未超過最大重試次數(shù):計算下次重試時間(指數(shù)退避:5秒×2^重試次數(shù))
        long retryInterval = INIT_RETRY_INTERVAL_SECONDS * (long) Math.pow(2, currentRetryCount);
        LocalDateTime nextRetryTime = LocalDateTime.now().plusSeconds(retryInterval);

        // 更新重試次數(shù)和下次重試時間,狀態(tài)重置為“待處理”
        LocalMessage updateMsg = new LocalMessage();
        updateMsg.setId(message.getId());
        updateMsg.setRetryCount(currentRetryCount);
        updateMsg.setNextRetryTime(nextRetryTime);
        updateMsg.setStatus(LocalMessage.STATUS_PENDING);
        updateMsg.setUpdateTime(LocalDateTime.now());
        localMessageMapper.updateById(updateMsg);

        log.warn("消息重試處理:msgId={}, 當前重試次數(shù)={}, 下次重試時間={}",
                message.getMsgId(), currentRetryCount, nextRetryTime);
    }

    // ------------------- 消息狀態(tài)更新工具方法 -------------------
    private void updateMessageStatus(Long id, Integer status) {
        LocalMessage updateMsg = new LocalMessage();
        updateMsg.setId(id);
        updateMsg.setStatus(status);
        updateMsg.setUpdateTime(LocalDateTime.now());
        localMessageMapper.updateById(updateMsg);
    }

    private void updateMessageToCompleted(Long id) {
        LocalMessage updateMsg = new LocalMessage();
        updateMsg.setId(id);
        updateMsg.setStatus(LocalMessage.STATUS_COMPLETED);
        updateMsg.setUpdateTime(LocalDateTime.now());
        localMessageMapper.updateById(updateMsg);
    }

    private void updateMessageToFailed(Long id, Integer retryCount) {
        LocalMessage updateMsg = new LocalMessage();
        updateMsg.setId(id);
        updateMsg.setStatus(LocalMessage.STATUS_FAILED);
        updateMsg.setRetryCount(retryCount);
        updateMsg.setUpdateTime(LocalDateTime.now());
        localMessageMapper.updateById(updateMsg);
    }
}

5. Mapper 接口(MyBatis-Plus)

// 本地消息表Mapper
public interface LocalMessageMapper extends BaseMapper<LocalMessage> {

    /**
     * 查詢待處理/發(fā)送中且已到重試時間的消息(悲觀鎖)
     */
    @Select("SELECT * FROM local_message WHERE status IN (#{statusPending}, #{statusSending}) " +
            "AND next_retry_time <= #{currentTime} FOR UPDATE")
    List<LocalMessage> selectPendingMessages(
            @Param("statusPending") Integer statusPending,
            @Param("statusSending") Integer statusSending,
            @Param("currentTime") LocalDateTime currentTime);
}

// 訂單表Mapper
public interface OrderMapper extends BaseMapper<Order> {
    // 基礎CRUD由MyBatis-Plus自動生成
}

6. 目標服務冪等性實現(xiàn)(庫存服務)

庫存服務需保證接口冪等,防止重復扣減庫存:

@Service
public class StockService {

    @Autowired
    private StockMapper stockMapper;
    @Autowired
    private StockOperateLogMapper operateLogMapper;

    /**
     * 扣庫存(冪等實現(xiàn))
     */
    @Transactional(rollbackFor = Exception.class)
    public boolean deductStock(StockDeductDTO deductDTO) {
        String orderNo = deductDTO.getOrderNo();
        Long productId = deductDTO.getProductId();
        Integer quantity = deductDTO.getQuantity();

        // 1. 冪等校驗:查詢是否已處理該訂單的扣庫存請求
        StockOperateLog log = operateLogMapper.selectByOrderNo(orderNo);
        if (log != null) {
            // 已處理,直接返回成功
            return true;
        }

        // 2. 扣減庫存(悲觀鎖防止并發(fā)扣減)
        Stock stock = stockMapper.selectByProductIdForUpdate(productId);
        if (stock == null || stock.getQuantity() < quantity) {
            throw new RuntimeException("庫存不足,商品ID:" + productId);
        }
        stock.setQuantity(stock.getQuantity()-quantity);
        stockMapper.updateById(stock);

        // 3. 記錄操作日志(冪等標記)
        StockOperateLog operateLog = new StockOperateLog();
        operateLog.setOrderNo(orderNo);
        operateLog.setProductId(productId);
        operateLog.setQuantity(quantity);
        operateLog.setOperateType("DEDUCT");
        operateLog.setCreateTime(LocalDateTime.now());
        operateLogMapper.insert(operateLog);

        return true;
    }
}

三、生產(chǎn)環(huán)境優(yōu)化與最佳實踐

1. 性能優(yōu)化

  • 索引優(yōu)化:消息表添加idx_status_next_retry_time復合索引,提升掃描效率;
  • 批量處理:定時任務批量讀取消息(如每次 100 條),減少數(shù)據(jù)庫連接開銷;
  • 分庫分表:高并發(fā)場景下,按消息類型或業(yè)務 ID 分片,避免消息表成為瓶頸;
  • 異步化:消息發(fā)送邏輯異步執(zhí)行,避免阻塞定時任務主線程。

2. 可靠性增強

  • 死信處理:死信消息存入死信表,提供可視化界面人工重試;
  • 監(jiān)控告警:接入 Prometheus+Grafana,監(jiān)控消息發(fā)送成功率、重試次數(shù)、死信數(shù)量,異常時告警;
  • 分布式鎖:定時任務執(zhí)行時加分布式鎖(如 Redis 鎖),防止多實例重復掃描;
  • 日志鏈路追蹤:通過traceId串聯(lián)業(yè)務操作與消息發(fā)送日志,便于問題排查。

3. 冪等性進階方案

冪等實現(xiàn)方式適用場景實現(xiàn)要點
唯一標識訂單、支付等有唯一 ID 的場景消息 ID / 訂單號作為唯一鍵,查詢操作日志判斷是否已處理
樂觀鎖庫存扣減、余額更新等數(shù)值操作基于版本號更新,UPDATE ... SET version=version+1 WHERE version=#{version}
狀態(tài)機校驗流程類業(yè)務(如訂單狀態(tài)流轉)狀態(tài)變更需符合預設流程,如 “待支付”→“已支付”

4. 與其他方案對比

方案依賴復雜度實時性適用場景
本地消息表中(取決于掃描間隔)小型團隊、資源有限、輕量級場景
RocketMQ 事務消息消息中間件中大型團隊、高并發(fā)、需高可靠場景
Seata分布式事務框架強一致性需求、復雜業(yè)務場景

四、總結與演進方向

本地消息表方案以無中間件依賴、實現(xiàn)簡單、可靠性高的特點,成為小型團隊解決分布式最終一致性問題的首選。核心是通過本地事務保證業(yè)務與消息的原子性,定時補償實現(xiàn)消息必達,冪等設計防止重復處理。

適用場景

  • 對實時性要求不高的場景(如積分同步、物流狀態(tài)更新);
  • 不想引入復雜中間件的輕量級微服務系統(tǒng);
  • 訂單創(chuàng)建、庫存扣減、支付回調等核心業(yè)務場景。

演進方向

  1. 接入消息中間件:當業(yè)務規(guī)模擴大,可平滑遷移至 RabbitMQ/RocketMQ/Kafka 事務消息,提升實時性與吞吐量;
  2. 引入 Saga 模式:處理長鏈路業(yè)務,支持正向流程與反向補償;
  3. 云原生適配:結合 Serverless 架構,實現(xiàn)消息處理的彈性擴縮容。

到此這篇關于SpringBoot+本地消息表實現(xiàn)分布式最終一致性的文章就介紹到這了,更多相關SpringBoot 分布式最終一致性內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!

相關文章

最新評論

蒙城县| 麻江县| 峨边| 临沂市| 聂荣县| 泰兴市| 吉隆县| 会宁县| 柳河县| 万全县| 元朗区| 泸西县| 伊宁市| 保定市| 清水河县| 沙洋县| 齐河县| 简阳市| 秦安县| 东安县| 灵武市| 盘锦市| 洮南市| 玉溪市| 齐齐哈尔市| 大邑县| 张家界市| 浪卡子县| 井冈山市| 井陉县| 崇阳县| 杨浦区| 包头市| 太保市| 广元市| 长春市| 周口市| 通河县| 中西区| 丽水市| 大悟县|