- 消息队列
- 后端
- 微服务
- 流处理
【免费下载链接】rocketmq
Apache RocketMQ is a cloud native messaging and streaming platform, making it simple to build event-driven applications.
本篇文章基于 docs/en/QuorumACK.md 与 docs/cn/QuorumACK.md 展开,深入讲解 RocketMQ 5 在 Master-Slave 复制架构中引入的 Quorum Write(法定人数写入)与自适应降级机制。文章以该文档为骨架,结合当前仓库中store模块的源码实现(如CommitLog、DefaultHAService、MessageStoreConfig)进行佐证与补充,帮助读者理解:如何在 broker 端精确指定一条消息写入成功后至少需要多少副本确认,以及当副本掉线或落后过多时系统如何自动降级、保证可用性,同时明确各参数的含义、默认值、生效条件与配置方法。
背景:同步复制与异步复制的取舍
在 RocketMQ 的 Master-Slave(主备)架构中,主备之间的数据复制主要有两种模式:
- 同步复制(Synchronous Replication):Master 需要等待 Slave 成功复制消息并确认后,才向 Producer 返回写入成功。同步复制可以保证 Master 失效后,数据仍然能在 Slave 中找到,适合可靠性要求较高的场景。
- 异步复制(Asynchronous Replication):Master 不需要等待 Slave 的响应即返回成功。异步复制虽然可能丢失消息,但由于无需等待 Slave 确认,效率高于同步复制,适合对效率有一定要求的场景。
在消息发送过程中,客户端最终会收到如下几种结果状态:
| 状态 | 含义 |
|---|---|
PUT_OK | 一切顺利,消息写入成功 |
FLUSH_SLAVE_TIMEOUT | Slave 同步超时 |
SLAVE_NOT_AVAILABLE | Slave 不可用,或 Slave 与 Master 的 CommitLog 差距超过一定值(默认 256MB) |
其中后两种状态并不会导致系统异常而无法写入下一条消息,但它们都意味着当前副本组的同步状态并不健康。
然而,只有同步和异步两种模式在灵活性上存在明显不足:
- 在三副本甚至五副本且可靠性要求高的场景中,异步复制无法满足要求;
- 而同步复制需要每一个副本确认后才返回,副本数多时严重拖慢写入效率;
- 在同步复制模式下,如果副本组中某一个 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 的副本数。通过灵活设置totalReplicas和inSyncReplicas,可以满足各类场景对可靠性(副本确认数)与写入性能(等待确认的副本数)之间的平衡需求。
值得注意的边界语义
从源码注释可以确认以下几点细节:
- Master 被计入 in-sync 副本数:
inSyncReplicas的计数包含 Master 自身; - 异步 Master 忽略
inSyncReplicas:如果 Master 是ASYNC_MASTER(异步刷盘/异步复制角色),inSyncReplicas会被忽略; - 控制器模式下的特殊行为:如果
enableControllerMode=true且allAckInSyncStateSet=true,inSyncReplicas会被忽略,此时要求消息写入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=false且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(单条消息异步写入)与批量写入路径都会在真正追加消息之前计算本次发送需要多少副本确认(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; }这段逻辑的关键点在于:
inSyncReplicas(实际可用的 in-sync 副本数)取存活副本数与HA 服务统计的 in-sync Slave 数 + 1(Master 自身)两者中的较小值;- 当
enableAutoInSyncReplicas=true时,期望的needAckNums被限制在[minInSyncReplicas, inSyncReplicas]区间内——即“能等多少就等多少,但最低不低于 minInSyncReplicas”; - 当
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的常驻线程)完成,它维护requestsWrite与requestsRead两个请求链表,周期性地检查 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=2、inSyncReplicas=2、minInSyncReplicas=1、enableAutoInSyncReplicas=true。正常情况下两个副本均处于同步复制,消息需 Master 与 Slave 都确认;当 Slave 下线或假死时,系统进行自适应降级,Producer 只需发送到 Master 即成功; - 三副本场景:设置
totalReplicas=3、inSyncReplicas=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 的行为变化
为了保证向后兼容,用户升级后必须设置正确的参数。默认情况下,totalReplicas与inSyncReplicas均为 1,这意味着:
- 假设用户原集群为两副本同步复制(Master + Slave,同步等待 Slave 确认),在不修改任何参数的情况下升级到 RocketMQ 5,由于
totalReplicas、inSyncReplicas默认都为 1,将降级为异步复制(只需 Master 确认); - 如果希望保持与以前一致的行为(两副本均确认后才返回),则需要将
totalReplicas和inSyncReplicas均设置为 2。
同理,原三副本同步复制的集群升级后,若要保持“三副本全确认”的行为,需要设置totalReplicas=3、inSyncReplicas=3;若希望采用 Quorum Write 的“任一两个副本确认即可”策略,则设置totalReplicas=3、inSyncReplicas=2。所有参数均在 broker 端(broker.conf 或启动命令行-c指定配置文件)配置。
参考
- docs/en/QuorumACK.md:本文英文原始文档
- docs/cn/QuorumACK.md:本文中文原始文档
- MessageStoreConfig.java:
totalReplicas、inSyncReplicas、minInSyncReplicas、enableAutoInSyncReplicas、haMaxGapNotInSync等参数定义 - CommitLog.java:
needAckNums计算与calcNeedAckNums实现 - DefaultHAService.java:
inSyncReplicasNums与isInSyncSlave实现 - 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.
相关推荐
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. 痛点与解决方案概述 在分布式系统中,消息中间件的高可用性直接决定了业务连续性。
消息队列后端微服务流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考