分布式系统中重补偿机制与最终一致性实现讲解
分布式系统中重补偿机制与最终一致性实现讲解一、核心概念什么是最终一致性在分布式系统中强一致性所有节点同时看到相同数据的代价极高。最终一致性是一种妥协系统允许短暂的数据不一致但保证在有限时间内所有数据达到一致状态。与强一致性的对比维度强一致性最终一致性数据可见性写入后立即可见写入后存在延迟窗口性能低需同步等待高异步处理可用性受限需所有节点在线高允许部分节点暂时不可用实现复杂度低高需补偿机制什么是补偿机制当某个操作因为依赖条件未满足、网络异常、服务不可用等原因失败时系统通过某种策略在之后重新执行该操作最终达到预期状态。补偿不是回滚而是正向重试直到成功。二、补偿策略分级一个健壮的系统通常采用多级补偿响应速度逐级递减覆盖范围逐级增大┌─────────────────────────────────────────────────────────────┐ │ 第一级即时触发毫秒~秒级 │ │ → 事件驱动条件满足时立即执行 │ ├─────────────────────────────────────────────────────────────┤ │ 第二级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)NOTNULLDEFAULTOCOMMENTO-待处理 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 任务接收与落库ServiceSlf4jpublicclassAsyncTaskService{ResourceprivateAsyncTaskLogRepositorytaskLogRepository;ResourceprivateTaskMqSendertaskMqSender;/** * 接收外部请求落库后异步处理. */Transactional(rollbackForException.class)publicvoidreceiveTask(StringtaskType,StringbizId,Objectpayload){// 幂等已成功则直接返回AsyncTaskLogexistingtaskLogRepository.findByTaskTypeAndBizId(taskType,bizId);if(existing!nullY.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);// 事务提交后投递 MQLonglogIdtaskLog.getId();TransactionSynchronizationManager.registerSynchronization(newTransactionSynchronization(){OverridepublicvoidafterCommit(){taskMqSender.send(taskType,logId);}});}}4.3 第一级即时触发事件驱动当依赖条件被满足时主动查找并触发待处理的任务ServiceSlf4jpublicclassOrderService{ResourceprivateAsyncTaskLogRepositorytaskLogRepository;ResourceprivateTaskMqSendertaskMqSender;/** * 订单完成后主动触发依赖该订单的待处理任务. */Transactional(rollbackForException.class)publicvoidcompleteOrder(LongorderId){// 核心业务逻辑doCompleteOrder(orderId);// 主动触发查找依赖该订单的待处理任务triggerPendingTasks(orderId);}privatevoidtriggerPendingTasks(LongorderId){StringbizIdORDER_orderId;AsyncTaskLogpendingTasktaskLogRepository.findByTaskTypeAndBizId(BARCODE_SCAN,bizId);// 不存在或已成功无需触发if(pendingTasknull||Y.equals(pendingTask.getStatus())){return;}LongtaskIdpendingTask.getId();// 事务提交后投递 MQ保证当前事务数据对消费者可见TransactionSynchronizationManager.registerSynchronization(newTransactionSynchronization(){OverridepublicvoidafterCommit(){taskMqSender.send(BARCODE_SCAN,taskId);}});}}4.4 第二级MQ 消费与失败标记ComponentSlf4jpublicclassTaskMqConsumer{ResourceprivateAsyncTaskLogRepositorytaskLogRepository;ResourceprivateDistributedLockProviderlockProvider;ResourceprivateTaskProcessortaskProcessor;RabbitListener(queues${mq.queue.async-task})publicvoidconsume(LonglogId){AsyncTaskLogtaskLogtaskLogRepository.findById(logId).orElse(null);if(taskLognull){return;}// 幂等已成功直接跳过if(Y.equals(taskLog.getStatus())){return;}// 分布式锁防止并发消费定时任务重投 主动触发可能同时到达StringlockKeytask:lock:taskLog.getTaskType():taskLog.getBizId();DistributedLocklocklockProvider.getLock(lockKey,60,TimeUnit.SECONDS);if(!lock.tryLock(30,TimeUnit.SECONDS)){log.warn(获取锁失败, logId{},logId);return;}try{// 再次检查状态获取锁期间可能已被其他消费者处理taskLogtaskLogRepository.findById(logId).orElse(null);if(taskLognull||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 第三级定时任务兜底ComponentSlf4jpublicclassTaskCompensationJob{ResourceprivateAsyncTaskLogRepositorytaskLogRepository;ResourceprivateTaskMqSendertaskMqSender;/** * 定时扫描未完成的任务重新投递MQ. * 建议执行间隔5~10 分钟 */Scheduled(cron0 */5 * * * ?)publicvoidcompensate(){DatestartTimeDateUtils.addHours(newDate(),-24);// 只扫最近24小时DateendTimenewDate();ListAsyncTaskLogpendingTaskstaskLogRepository.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{privatefinalListRunnablecommitActionsnewArrayList();privatefinalListRunnablerollbackActionsnewArrayList();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(statusSTATUS_ROLLED_BACK){for(Runnableaction:rollbackActions){try{action.run();}catch(Exceptione){log.warn(事务回滚后动作执行异常,e);}}}}}使用方式AfterTransactionActionCollectorcollectornewAfterTransactionActionCollector();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 去重查询法// 执行业务前查询是否已处理ListStringexistingBarcodesbarcodeMapper.listExistingBarcodes(orderId);ListStringtoInsertnewBarcodes.stream().filter(b-!existingBarcodes.contains(b)).collect(Collectors.toList());if(toInsert.isEmpty()){return;}5.4 分布式锁串行化// 同一业务 ID 同一时刻只有一个消费者在处理StringlockKeyprocess:bizId;if(lock.tryLock()){try{// 获取锁后再次检查状态double-checkif(!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 丢了 → 定时任务兜底扫描statusO的记录重新投递。这就是为什么需要多级补偿。七、适用场景场景补偿策略选择两个接口时序不确定A 依赖 B 的结果后到的一方主动触发 定时任务兜底第三方回调可能丢失主动轮询 超时重试跨服务数据同步MQ 异步 状态表 定时对账批量任务部分失败逐条记录状态 失败的单独重试支付回调与订单状态同步回调处理 主动查询补偿 对账八、注意事项重试次数上限避免无限重试消耗资源超限后转人工或告警退避策略定时任务不要过于频繁指数退避1min → 5min → 30min更合理监控告警statusP且retry_count接近上限时应告警数据清理已成功的历史记录定期归档避免表膨胀无效重试识别如果失败原因是业务层面不可恢复的如数据不存在且永远不会存在应标记为终态而非持续重试