SpringBoot+本地消息表實現分布式最終一致性
在分布式系統(tǒng)中,跨服務數據一致性是核心難題 —— 例如電商下單時,需同時完成訂單創(chuàng)建、庫存扣減、積分增加等操作,若某一步失敗,可能導致數據不一致。主流的分布式事務方案(如 Seata、RocketMQ 事務消息)需依賴額外中間件,運維成本高,不適用于小型團隊或資源有限的場景。本文詳細拆解 SpringBoot + 本地消息表 + 定時補償 的輕量級方案,無需額外中間件,通過 “本地事務保障 + 定時重試” 實現分布式最終一致性。
一、分布式一致性痛點與方案選型
1. 核心痛點
分布式場景下,跨服務調用面臨三大問題,導致數據不一致:
- 網絡異常:服務間通信超時、斷連,部分操作執(zhí)行成功部分失??;
- 服務宕機:某服務執(zhí)行中宕機,未完成后續(xù)操作;
- 數據沖突:并發(fā)場景下,多服務同時操作同一數據導致沖突。
傳統(tǒng) “同步調用” 方案(下單→扣庫存→加積分)一旦中間環(huán)節(jié)失敗,需手動回滾所有已執(zhí)行操作,代碼復雜且可靠性低。
2. 方案設計理念
本地消息表方案的核心是 “本地事務原子性 + 消息異步補償”:
- 將跨服務操作轉化為 “業(yè)務操作 + 記錄消息” 的本地事務,確保兩者同時成功或同時回滾;
- 通過定時任務掃描未完成的消息,異步調用目標服務;
- 采用重試機制處理臨時失敗,達到最大重試次數后標記為死信,人工介入處理。
整體流程示意圖:
服務A(訂單):
1. 開啟本地事務 → 2. 創(chuàng)建訂單 → 3. 記錄消息到本地消息表 → 4. 提交事務
5. 定時任務掃描消息表 → 6. 發(fā)送消息給服務B(庫存)→ 7. 服務B執(zhí)行扣庫存 → 8. 回調更新消息狀態(tài)
3. 技術選型與優(yōu)勢
| 組件 | 選型理由 |
|---|---|
| 開發(fā)框架 | SpringBoot(快速整合組件,簡化配置) |
| 數據庫 | MySQL(支持事務、行鎖,滿足本地消息表存儲需求) |
| ORM 框架 | MyBatis-Plus(簡化 CRUD 操作,支持樂觀鎖 / 悲觀鎖) |
| 定時任務 | Spring Scheduled(輕量級,無需額外部署,滿足定時掃描需求) |
| 重試策略 | 指數退避算法(避免頻繁重試導致服務壓力,適配臨時故障場景) |
| 冪等保障 | 消息唯一標識(msgId)+ 目標服務接口冪等設計 |
核心優(yōu)勢:
- 無額外依賴:無需部署消息隊列、分布式事務中間件,降低運維成本;
- 實現簡單:基于本地事務和定時任務,開發(fā)門檻低,易落地;
- 可靠性高:消息持久化存儲,宕機后可恢復,通過重試保障最終一致性;
- 適配場景廣:適用于訂單創(chuàng)建、支付回調、庫存同步等非實時強一致場景。
二、核心實現細節(jié)
1. 數據庫設計(本地消息表)
本地消息表與業(yè)務表在同一數據庫,確保業(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 '重試次數', `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 '購買數量', `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;
// 最大重試次數
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、數量、訂單號)
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);
// 初始重試次數: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ā)送邏輯,失敗則更新重試次數和下次重試時間:
@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("處理消息失敗:msgId={}, 異常={}", 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;
}
}
/**
* 處理重試邏輯(指數退避)
*/
private void handleRetry(LocalMessage message) {
int currentRetryCount = message.getRetryCount() + 1;
// 超過最大重試次數:標記為死信
if (currentRetryCount >= MAX_RETRY_COUNT) {
updateMessageToFailed(message.getId(), currentRetryCount);
log.warn("消息達到最大重試次數,標記為死信:msgId={}, retryCount={}", message.getMsgId(), currentRetryCount);
return;
}
// 未超過最大重試次數:計算下次重試時間(指數退避:5秒×2^重試次數)
long retryInterval = INIT_RETRY_INTERVAL_SECONDS * (long) Math.pow(2, currentRetryCount);
LocalDateTime nextRetryTime = LocalDateTime.now().plusSeconds(retryInterval);
// 更新重試次數和下次重試時間,狀態(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={}, 當前重試次數={}, 下次重試時間={}",
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. 目標服務冪等性實現(庫存服務)
庫存服務需保證接口冪等,防止重復扣減庫存:
@Service
public class StockService {
@Autowired
private StockMapper stockMapper;
@Autowired
private StockOperateLogMapper operateLogMapper;
/**
* 扣庫存(冪等實現)
*/
@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;
}
}
三、生產環(huán)境優(yōu)化與最佳實踐
1. 性能優(yōu)化
- 索引優(yōu)化:消息表添加
idx_status_next_retry_time復合索引,提升掃描效率; - 批量處理:定時任務批量讀取消息(如每次 100 條),減少數據庫連接開銷;
- 分庫分表:高并發(fā)場景下,按消息類型或業(yè)務 ID 分片,避免消息表成為瓶頸;
- 異步化:消息發(fā)送邏輯異步執(zhí)行,避免阻塞定時任務主線程。
2. 可靠性增強
- 死信處理:死信消息存入死信表,提供可視化界面人工重試;
- 監(jiān)控告警:接入 Prometheus+Grafana,監(jiān)控消息發(fā)送成功率、重試次數、死信數量,異常時告警;
- 分布式鎖:定時任務執(zhí)行時加分布式鎖(如 Redis 鎖),防止多實例重復掃描;
- 日志鏈路追蹤:通過
traceId串聯業(yè)務操作與消息發(fā)送日志,便于問題排查。
3. 冪等性進階方案
| 冪等實現方式 | 適用場景 | 實現要點 |
|---|---|---|
| 唯一標識 | 訂單、支付等有唯一 ID 的場景 | 消息 ID / 訂單號作為唯一鍵,查詢操作日志判斷是否已處理 |
| 樂觀鎖 | 庫存扣減、余額更新等數值操作 | 基于版本號更新,UPDATE ... SET version=version+1 WHERE version=#{version} |
| 狀態(tài)機校驗 | 流程類業(yè)務(如訂單狀態(tài)流轉) | 狀態(tài)變更需符合預設流程,如 “待支付”→“已支付” |
4. 與其他方案對比
| 方案 | 依賴 | 復雜度 | 實時性 | 適用場景 |
|---|---|---|---|---|
| 本地消息表 | 無 | 低 | 中(取決于掃描間隔) | 小型團隊、資源有限、輕量級場景 |
| RocketMQ 事務消息 | 消息中間件 | 中 | 高 | 中大型團隊、高并發(fā)、需高可靠場景 |
| Seata | 分布式事務框架 | 高 | 高 | 強一致性需求、復雜業(yè)務場景 |
四、總結與演進方向
本地消息表方案以無中間件依賴、實現簡單、可靠性高的特點,成為小型團隊解決分布式最終一致性問題的首選。核心是通過本地事務保證業(yè)務與消息的原子性,定時補償實現消息必達,冪等設計防止重復處理。
適用場景
- 對實時性要求不高的場景(如積分同步、物流狀態(tài)更新);
- 不想引入復雜中間件的輕量級微服務系統(tǒng);
- 訂單創(chuàng)建、庫存扣減、支付回調等核心業(yè)務場景。
演進方向
- 接入消息中間件:當業(yè)務規(guī)模擴大,可平滑遷移至 RabbitMQ/RocketMQ/Kafka 事務消息,提升實時性與吞吐量;
- 引入 Saga 模式:處理長鏈路業(yè)務,支持正向流程與反向補償;
- 云原生適配:結合 Serverless 架構,實現消息處理的彈性擴縮容。
到此這篇關于SpringBoot+本地消息表實現分布式最終一致性的文章就介紹到這了,更多相關SpringBoot 分布式最終一致性內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!
相關文章
intellij idea 2021.2 打包并上傳運行spring boot項目的詳細過程(spring boot 2
這篇文章主要介紹了intellij idea 2021.2 打包并上傳運行一個spring boot項目(spring boot 2.5.4),本文通過圖文并茂的形式給大家介紹的非常詳細,需要的朋友可以參考下2021-09-09
Java遠程執(zhí)行shell命令出現java: command not found問題及解決
這篇文章主要介紹了Java遠程執(zhí)行shell命令出現java: command not found問題及解決方案,具有很好的參考價值,希望對大家有所幫助。2023-07-07
SpringMVC中事務是否可以加在Controller層的問題
這篇文章主要介紹了SpringMVC中事務是否可以加在Controller層的問題,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教2022-02-02

