分布式系统中重补偿机制与最终一致性实现讲解
分布式系统中重补偿机制与最终一致性实现讲解一、核心概念什么是最终一致性在分布式系统中强一致性所有节点同时看到相同数据的代价极高。最终一致性是一种妥协系统允许短暂的数据不一致但保证在有限时间内所有数据达到一致状态。与强一致性的对比维度强一致性最终一致性数据可见性写入后立即可见写入后存在延迟窗口性能低需同步等待高异步处理可用性受限需所有节点在线高允许部分节点暂时不可用实现复杂度低高需补偿机制什么是补偿机制当某个操作因为依赖条件未满足、网络异常、服务不可用等原因失败时系统通过某种策略在之后重新执行该操作最终达到预期状态。补偿不是回滚而是正向重试直到成功。二、补偿策略分级一个健壮的系统通常采用多级补偿响应速度逐级递减覆盖范围逐级增大┌─────────────────────────────────────────────────────────────┐ │ 第一级即时触发毫秒~秒级 │ │ → 事件驱动条件满足时立即执行 │ ├─────────────────────────────────────────────────────────────┤ │ 第二级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接近上限时应告警数据清理已成功的历史记录定期归档避免表膨胀无效重试识别如果失败原因是业务层面不可恢复的如数据不存在且永远不会存在应标记为终态而非持续重试

相关新闻

【后端实战】超详细用户在线状态系统设计(心跳机制+Redis+分布式+多端登录+避坑指南)

【后端实战】超详细用户在线状态系统设计(心跳机制+Redis+分布式+多端登录+避坑指南)

文章简介:用户在线状态是IM聊天、社交软件、协同办公、客服系统的核心基础功能。90%的新手都会踩坑:闪退、断网、杀进程导致状态卡死在线。本文从需求分析、状态定义、技术方案演进、核心心跳原理、Redis存储设计、多端登录、分布式适配、推送策略、生产…

2026/7/28 23:19:53 阅读更多 →
贡献指南:如何参与 go-nebulas 开源项目并成为核心开发者

贡献指南:如何参与 go-nebulas 开源项目并成为核心开发者

贡献指南:如何参与 go-nebulas 开源项目并成为核心开发者 【免费下载链接】go-nebulas Official Go implementation of the Nebulas protocol. 项目地址: https://gitcode.com/gh_mirrors/go/go-nebulas go-nebulas 是Nebulas协议的官方Go实现,当…

2026/7/28 23:19:53 阅读更多 →
C++ vector内存陷阱:浅拷贝与迭代器失效的深度解析与实战解决方案

C++ vector内存陷阱:浅拷贝与迭代器失效的深度解析与实战解决方案

1. 项目概述:从一次内存泄漏事故说起那天下午,线上服务突然告警,内存使用率在几分钟内飙升到90%以上,服务响应变得极其缓慢。经过紧急排查和内存Dump分析,罪魁祸首锁定在一段处理批量数据的C代码上。代码里大量使用了s…

2026/7/28 23:19:53 阅读更多 →

最新新闻

密码学与逆向工程入门,攻克 CTF 难点题型

密码学与逆向工程入门,攻克 CTF 难点题型

破局 CTF:从密码学迷雾到逆向工程深处在 CTF(Capture The Flag)的赛场上,Web 和 Misc 方向往往因为上手快、反馈直接而成为新手的“舒适区”。然而,真正决定比赛上限的,往往是那些让人望而生畏的 Crypto&am…

2026/7/28 23:27:59 阅读更多 →
Kali Linux 工具链详解,渗透测试必备神器清单

Kali Linux 工具链详解,渗透测试必备神器清单

为什么 Kali Linux 是渗透测试的“瑞士军刀”在网络安全的学习路径中,工具的选择往往决定了实战的效率。对于初学者而言,面对琳琅满目的安全软件,最容易陷入“收藏党”的误区——下载了一堆工具却不知从何下手。Kali Linux 之所以成为行业标准…

2026/7/28 23:27:59 阅读更多 →
从 CTF 做题到挖洞赚钱,新手变现实操指南

从 CTF 做题到挖洞赚钱,新手变现实操指南

CTF:从解题到挖洞的低成本试错场很多刚接触网络安全的朋友,往往被“黑客”、“渗透”这些词汇背后的法律风险和高门槛劝退。大家心里都有个疑问:我想靠技术赚点零花钱,但还没本事去碰真实网站,怎么办?其实&…

2026/7/28 23:27:59 阅读更多 →
RulEuler规则引擎:基于Rete算法的高效业务规则处理

RulEuler规则引擎:基于Rete算法的高效业务规则处理

1. RulEuler规则引擎概述RulEuler是一款基于Rete算法实现的开源规则引擎,它通过高效的模式匹配机制,为业务规则与数据处理的解耦提供了专业解决方案。我在金融风控系统开发中首次接触这个项目时,发现其执行效率比传统if-else规则链提升了近20…

2026/7/28 23:27:59 阅读更多 →
Minerva核心组件解析:DAG调度器如何让多GPU协同工作如丝般顺滑

Minerva核心组件解析:DAG调度器如何让多GPU协同工作如丝般顺滑

Minerva核心组件解析:DAG调度器如何让多GPU协同工作如丝般顺滑 【免费下载链接】minerva Minerva: a fast and flexible tool for deep learning on multi-GPU. It provides ndarray programming interface, just like Numpy. Python bindings and C bindings are b…

2026/7/28 23:27:59 阅读更多 →
话咽回去时听《别说话 让我抱一下》

话咽回去时听《别说话 让我抱一下》

话到嘴边又咽回去,拥抱有时比解释清楚。《别说话 让我抱一下》把这个动作写成可跟随的听感:先靠近,再谈道理。 歌名像一句指令,听进去却更像请求。情绪救援队长&添火乐队把请求写得很克制,避免演变成煽情口号。完整…

2026/7/28 23:26:59 阅读更多 →

日新闻

告别臃肿!3步让你的暗影精灵笔记本重获新生

告别臃肿!3步让你的暗影精灵笔记本重获新生

告别臃肿!3步让你的暗影精灵笔记本重获新生 【免费下载链接】OmenSuperHub Control Omen laptop performance, fan speeds, and keyboard lighting, and unlock power limits. 项目地址: https://gitcode.com/gh_mirrors/om/OmenSuperHub 你是否也曾为官方Om…

2026/7/28 0:00:43 阅读更多 →
RAG必踩坑!财报法规检索不准?这款开源工具让答案浮出水面,准确率飙升98.7%!

RAG必踩坑!财报法规检索不准?这款开源工具让答案浮出水面,准确率飙升98.7%!

做 RAG 的人应该都踩过这个致命的坑:把几百页的财报、法规、技术手册扔给向量库,问一个具体问题,搜出来的全是沾边但没用的内容 —— 关键信息要么被硬切块拆碎了,要么藏在几十条结果的最下面。语义相似≠真正相关,这个…

2026/7/28 0:00:43 阅读更多 →
抖音视频文案提取工具全指南:免费2026版、手机App、在线工具一网打尽

抖音视频文案提取工具全指南:免费2026版、手机App、在线工具一网打尽

2026年做短视频运营,从抖音上扒文案早就不是偷偷抄笔记的事了。我刚开始做内容的时候,每天刷半小时抖音,手动把爆款视频的口播敲进备忘录,一条2分钟的视频得花十来分钟,碰到语速快的还要反复回听。后来试了一圈工具&am…

2026/7/28 0:00:43 阅读更多 →

周新闻

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 数据集6000张 完整源码已标注数据集训练好的模型环境配置教程程序运行说明文档,可以直接使用!系统支持图片、视频、摄像头等多种方式检测裂缝,功能强大实用。 1数据集6000张 8各类别

2026/7/28 12:04:22 阅读更多 →
深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

pubg数据集 精选原图1.42万数据 1.49万标签 无任何重复、算法增强或冗余图像! pubg绝地求生目标检测数据集 1分类:e_body,14905个标签,txt格式 共计14244张图,99%为640*640尺寸图像 适合yolo目标检测、AI训练关键词&am…

2026/7/28 8:29:16 阅读更多 →
Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex检测数据集数据集详情检测类别: allies enemy tag图片总量:7247张训练集:5139张验证集:1425张测试集:683张标注状态:全部已标注,即拿即用数据格式:支持YOLO格式及其他格式&#…

2026/7/28 5:03:42 阅读更多 →

月新闻