供应链的分布式事务处理:跨仓库调拨的最终一致性方案

📅 2026/7/25 3:33:42
供应链的分布式事务处理:跨仓库调拨的最终一致性方案
供应链的分布式事务处理跨仓库调拨的最终一致性方案一、货从A仓发出了B仓库存没加分布式调拨的数据黑洞某电商平台的供应链系统在双11期间遇到一个诡异问题华南仓向华东仓调拨了5000件商品华南仓已经扣减库存并出库物流显示运输中但华东仓的库存没有增加。原因是跨仓库调拨的分布式事务中华南仓的扣减事务提交成功华东仓的入库事务因为网络超时回滚了——两个操作不在同一个本地事务中。这就是供应链中最经典的分布式问题A仓扣减 B仓增加 两个独立的数据库事务天然的分布式事务场景。二、TCC SAGA 事件溯源的三重保障三、SAGA模式的跨仓库调拨实现Service public class CrossWarehouseTransferService { private final JdbcTemplate warehouseA; private final JdbcTemplate warehouseB; private final KafkaTemplateString, TransferEvent kafka; Transactional(warehouseATransactionManager) public void transferOut(String transferId, String skuId, int quantity, String targetWarehouse) { try { // Step 1: A仓扣减并记录调拨单 int rows warehouseA.update( UPDATE inventory SET stock stock - ?, locked locked ? WHERE sku_id ? AND stock ?, quantity, quantity, skuId, quantity ); if (rows 0) { throw new InsufficientStockException( skuId 库存不足 ); } warehouseA.update( INSERT INTO transfer_out_records (transfer_id, sku_id, quantity, target_warehouse, status) VALUES (?, ?, ?, ?, OUTBOUND), transferId, skuId, quantity, targetWarehouse ); // Step 2: 发布调拨事件 kafka.send(transfer_events, new TransferEvent( transferId, skuId, quantity, WAREHOUSE_A, targetWarehouse, OUTBOUND_CONFIRMED )); } catch (Exception e) { kafka.send(transfer_events, new TransferEvent( transferId, skuId, quantity, WAREHOUSE_A, targetWarehouse, OUTBOUND_FAILED )); throw new TransferException(调拨出库失败, e); } } KafkaListener(topics transfer_events) public void handleTransferEvent(TransferEvent event) { if (!OUTBOUND_CONFIRMED.equals(event.getStatus())) { return; } try { // Step 3: B仓入库 warehouseB.update( INSERT INTO transfer_in_records (transfer_id, sku_id, quantity, source_warehouse, status) VALUES (?, ?, ?, ?, INBOUND) ON DUPLICATE KEY UPDATE retry_count retry_count 1, event.getTransferId(), event.getSkuId(), event.getQuantity(), event.getSourceWarehouse() ); warehouseB.update( UPDATE inventory SET stock stock ? WHERE sku_id ?, event.getQuantity(), event.getSkuId() ); // Step 4: 确认入库成功 kafka.send(transfer_events, new TransferEvent( event.getTransferId(), event.getSkuId(), event.getQuantity(), event.getSourceWarehouse(), event.getTargetWarehouse(), INBOUND_CONFIRMED )); } catch (Exception e) { // B仓入库失败 → 触发补偿 kafka.send(transfer_events, TransferEvent.compensate(event, INBOUND_FAILED_NEED_COMPENSATION) ); } } KafkaListener(topics transfer_events) public void handleCompensation(TransferEvent event) { if (!INBOUND_FAILED_NEED_COMPENSATION.equals(event.getStatus())) { return; } // 补偿策略将货物发往异常处理仓人工介入 try { warehouseA.update( INSERT INTO exception_warehouse_records (transfer_id, sku_id, quantity, reason) VALUES (?, ?, ?, 目标仓库入库失败), event.getTransferId(), event.getSkuId(), event.getQuantity() ); // 释放A仓locked库存 warehouseA.update( UPDATE inventory SET locked locked - ? WHERE sku_id ?, event.getQuantity(), event.getSkuId() ); } catch (Exception ex) { // 补偿也失败了 → 进入最终人工处理队列 finalFallbackQueue.add(event); } } }四、供应链分布式事务的三个现实妥协妥协一最终一致性的时间窗口。A仓扣减到B仓入库之间可能有2-5分钟的库存不对称期。业务上必须接受这个事实——如果运营在这个窗口内查询全公司总库存应该标注含在途库存。妥协二幂等性保障的复杂度。B仓入库的消息可能因为Kafka重试被消费两次。ON DUPLICATE KEY UPDATE retry_count retry_count 1这种幂等设计依赖transfer_id的唯一性。如果transfer_id是雪花算法生成的必须保证在分布式环境下的全局唯一。妥协三长事务的超时处理。调拨可能因为物流中断而卡在运输中状态数天。SAGA的补偿策略需要支持超时自动补偿——超过72小时仍未确认入库系统自动发起退货流程。五、总结跨仓库调拨的分布式事务处理TCC提供预留-确认-取消的原子性保证SAGA提供正向执行反向补偿的最终一致性事件溯源提供消息丢了也能重建的容灾兜底。三者叠加构成供应链分布式事务的完整保障体系。在供应链领域100%的强一致性是不现实的目标。务实的设计是把不一致的概率降到可接受范围0.01%并为剩余的不一致设计自动补偿人工处理的兜底机制。本文属于「行业场景与项目复盘」系列深入探讨供应链跨仓库调拨的分布式事务与最终一致性方案。