供应链的分布式事务处理:跨仓库调拨的最终一致性方案
供应链的分布式事务处理:跨仓库调拨的最终一致性方案
一、"货从A仓发出了,B仓库存没加":分布式调拨的数据黑洞
某电商平台的供应链系统在"双11"期间遇到一个诡异问题:华南仓向华东仓调拨了5000件商品,华南仓已经扣减库存并出库,物流显示"运输中",但华东仓的库存没有增加。原因是跨仓库调拨的分布式事务中,华南仓的扣减事务提交成功,华东仓的入库事务因为网络超时回滚了——两个操作不在同一个本地事务中。
这就是供应链中最经典的分布式问题:A仓扣减 + B仓增加 = 两个独立的数据库事务,天然的分布式事务场景。
二、TCC + SAGA + 事件溯源的三重保障
三、SAGA模式的跨仓库调拨实现
@Service public class CrossWarehouseTransferService { private final JdbcTemplate warehouseA; private final JdbcTemplate warehouseB; private final KafkaTemplate<String, 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%),并为剩余的不一致设计自动补偿+人工处理的兜底机制。
本文属于「行业场景与项目复盘」系列,深入探讨供应链跨仓库调拨的分布式事务与最终一致性方案。
