news 2026/9/20 18:11:31

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache RocketMQ 副本组 Quorum Write 与自适应降级(Adaptive Downgrade)实现解析
  • 消息队列
  • 后端
  • 微服务
  • 流处理

【免费下载链接】rocketmq

Apache 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模块的源码实现(如CommitLogDefaultHAServiceMessageStoreConfig)进行佐证与补充,帮助读者理解:如何在 broker 端精确指定一条消息写入成功后至少需要多少副本确认,以及当副本掉线或落后过多时系统如何自动降级、保证可用性,同时明确各参数的含义、默认值、生效条件与配置方法。

背景:同步复制与异步复制的取舍

在 RocketMQ 的 Master-Slave(主备)架构中,主备之间的数据复制主要有两种模式:

  • 同步复制(Synchronous Replication):Master 需要等待 Slave 成功复制消息并确认后,才向 Producer 返回写入成功。同步复制可以保证 Master 失效后,数据仍然能在 Slave 中找到,适合可靠性要求较高的场景。
  • 异步复制(Asynchronous Replication):Master 不需要等待 Slave 的响应即返回成功。异步复制虽然可能丢失消息,但由于无需等待 Slave 确认,效率高于同步复制,适合对效率有一定要求的场景。

在消息发送过程中,客户端最终会收到如下几种结果状态:

状态含义
PUT_OK一切顺利,消息写入成功
FLUSH_SLAVE_TIMEOUTSlave 同步超时
SLAVE_NOT_AVAILABLESlave 不可用,或 Slave 与 Master 的 CommitLog 差距超过一定值(默认 256MB)

其中后两种状态并不会导致系统异常而无法写入下一条消息,但它们都意味着当前副本组的同步状态并不健康。

然而,只有同步和异步两种模式在灵活性上存在明显不足:

  1. 在三副本甚至五副本且可靠性要求高的场景中,异步复制无法满足要求;
  2. 而同步复制需要每一个副本确认后才返回,副本数多时严重拖慢写入效率;
  3. 在同步复制模式下,如果副本组中某一个 Slave 出现假死,整个发送会一直失败,直到人工介入处理。

因此,RocketMQ 5 提出了副本组的Quorum Write(法定人数写入)机制:在同步复制模式下,用户可以在 broker 端指定发送后至少需要写入多少副本数后才能返回;同时提供**自适应降级(Adaptive Downgrade)**能力,根据存活的副本数以及 CommitLog 差距自动完成降级。

该特性在社区通过 RIP-34(Support quorum write and adaptive degradation in master-slave architecture)提出并落地。

Quorum Write:通过 totalReplicas 与 inSyncReplicas 灵活指定 ACK 副本数

Quorum Write 通过增加两个 broker 端参数实现:

参数含义默认值
totalReplicas副本组 broker 总数1
inSyncReplicas正常情况下需保持同步的副本组数量1

这两个参数定义在 MessageStoreConfig.java 中,均标注为@ImportantField

@ImportantField 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 it's 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 的副本数,例如:

  • 两副本:设置inSyncReplicas=2,则该条消息需要在 Master 和 Slave 中均复制完成后才返回给客户端;
  • 三副本:设置inSyncReplicas=2,则该条消息除了需要复制在 Master 上,还需要复制到任意一个 Slave 上才返回给客户端;
  • 四副本:设置inSyncReplicas=3,则该条消息除了需要复制在 Master 上,还需要复制到任意两个 Slave 上才返回给客户端。

即:inSyncReplicas指定的是包括 Master 自身在内需要 ACK 的副本数。通过灵活设置totalReplicasinSyncReplicas,可以满足各类场景对可靠性(副本确认数)与写入性能(等待确认的副本数)之间的平衡需求。

值得注意的边界语义

从源码注释可以确认以下几点细节:

  1. Master 被计入 in-sync 副本数inSyncReplicas的计数包含 Master 自身;
  2. 异步 Master 忽略inSyncReplicas:如果 Master 是ASYNC_MASTER(异步刷盘/异步复制角色),inSyncReplicas会被忽略;
  3. 控制器模式下的特殊行为:如果enableControllerMode=trueallAckInSyncStateSet=trueinSyncReplicas会被忽略,此时要求消息写入SyncStateSet 中的所有副本(详见后文“与 Controller 模式的协同”一节)。

自动降级:根据存活副本数与 CommitLog 高度差动态调整

Quorum Write 解决了“指定 ACK 副本数”的问题,但同步复制下“某个 Slave 假死导致整个发送失败”的问题依然存在。为此,RocketMQ 5 提供了自动降级能力。

自动降级依据两个标准:

  • 当前副本组的存活副本数
  • Master CommitLog 与 Slave CommitLog 的高度差

注意:自动降级只在slaveActingMaster模式开启后才生效。

slaveActingMaster开关定义在 BrokerConfig.java 中,默认值为false

private boolean enableSlaveActingMaster = false;

新增的三个参数

自动降级引入以下三个参数(定义于 MessageStoreConfig.java):

参数含义默认值生效条件
minInSyncReplicas最小需保持同步的副本组数量1仅在enableAutoInSyncReplicas=true时生效
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:自动降级时允许降到的最小副本数。例如设置minInSyncReplicas=1,最坏情况下消息只需写入 Master 即可成功;
  • enableAutoInSyncReplicas:自动同步降级开关。开启后,若当前副本组处于同步状态的 broker 数量(包括 Master 自身)不满足inSyncReplicas指定的数量,则按照minInSyncReplicas进行同步;
  • haMaxGapNotInSync:Slave 是否与 Master 处于 in-sync 状态的判断阈值。若 Slave 的 CommitLog 落后 Master 长度超过该值,则认为该 Slave 已处于非同步状态。

关于haMaxGapNotInSync的默认值需要特别说明:文档中记载的默认值为256K,而当前仓库源码 MessageStoreConfig.java 中实际定义为:

private int haMaxGapNotInSync = 1024 * 1024 * 256;

即当前源码中的默认值是256MB(1024×1024×256 字节),与文档记载存在差异,读者在实际部署时请以所用版本源码为准。该参数的调优方向如下:

  • enableAutoInSyncReplicas=true时,该值越小越容易触发 Master 的自动降级(Slave 稍一落后就被判定为 out-of-sync);
  • enableAutoInSyncReplicas=falsetotalReplicas == 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(单条消息异步写入)与批量写入路径都会在真正追加消息之前计算本次发送需要多少副本确认(needAckNums):

int 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, don't have enough slaves to handle the send request return CompletableFuture.completedFuture(new PutMessageResult( PutMessageStatus.IN_SYNC_REPLICAS_NOT_ENOUGH, null)); } }

对应文档中的代码骨架,calcNeedAckNums实现了自适应降级的核心计算(见 CommitLog.java):

private 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; }

这段逻辑的关键点在于:

  1. inSyncReplicas(实际可用的 in-sync 副本数)取存活副本数HA 服务统计的 in-sync Slave 数 + 1(Master 自身)两者中的较小值;
  2. enableAutoInSyncReplicas=true时,期望的needAckNums被限制在[minInSyncReplicas, inSyncReplicas]区间内——即“能等多少就等多少,但最低不低于 minInSyncReplicas”;
  3. needAckNums > inSyncReplicas(即实际可用副本数不足以满足要求)时,直接向 Producer 返回IN_SYNC_REPLICAS_NOT_ENOUGH,拒绝写入。

IN_SYNC_REPLICAS_NOT_ENOUGH状态定义在 PutMessageStatus.java。在 DLedgerCommitLog.java(DLedger 模式)中同样会返回该状态。

2. 存活副本数与 in-sync 副本数的来源

  • 存活副本数(AliveReplicaNumInGroup):定义于 DefaultMessageStore.java,初始值为 1,在构造时被设置为totalReplicas
private 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.java):

if (needAckNums >= 0 && needAckNums <= 1) { // 无需等待副本确认 } GroupCommitRequest request = new GroupCommitRequest(nextOffset, this.defaultMessageStore.getMessageStoreConfig().getSlaveTimeout(), needAckNums);

needAckNums <= 1(即只需 Master 自身确认)时,可以跳过等待逻辑直接成功——这正是降级后“只写 Master 即成功”的实现基础。等待与唤醒逻辑由 GroupTransferService.java(继承ServiceThread的常驻线程)完成,它维护requestsWriterequestsRead两个请求链表,周期性地检查 Slave 的复制位点是否满足各请求所需的确认副本数。

配置示例:在 broker.conf 中启用 Quorum Write 与自动降级

以下是一个三副本场景下,开启 Quorum Write 与自适应降级的 broker 配置示例(可参考 distribution/conf/broker.conf 的配置方式):

# 副本组 broker 总数:Master + 2 个 Slave totalReplicas=3 # 正常情况下需保持同步的副本数量(含 Master 自身) # 即:消息需写入 Master 和任意 1 个 Slave 后才返回 inSyncReplicas=2 # 最小需保持同步的副本数量(自动降级下限) minInSyncReplicas=1 # 自动同步降级开关 enableAutoInSyncReplicas=true # Slave 落后 Master 超过该字节数则判定为 out-of-sync # 注意:当前仓库源码默认值为 1024*1024*256(256MB),请以实际版本为准 haMaxGapNotInSync=268435456 # 关键前提:自动降级仅在 slaveActingMaster 模式开启后生效 enableSlaveActingMaster=true

典型行为验证:

  • 两副本场景:设置totalReplicas=2inSyncReplicas=2minInSyncReplicas=1enableAutoInSyncReplicas=true。正常情况下两个副本均处于同步复制,消息需 Master 与 Slave 都确认;当 Slave 下线或假死时,系统进行自适应降级,Producer 只需发送到 Master 即成功;
  • 三副本场景:设置totalReplicas=3inSyncReplicas=2,消息需 Master 与任意一个 in-sync 的 Slave 确认后返回,兼顾可靠性与吞吐。

与 Controller(DLedger 自动切换)模式的协同

从 CommitLog.java 可以看出,自动降级逻辑在控制器模式enableControllerMode=true)下走另一条分支:此时用 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 的行为变化

为了保证向后兼容,用户升级后必须设置正确的参数。默认情况下,totalReplicasinSyncReplicas均为 1,这意味着:

  • 假设用户原集群为两副本同步复制(Master + Slave,同步等待 Slave 确认),在不修改任何参数的情况下升级到 RocketMQ 5,由于totalReplicasinSyncReplicas默认都为 1,将降级为异步复制(只需 Master 确认);
  • 如果希望保持与以前一致的行为(两副本均确认后才返回),则需要将totalReplicasinSyncReplicas均设置为 2

同理,原三副本同步复制的集群升级后,若要保持“三副本全确认”的行为,需要设置totalReplicas=3inSyncReplicas=3;若希望采用 Quorum Write 的“任一两个副本确认即可”策略,则设置totalReplicas=3inSyncReplicas=2。所有参数均在 broker 端(broker.conf 或启动命令行-c指定配置文件)配置。

参考

  • docs/en/QuorumACK.md:本文英文原始文档
  • docs/cn/QuorumACK.md:本文中文原始文档
  • MessageStoreConfig.java:totalReplicasinSyncReplicasminInSyncReplicasenableAutoInSyncReplicashaMaxGapNotInSync等参数定义
  • CommitLog.java:needAckNums计算与calcNeedAckNums实现
  • DefaultHAService.java:inSyncReplicasNumsisInSyncSlave实现
  • GroupTransferService.java:组提交等待服务
  • DefaultMessageStore.java:存活副本数维护
  • BrokerConfig.java:enableSlaveActingMaster开关
  • distribution/conf/broker.conf:broker 配置文件示例
  • 消息队列
  • 后端
  • 微服务
  • 流处理

【免费下载链接】rocketmq

Apache RocketMQ is a cloud native messaging and streaming platform, making it simple to build event-driven applications.

项目地址:https://gitcode.com/gh_mirrors/ro/rocketmq
点击查看免费下载

相关推荐

上一篇:LovyanGFX实战教程:从基础绘制到高级动画的完整示例
下一篇:Windows 风扇控制实战:用 FanControl 走完 5 个排障关卡

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/20 18:10:35

Nvivo12自动编码语言包实战:中文文本编码的规则与技巧

简介&#xff1a;Nvivo12自动编码语言包是面向定性数据分析研究者的语言支持组件&#xff0c;目的是帮助处理英语&#xff08;美国&#xff09;语料的用户提升自动编码识别准确率&#xff0c;解决大规模文本中主题分类耗时、人工编码负担重的问题。资源共计144个文件&#xff0…

作者头像 李华
网站建设 2026/9/20 18:08:57

Vortex模组管理器完全指南:从安装部署到冲突排查与性能优化

/* 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:08:10

安桥TX-NR636说明书实战:接线、AccuEQ校准与常见故障排查

简介&#xff1a;这是一份安桥TX-NR636功放的中文高级使用说明书&#xff0c;面向拥有该型号功放、希望充分挖掘其功能的中高级用户及家庭影院爱好者。内容涵盖AM/FM自动与手动调台、RDS电台信息显示、USB存储设备音乐播放、网络收音机&#xff08;TuneIn&#xff09;与DLNA串流…

作者头像 李华
网站建设 2026/9/20 18:07:57

MicroPython pyboard 入门指南:硬件布局、供电方式与首次上电

嵌入式语言运行时编程语言解释器编译器物联网系统编程 【免费下载链接】micropython MicroPython - a lean and efficient Python implementation for microcontrollers and constrained systems 项目地址&#xff1a; https://gitcode.com/gh_mirrors/mi/micropython 点击查看…

作者头像 李华
网站建设 2026/9/20 18:07:50

Aider 实战:TaoToken 跑通 Django 迁移脚本仓库任务

/* 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:04:52

答案格式总不对?TaoToken 这样查 GSM8K 评测请求链路

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

作者头像 李华