ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

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

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),仅供参考
返回列表