分布式系统中重补偿机制与最终一致性实现讲解

📅 2026/7/28 23:21:05 👁️ 阅读次数 📝 编程学习
分布式系统中重补偿机制与最终一致性实现讲解

分布式系统中重补偿机制与最终一致性实现讲解

一、核心概念

什么是最终一致性

在分布式系统中,强一致性(所有节点同时看到相同数据)的代价极高。最终一致性是一种妥协:系统允许短暂的数据不一致,但保证在有限时间内所有数据达到一致状态

与强一致性的对比:

维度强一致性最终一致性
数据可见性写入后立即可见写入后存在延迟窗口
性能低(需同步等待)高(异步处理)
可用性受限(需所有节点在线)高(允许部分节点暂时不可用)
实现复杂度高(需补偿机制)

什么是补偿机制

当某个操作因为依赖条件未满足、网络异常、服务不可用等原因失败时,系统通过某种策略在之后重新执行该操作,最终达到预期状态。补偿不是回滚,而是正向重试直到成功


二、补偿策略分级

一个健壮的系统通常采用多级补偿,响应速度逐级递减,覆盖范围逐级增大:

┌─────────────────────────────────────────────────────────────┐ │ 第一级:即时触发(毫秒~秒级) │ │ → 事件驱动,条件满足时立即执行 │ ├─────────────────────────────────────────────────────────────┤ │ 第二级:MQ 重试(秒~分钟级) │ │ → 消费失败后 MQ 自带的重投机制 │ ├─────────────────────────────────────────────────────────────┤ │ 第三级:定时任务兜底(分钟~小时级) │ │ → 扫描未完成记录批量重新投递 │ ├─────────────────────────────────────────────────────────────┤ │ 第四级:人工介入 │ │ → 告警通知 + 运维工具手动触发 │ └─────────────────────────────────────────────────────────────┘

设计原则

  • 快速路径优先:能即时补偿就不等定时任务
  • 兜底必须存在:即时触发可能因自身异常失效,定时任务作为最终保障
  • 多级不冲突:幂等设计保证多个级别同时触发同一条数据不会产生副作用

注:

博客:

https://blog.csdn.net/badao_liumang_qizhi

三、状态机驱动模型

补偿机制通常配合状态机使用,用一张状态表记录每个操作的处理进度:

┌───────┐ 处理中 ┌───────┐ 成功 ┌───────┐ │PENDING├──────────────────►│RUNNING├─────────────►│SUCCESS│ └───┬───┘ └───┬───┘ └───────┘ │ │ │ 超时/失败 │ 失败 │◄──────────────────────────┘ │ ▼ ┌───────┐ 达到最大重试次数 ┌───────┐ │FAILED ├──────────────────────►│ DEAD │ └───────┘ └───────┘ public enum TaskStatus { PENDING("O"), // 待处理 FAILED("P"), // 处理失败,等待重试 SUCCESS("Y"); // 处理成功 private final String code; }

四、通用示例代码

4.1 状态表设计

CREATETABLEasync_task_log(idBIGINTAUTO_INCREMENTPRIMARYKEY,task_typeVARCHAR(32)NOTNULLCOMMENT'任务类型',biz_idVARCHAR(64)NOTNULLCOMMENT'业务标识',payloadTEXTNOTNULLCOMMENT'任务数据(JSON)',statusCHAR(1)NOTNULLDEFAULT'O'COMMENT'O-待处理 P-失败 Y-成功',retry_countINTNOTNULLDEFAULT0COMMENT'已重试次数',max_retryINTNOTNULLDEFAULT10COMMENT'最大重试次数',error_msgVARCHAR(512)COMMENT'最近一次失败原因',create_timeDATETIMENOTNULLDEFAULTCURRENT_TIMESTAMP,update_timeDATETIMENOTNULLDEFAULTCURRENT_TIMESTAMPONUPDATECURRENT_TIMESTAMP,INDEXidx_status_type(status,task_type),UNIQUEINDEXuk_biz(task_type,biz_id))COMMENT'异步任务日志表';

4.2 任务接收与落库

@Service@Slf4jpublicclassAsyncTaskService{@ResourceprivateAsyncTaskLogRepositorytaskLogRepository;@ResourceprivateTaskMqSendertaskMqSender;/** * 接收外部请求,落库后异步处理. */@Transactional(rollbackFor=Exception.class)publicvoidreceiveTask(StringtaskType,StringbizId,Objectpayload){// 幂等:已成功则直接返回AsyncTaskLogexisting=taskLogRepository.findByTaskTypeAndBizId(taskType,bizId);if(existing!=null&&"Y".equals(existing.getStatus())){return;}// 落库AsyncTaskLogtaskLog=(existing!=null)?existing:newAsyncTaskLog();taskLog.setTaskType(taskType);taskLog.setBizId(bizId);taskLog.setPayload(JsonUtil.toJson(payload));taskLog.setStatus("O");taskLogRepository.saveAndFlush(taskLog);// 事务提交后投递 MQLonglogId=taskLog.getId();TransactionSynchronizationManager.registerSynchronization(newTransactionSynchronization(){@OverridepublicvoidafterCommit(){taskMqSender.send(taskType,logId);}});}}

4.3 第一级:即时触发(事件驱动)

当依赖条件被满足时,主动查找并触发待处理的任务:

@Service@Slf4jpublicclassOrderService{@ResourceprivateAsyncTaskLogRepositorytaskLogRepository;@ResourceprivateTaskMqSendertaskMqSender;/** * 订单完成后,主动触发依赖该订单的待处理任务. */@Transactional(rollbackFor=Exception.class)publicvoidcompleteOrder(LongorderId){// 核心业务逻辑doCompleteOrder(orderId);// 主动触发:查找依赖该订单的待处理任务triggerPendingTasks(orderId);}privatevoidtriggerPendingTasks(LongorderId){StringbizId="ORDER_"+orderId;AsyncTaskLogpendingTask=taskLogRepository.findByTaskTypeAndBizId("BARCODE_SCAN",bizId);// 不存在或已成功,无需触发if(pendingTask==null||"Y".equals(pendingTask.getStatus())){return;}LongtaskId=pendingTask.getId();// 事务提交后投递 MQ(保证当前事务数据对消费者可见)TransactionSynchronizationManager.registerSynchronization(newTransactionSynchronization(){@OverridepublicvoidafterCommit(){taskMqSender.send("BARCODE_SCAN",taskId);}});}}

4.4 第二级:MQ 消费与失败标记

@Component@Slf4jpublicclassTaskMqConsumer{@ResourceprivateAsyncTaskLogRepositorytaskLogRepository;@ResourceprivateDistributedLockProviderlockProvider;@ResourceprivateTaskProcessortaskProcessor;@RabbitListener(queues="${mq.queue.async-task}")publicvoidconsume(LonglogId){AsyncTaskLogtaskLog=taskLogRepository.findById(logId).orElse(null);if(taskLog==null){return;}// 幂等:已成功直接跳过if("Y".equals(taskLog.getStatus())){return;}// 分布式锁:防止并发消费(定时任务重投 + 主动触发可能同时到达)StringlockKey="task:lock:"+taskLog.getTaskType()+":"+taskLog.getBizId();DistributedLocklock=lockProvider.getLock(lockKey,60,TimeUnit.SECONDS);if(!lock.tryLock(30,TimeUnit.SECONDS)){log.warn("获取锁失败, logId={}",logId);return;}try{// 再次检查状态(获取锁期间可能已被其他消费者处理)taskLog=taskLogRepository.findById(logId).orElse(null);if(taskLog==null||"Y".equals(taskLog.getStatus())){return;}// 执行业务逻辑taskProcessor.process(taskLog);// 标记成功taskLog.setStatus("Y");taskLog.setErrorMsg(null);taskLogRepository.saveAndFlush(taskLog);}catch(Exceptione){log.warn("任务消费失败, logId={}",logId,e);// 标记失败taskLog.setStatus("P");taskLog.setRetryCount(taskLog.getRetryCount()+1);taskLog.setErrorMsg(e.getMessage());taskLogRepository.saveAndFlush(taskLog);}finally{lock.unlock();}}}

4.5 第三级:定时任务兜底

@Component@Slf4jpublicclassTaskCompensationJob{@ResourceprivateAsyncTaskLogRepositorytaskLogRepository;@ResourceprivateTaskMqSendertaskMqSender;/** * 定时扫描未完成的任务,重新投递MQ. * 建议执行间隔:5~10 分钟 */@Scheduled(cron="0 */5 * * * ?")publicvoidcompensate(){DatestartTime=DateUtils.addHours(newDate(),-24);// 只扫最近24小时DateendTime=newDate();List<AsyncTaskLog>pendingTasks=taskLogRepository.findByStatusInAndCreateTimeBetween(Arrays.asList("O","P"),startTime,endTime);for(AsyncTaskLogtask:pendingTasks){// 超过最大重试次数,跳过(转人工处理)if(task.getRetryCount()>=task.getMaxRetry()){log.warn("任务超过最大重试次数, id={}, bizId={}",task.getId(),task.getBizId());continue;}taskMqSender.send(task.getTaskType(),task.getId());}log.info("补偿任务扫描完成, 待处理数量={}",pendingTasks.size());}}

4.6 事务后置动作工具类

/** * 事务同步回调收集器. * 收集需要在事务提交/回滚后执行的动作。 */publicclassAfterTransactionActionCollectorimplementsTransactionSynchronization{privatefinalList<Runnable>commitActions=newArrayList<>();privatefinalList<Runnable>rollbackActions=newArrayList<>();publicvoidaddCommitAction(Runnableaction){commitActions.add(action);}publicvoidaddRollbackAction(Runnableaction){rollbackActions.add(action);}@OverridepublicvoidafterCommit(){for(Runnableaction:commitActions){try{action.run();}catch(Exceptione){// 提交后动作失败不影响已提交的事务log.warn("事务提交后动作执行异常",e);}}}@OverridepublicvoidafterCompletion(intstatus){if(status==STATUS_ROLLED_BACK){for(Runnableaction:rollbackActions){try{action.run();}catch(Exceptione){log.warn("事务回滚后动作执行异常",e);}}}}}

使用方式:

AfterTransactionActionCollectorcollector=newAfterTransactionActionCollector();collector.addCommitAction(()->mqSender.send(taskId));collector.addCommitAction(()->lock.unlock());collector.addRollbackAction(()->lock.unlock());TransactionSynchronizationManager.registerSynchronization(collector);

五、幂等性保障

补偿机制的前提是重复执行不产生副作用。常见幂等策略:

5.1 状态判断法

// 消费前判断状态if("Y".equals(taskLog.getStatus())){return;// 已成功,跳过}

5.2 唯一约束法

-- 数据库层面保证不会插入重复数据UNIQUEINDEXuk_biz(task_type,biz_id)

5.3 去重查询法

// 执行业务前查询是否已处理List<String>existingBarcodes=barcodeMapper.listExistingBarcodes(orderId);List<String>toInsert=newBarcodes.stream().filter(b->!existingBarcodes.contains(b)).collect(Collectors.toList());if(toInsert.isEmpty()){return;}

5.4 分布式锁串行化

// 同一业务 ID 同一时刻只有一个消费者在处理StringlockKey="process:"+bizId;if(lock.tryLock()){try{// 获取锁后再次检查状态(double-check)if(!"Y".equals(reload().getStatus())){doProcess();}}finally{lock.unlock();}}

六、关键时序问题

为什么 MQ 必须在事务提交后发送?

错误做法:事务内发送 MQ ┌───事务开始──────────────────────事务提交───┐ │ 写入数据A 发送MQ 写入数据B │ └─────────────────────────────────────────────┘ ↓ 消费者收到消息 查询数据A → 可能查到(取决于隔离级别) 查询数据B → 未提交,查不到 ❌ 正确做法:事务提交后发送 MQ ┌───事务开始──────────────────事务提交───┐ │ 写入数据A 写入数据B │ └────────────────────────────────────────┘ ↓ afterCommit() 发送MQ ↓ 消费者收到消息 查询数据A → ✅ 查询数据B → ✅

如果afterCommit中 MQ 发送失败怎么办?

数据已经提交,MQ 丢了 → 定时任务兜底扫描status=O的记录重新投递。这就是为什么需要多级补偿。


七、适用场景

场景补偿策略选择
两个接口时序不确定(A 依赖 B 的结果)后到的一方主动触发 + 定时任务兜底
第三方回调可能丢失主动轮询 + 超时重试
跨服务数据同步MQ 异步 + 状态表 + 定时对账
批量任务部分失败逐条记录状态 + 失败的单独重试
支付回调与订单状态同步回调处理 + 主动查询补偿 + 对账

八、注意事项

  1. 重试次数上限:避免无限重试消耗资源,超限后转人工或告警
  2. 退避策略:定时任务不要过于频繁,指数退避(1min → 5min → 30min)更合理
  3. 监控告警status=Pretry_count接近上限时应告警
  4. 数据清理:已成功的历史记录定期归档,避免表膨胀
  5. 无效重试识别:如果失败原因是业务层面不可恢复的(如数据不存在且永远不会存在),应标记为终态而非持续重试