Apache RocketMQ 副本组 Quorum Write 与自适应降级(Adaptive Downgrade)实现解析

发布时间:2026/9/20 18:11:32

Apache RocketMQ 副本组 Quorum Write 与自适应降级(Adaptive Downgrade)实现解析 消息队列后端微服务流处理【免费下载链接】rocketmqApache RocketMQ is a cloud native messaging and streaming platform, making it simple to build event-driven applications.项目地址https://gitcode.com/gh_mirrors/ro/rocketmq点击查看免费下载本篇文章基于 docs/en/QuorumACK.md 与 docs/cn/QuorumACK.md 展开深入讲解 RocketMQ 5 在 Master-Slave 复制架构中引入的 Quorum Write法定人数写入与自适应降级机制。文章以该文档为骨架结合当前仓库中store模块的源码实现如CommitLog、DefaultHAService、MessageStoreConfig进行佐证与补充帮助读者理解如何在 broker 端精确指定一条消息写入成功后至少需要多少副本确认以及当副本掉线或落后过多时系统如何自动降级、保证可用性同时明确各参数的含义、默认值、生效条件与配置方法。背景同步复制与异步复制的取舍在 RocketMQ 的 Master-Slave主备架构中主备之间的数据复制主要有两种模式同步复制Synchronous ReplicationMaster 需要等待 Slave 成功复制消息并确认后才向 Producer 返回写入成功。同步复制可以保证 Master 失效后数据仍然能在 Slave 中找到适合可靠性要求较高的场景。异步复制Asynchronous ReplicationMaster 不需要等待 Slave 的响应即返回成功。异步复制虽然可能丢失消息但由于无需等待 Slave 确认效率高于同步复制适合对效率有一定要求的场景。在消息发送过程中客户端最终会收到如下几种结果状态状态含义PUT_OK一切顺利消息写入成功FLUSH_SLAVE_TIMEOUTSlave 同步超时SLAVE_NOT_AVAILABLESlave 不可用或 Slave 与 Master 的 CommitLog 差距超过一定值默认 256MB其中后两种状态并不会导致系统异常而无法写入下一条消息但它们都意味着当前副本组的同步状态并不健康。然而只有同步和异步两种模式在灵活性上存在明显不足在三副本甚至五副本且可靠性要求高的场景中异步复制无法满足要求而同步复制需要每一个副本确认后才返回副本数多时严重拖慢写入效率在同步复制模式下如果副本组中某一个 Slave 出现假死整个发送会一直失败直到人工介入处理。因此RocketMQ 5 提出了副本组的Quorum Write法定人数写入机制在同步复制模式下用户可以在 broker 端指定发送后至少需要写入多少副本数后才能返回同时提供**自适应降级Adaptive Downgrade**能力根据存活的副本数以及 CommitLog 差距自动完成降级。该特性在社区通过 RIP-34Support quorum write and adaptive degradation in master-slave architecture提出并落地。Quorum Write通过 totalReplicas 与 inSyncReplicas 灵活指定 ACK 副本数Quorum Write 通过增加两个 broker 端参数实现参数含义默认值totalReplicas副本组 broker 总数1inSyncReplicas正常情况下需保持同步的副本组数量1这两个参数定义在 MessageStoreConfig.java 中均标注为ImportantFieldImportantField private int totalReplicas 1; /** * Each message must be written successfully to at least in-sync replicas. * The master broker is considered one of the in-sync replicas, and its included in the count of total. * If a master broker is ASYNC_MASTER, inSyncReplicas will be ignored. * If enableControllerMode is true and ackAckInSyncStateSet is true, inSyncReplicas will be ignored. */ ImportantField private int inSyncReplicas 1;通过这两个参数可以在同步复制模式下灵活指定需要 ACK 的副本数例如两副本设置inSyncReplicas2则该条消息需要在 Master 和 Slave 中均复制完成后才返回给客户端三副本设置inSyncReplicas2则该条消息除了需要复制在 Master 上还需要复制到任意一个 Slave 上才返回给客户端四副本设置inSyncReplicas3则该条消息除了需要复制在 Master 上还需要复制到任意两个 Slave 上才返回给客户端。即inSyncReplicas指定的是包括 Master 自身在内需要 ACK 的副本数。通过灵活设置totalReplicas和inSyncReplicas可以满足各类场景对可靠性副本确认数与写入性能等待确认的副本数之间的平衡需求。值得注意的边界语义从源码注释可以确认以下几点细节Master 被计入 in-sync 副本数inSyncReplicas的计数包含 Master 自身异步 Master 忽略inSyncReplicas如果 Master 是ASYNC_MASTER异步刷盘/异步复制角色inSyncReplicas会被忽略控制器模式下的特殊行为如果enableControllerModetrue且allAckInSyncStateSettrueinSyncReplicas会被忽略此时要求消息写入SyncStateSet 中的所有副本详见后文“与 Controller 模式的协同”一节。自动降级根据存活副本数与 CommitLog 高度差动态调整Quorum Write 解决了“指定 ACK 副本数”的问题但同步复制下“某个 Slave 假死导致整个发送失败”的问题依然存在。为此RocketMQ 5 提供了自动降级能力。自动降级依据两个标准当前副本组的存活副本数Master CommitLog 与 Slave CommitLog 的高度差。注意自动降级只在slaveActingMaster模式开启后才生效。slaveActingMaster开关定义在 BrokerConfig.java 中默认值为falseprivate boolean enableSlaveActingMaster false;新增的三个参数自动降级引入以下三个参数定义于 MessageStoreConfig.java参数含义默认值生效条件minInSyncReplicas最小需保持同步的副本组数量1仅在enableAutoInSyncReplicastrue时生效enableAutoInSyncReplicas自动同步降级开关false需同时开启slaveActingMaster模式haMaxGapNotInSync判定 Slave 是否与 Master in-sync 的高度差阈值见下文说明全局生效对应源码/** * Will be worked in auto multiple replicas mode, to provide minimum in-sync replicas. * It is still valid in controller mode. */ ImportantField private int minInSyncReplicas 1; /** * Dynamically adjust in-sync replicas to provide higher availability, the real time in-sync replicas * will smaller than inSyncReplicas config. */ ImportantField private boolean enableAutoInSyncReplicas false;各参数的行为说明minInSyncReplicas自动降级时允许降到的最小副本数。例如设置minInSyncReplicas1最坏情况下消息只需写入 Master 即可成功enableAutoInSyncReplicas自动同步降级开关。开启后若当前副本组处于同步状态的 broker 数量包括 Master 自身不满足inSyncReplicas指定的数量则按照minInSyncReplicas进行同步haMaxGapNotInSyncSlave 是否与 Master 处于 in-sync 状态的判断阈值。若 Slave 的 CommitLog 落后 Master 长度超过该值则认为该 Slave 已处于非同步状态。关于haMaxGapNotInSync的默认值需要特别说明文档中记载的默认值为256K而当前仓库源码 MessageStoreConfig.java 中实际定义为private int haMaxGapNotInSync 1024 * 1024 * 256;即当前源码中的默认值是256MB1024×1024×256 字节与文档记载存在差异读者在实际部署时请以所用版本源码为准。该参数的调优方向如下当enableAutoInSyncReplicastrue时该值越小越容易触发 Master 的自动降级Slave 稍一落后就被判定为 out-of-sync当enableAutoInSyncReplicasfalse且totalReplicas inSyncReplicas时该值越小越容易导致大流量时发送请求失败此时可适当调大haMaxGapNotInSync。与 RocketMQ 4.x 的差异haSlaveFallBehindMax 被取消在 RocketMQ 4.x 中存在haSlaveFallbehindMax参数默认值为256MB用于表示 Slave 与 Master 的 CommitLog 高度差达到多少后判定 Slave 不可用。该参数在 RIP-34 中被取消由上述haMaxGapNotInSync等参数取代。源码级实现从 CommitLog 写入到 HA 服务判定1. 写入路径上的 needAckNums 计算在消息写入的核心类 CommitLog.java 中asyncPutMessage单条消息异步写入与批量写入路径都会在真正追加消息之前计算本次发送需要多少副本确认needAckNumsint needAckNums this.defaultMessageStore.getMessageStoreConfig().getInSyncReplicas(); boolean needHandleHA needHandleHA(msg); if (needHandleHA this.defaultMessageStore.getBrokerConfig().isEnableControllerMode()) { // 控制器模式不足 minInSyncReplicas 直接拒绝 if (this.defaultMessageStore.getHaService().inSyncReplicasNums(currOffset) this.defaultMessageStore.getMessageStoreConfig().getMinInSyncReplicas()) { return CompletableFuture.completedFuture(new PutMessageResult( PutMessageStatus.IN_SYNC_REPLICAS_NOT_ENOUGH, null)); } if (this.defaultMessageStore.getMessageStoreConfig().isAllAckInSyncStateSet()) { // -1 means all ack in SyncStateSet needAckNums MixAll.ALL_ACK_IN_SYNC_STATE_SET; } } else if (needHandleHA this.defaultMessageStore.getBrokerConfig().isEnableSlaveActingMaster()) { int inSyncReplicas Math.min(this.defaultMessageStore.getAliveReplicaNumInGroup(), this.defaultMessageStore.getHaService().inSyncReplicasNums(currOffset)); needAckNums calcNeedAckNums(inSyncReplicas); if (needAckNums inSyncReplicas) { // Tell the producer, dont have enough slaves to handle the send request return CompletableFuture.completedFuture(new PutMessageResult( PutMessageStatus.IN_SYNC_REPLICAS_NOT_ENOUGH, null)); } }对应文档中的代码骨架calcNeedAckNums实现了自适应降级的核心计算见 CommitLog.javaprivate int calcNeedAckNums(int inSyncReplicas) { int needAckNums this.defaultMessageStore.getMessageStoreConfig().getInSyncReplicas(); if (this.defaultMessageStore.getMessageStoreConfig().isEnableAutoInSyncReplicas()) { needAckNums Math.min(needAckNums, inSyncReplicas); needAckNums Math.max(needAckNums, this.defaultMessageStore.getMessageStoreConfig().getMinInSyncReplicas()); } return needAckNums; }这段逻辑的关键点在于inSyncReplicas实际可用的 in-sync 副本数取存活副本数与HA 服务统计的 in-sync Slave 数 1Master 自身两者中的较小值当enableAutoInSyncReplicastrue时期望的needAckNums被限制在[minInSyncReplicas, inSyncReplicas]区间内——即“能等多少就等多少但最低不低于 minInSyncReplicas”当needAckNums inSyncReplicas即实际可用副本数不足以满足要求时直接向 Producer 返回IN_SYNC_REPLICAS_NOT_ENOUGH拒绝写入。IN_SYNC_REPLICAS_NOT_ENOUGH状态定义在 PutMessageStatus.java。在 DLedgerCommitLog.javaDLedger 模式中同样会返回该状态。2. 存活副本数与 in-sync 副本数的来源存活副本数AliveReplicaNumInGroup定义于 DefaultMessageStore.java初始值为 1在构造时被设置为totalReplicasprivate volatile int aliveReplicasNum 1; ... this.aliveReplicasNum messageStoreConfig.getTotalReplicas();随后通过setAliveReplicaNumInGroup/getAliveReplicaNumInGroup维护见 DefaultMessageStore.java。该存活信息可以通过Nameserver 的反向通知以及GetBrokerMemberGroup 请求获取并同步到副本组内。in-sync Slave 数inSyncReplicasNums由 HA 服务统计。在 DefaultHAService.java 中Override public int inSyncReplicasNums(final long masterPutWhere) { int inSyncNums 1; // 初始为 1代表 Master 自身 for (HAConnection conn : this.connectionList) { if (this.isInSyncSlave(masterPutWhere, conn)) { inSyncNums; } } return inSyncNums; } protected boolean isInSyncSlave(final long masterPutWhere, HAConnection conn) { if (masterPutWhere - conn.getSlaveAckOffset() this.defaultMessageStore.getMessageStoreConfig() .getHaMaxGapNotInSync()) { return true; } return false; }即Master 当前写入位点masterPutWhere与某个 Slave 已确认位点slaveAckOffset之差小于haMaxGapNotInSync就判定该 Slave 处于 in-sync 状态。这正是文档中“自动降级标准”中“Master 与 Slave CommitLog 高度差”的落地实现——高度差直接由 HA 服务中的位点记录计算得出。3. 组提交等待GroupTransferService计算出的needAckNums会被传入handleHA进而封装为GroupCommitRequest提交给组提交服务等待见 CommitLog.javaif (needAckNums 0 needAckNums 1) { // 无需等待副本确认 } GroupCommitRequest request new GroupCommitRequest(nextOffset, this.defaultMessageStore.getMessageStoreConfig().getSlaveTimeout(), needAckNums);当needAckNums 1即只需 Master 自身确认时可以跳过等待逻辑直接成功——这正是降级后“只写 Master 即成功”的实现基础。等待与唤醒逻辑由 GroupTransferService.java继承ServiceThread的常驻线程完成它维护requestsWrite与requestsRead两个请求链表周期性地检查 Slave 的复制位点是否满足各请求所需的确认副本数。配置示例在 broker.conf 中启用 Quorum Write 与自动降级以下是一个三副本场景下开启 Quorum Write 与自适应降级的 broker 配置示例可参考 distribution/conf/broker.conf 的配置方式# 副本组 broker 总数Master 2 个 Slave totalReplicas3 # 正常情况下需保持同步的副本数量含 Master 自身 # 即消息需写入 Master 和任意 1 个 Slave 后才返回 inSyncReplicas2 # 最小需保持同步的副本数量自动降级下限 minInSyncReplicas1 # 自动同步降级开关 enableAutoInSyncReplicastrue # Slave 落后 Master 超过该字节数则判定为 out-of-sync # 注意当前仓库源码默认值为 1024*1024*256256MB请以实际版本为准 haMaxGapNotInSync268435456 # 关键前提自动降级仅在 slaveActingMaster 模式开启后生效 enableSlaveActingMastertrue典型行为验证两副本场景设置totalReplicas2、inSyncReplicas2、minInSyncReplicas1、enableAutoInSyncReplicastrue。正常情况下两个副本均处于同步复制消息需 Master 与 Slave 都确认当 Slave 下线或假死时系统进行自适应降级Producer 只需发送到 Master 即成功三副本场景设置totalReplicas3、inSyncReplicas2消息需 Master 与任意一个 in-sync 的 Slave 确认后返回兼顾可靠性与吞吐。与 ControllerDLedger 自动切换模式的协同从 CommitLog.java 可以看出自动降级逻辑在控制器模式enableControllerModetrue下走另一条分支此时用 HA 服务统计的 in-sync 副本数与minInSyncReplicas比较不足则直接返回IN_SYNC_REPLICAS_NOT_ENOUGH若开启allAckInSyncStateSet见 MessageStoreConfig.java则要求写入SyncStateSet 中所有副本needAckNums MixAll.ALL_ACK_IN_SYNC_STATE_SET即 -1。minInSyncReplicas在控制器模式下同样有效。这与slaveActingMaster分支共同构成了两条独立的“副本确认数动态计算”路径。兼容性升级到 RocketMQ 5 的行为变化为了保证向后兼容用户升级后必须设置正确的参数。默认情况下totalReplicas与inSyncReplicas均为 1这意味着假设用户原集群为两副本同步复制Master Slave同步等待 Slave 确认在不修改任何参数的情况下升级到 RocketMQ 5由于totalReplicas、inSyncReplicas默认都为 1将降级为异步复制只需 Master 确认如果希望保持与以前一致的行为两副本均确认后才返回则需要将totalReplicas和inSyncReplicas均设置为 2。同理原三副本同步复制的集群升级后若要保持“三副本全确认”的行为需要设置totalReplicas3、inSyncReplicas3若希望采用 Quorum Write 的“任一两个副本确认即可”策略则设置totalReplicas3、inSyncReplicas2。所有参数均在 broker 端broker.conf 或启动命令行-c指定配置文件配置。参考docs/en/QuorumACK.md本文英文原始文档docs/cn/QuorumACK.md本文中文原始文档MessageStoreConfig.javatotalReplicas、inSyncReplicas、minInSyncReplicas、enableAutoInSyncReplicas、haMaxGapNotInSync等参数定义CommitLog.javaneedAckNums计算与calcNeedAckNums实现DefaultHAService.javainSyncReplicasNums与isInSyncSlave实现GroupTransferService.java组提交等待服务DefaultMessageStore.java存活副本数维护BrokerConfig.javaenableSlaveActingMaster开关distribution/conf/broker.confbroker 配置文件示例赞分享消息队列后端微服务流处理【免费下载链接】rocketmqApache RocketMQ is a cloud native messaging and streaming platform, making it simple to build event-driven applications.项目地址https://gitcode.com/gh_mirrors/ro/rocketmq点击查看免费下载相关推荐Apache RocketMQ 5 Quorum Write 与自适应降级主备副本组同步策略深度指南Apache RocketMQ 5 Quorum Write 与自适应降级主备副本组同步策略深度指南 导读 RocketMQ 主备复制一直面临同步复制保可靠消息队列后端微服务流处理Apache RocketMQ DLedger配置详解副本数与选举策略优化Apache RocketMQ DLedger配置详解副本数与选举策略优化 引言分布式系统的容灾痛点与DLedger解决方案 在分布式消息中间件领域保障消消息队列流处理后端Apache RocketMQ Controller高可用部署多副本方案Apache RocketMQ Controller高可用部署多副本方案 1. 痛点与解决方案概述 在分布式系统中消息中间件的高可用性直接决定了业务连续性。消息队列后端微服务流处理上一篇LovyanGFX实战教程从基础绘制到高级动画的完整示例下一篇Windows 风扇控制实战用 FanControl 走完 5 个排障关卡创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
延伸阅读

更多相关文章

2026/9/20 18:06:32

DeepSeek 论文逻辑漏洞检测,Base URL 填 TaoToken 的 API 地址

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/20 18:06:32

Lucky 三步接入巴法云、点灯科技,语音助手控制家里的设备

Lucky 三步接入巴法云、点灯科技,语音助手控制家里的设备 【免费下载链接】lucky 软硬路由公网神器,ipv6/ipv4 端口转发,反向代理,DDNS,WOL,ipv4 stun内网穿透,cron,acme,rclone,ftp,webdav,filebrowser 项目地址: https://gitcode.com/GitHub_Trending/luc/lucky…

2026/9/20 21:01:49

doocs/source-code-hunter:ArrayList 底层原理源码级剖析与面试指南

文档教程知识库 【免费下载链接】source-code-hunter 😱 从源码层面,剖析挖掘互联网行业主流技术的底层实现原理,为广大开发者 “提升技术深度” 提供便利。目前开放 Spring 全家桶,Mybatis、Netty、Dubbo 框架,及 Red…

2026/9/20 21:01:49

TVBoxOSC 电视盒子播放管理指南:3 步让盒子开始看片

TVBoxOSC 电视盒子播放管理指南:3 步让盒子开始看片 【免费下载链接】TVBoxOSC TVBoxOSC - 一个基于第三方项目的代码库,用于电视盒子的控制和管理。 项目地址: https://gitcode.com/GitHub_Trending/tv/TVBoxOSC 如果想在电视上播放收藏的片源&a…

2026/9/20 20:56:48

如何把微信聊天记录导出成文档:WeChatMsg 完整上手指南

如何把微信聊天记录导出成文档:WeChatMsg 完整上手指南 【免费下载链接】WeChatMsg 提取微信聊天记录,将其导出成HTML、Word、CSV文档永久保存,对聊天记录进行分析生成年度聊天报告 项目地址: https://gitcode.com/GitHub_Trending/we/WeCh…

2026/9/20 0:04:49

GAMP 5 基于风险的计算机化系统验证:软件分类与审计追踪实践

简介:《A Risk-Based Approach to Compliant GxP Computerized Systems》即业内熟知的GAMP 5指南,面向制药企业质量与IT合规人员、验证工程师及计算机化系统管理者,用于解决GxP法规环境下系统合规性难以科学落地的问题。文档以风险管理为主线…

2026/9/20 0:04:49

安全托管MSSP实战:从静态防御到人机协同的攻防运营与应急响应

简介:这份PPT围绕互联网业务安全托管服务展开,面向企业安全负责人、IT运维人员及关注MSSP/MSS选型的读者,重点回应传统安全过度依赖人工、碎片化静态防御难以对抗产业化攻击等痛点。资源共1个pptx文件,包体约30.63MB,以…

2026/9/20 0:04:49

GAMP 5 基于风险的计算机化系统验证:软件分类与审计追踪实践

简介:《A Risk-Based Approach to Compliant GxP Computerized Systems》即业内熟知的GAMP 5指南,面向制药企业质量与IT合规人员、验证工程师及计算机化系统管理者,用于解决GxP法规环境下系统合规性难以科学落地的问题。文档以风险管理为主线…

2026/9/20 0:04:49

安全托管MSSP实战:从静态防御到人机协同的攻防运营与应急响应

简介:这份PPT围绕互联网业务安全托管服务展开,面向企业安全负责人、IT运维人员及关注MSSP/MSS选型的读者,重点回应传统安全过度依赖人工、碎片化静态防御难以对抗产业化攻击等痛点。资源共1个pptx文件,包体约30.63MB,以…

2026/9/20 4:54:47

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/20 5:01:23

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/20 5:09:33

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码