ARTICLE DETAIL

资讯详情

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

Raft在消息队列中的角色:从多数派确认到流处理可靠性

Raft在消息队列中的角色:从多数派确认到流处理可靠性 1. 消息队列的可靠性承诺靠的是多副本而不是运气凌晨两点零三分值班群里跳出一条告警某个消息队列Broker节点掉电。我盯着监控面板Producer写入量没掉消费者消费正常分区副本已经自动切到另一台机器继续服务。整个过程不到三十秒业务无感知。这种场景干过消息队列运维的人都懂它是Raft这类一致性协议在消息队列里真正落地后才能有的底气。今天想聊的就是Raft在消息队列中的应用尤其是在大数据流处理链路里它为什么被称为基石。很多人对Raft的印象停留在“一个分布式共识算法”知道它能选主、能同步日志但并不知道它在Kafka、RocketMQ这类消息系统里到底承担了什么角色也不知道为什么有了Raft流处理就能放心地依赖有序、不丢的消息流。这篇文章不讲虚的直接拆解Raft在消息队列里的工作机制、关键配置、还有那些文档里不会告诉你的坑。先说结论消息队列的可靠性从来不是靠“写盘了”这句话撑起来的而是靠多副本之间的一致协议撑起来的。单机写盘磁盘一坏全完蛋。多副本如果不讲一致性主节点确认了、从节点没同步完主节点一挂消息照样丢。真正的可靠性承诺来自“多数派副本都确认写入才算成功”这一条规则而Raft就是把这个规则工程化的最成功方案。1.1 单机消息队列为什么让人睡不踏实早期的消息队列很多就是单机部署。消息发过来写入本机磁盘返回成功。进程崩了还好数据还在磁盘上重启之后还能读出来。但怕的是磁盘物理故障那一块盘上的所有消息就跟人间蒸发一样没有任何挽救办法。生产环境里单机队列就是给自己埋雷无论文档里写得多好听。后来有了主从复制。一台主节点负责写一台从节点负责同步。看起来没问题但细想就慌了主节点收到一条消息先写进自己的磁盘然后返回Producer“成功”此时从节点可能还没同步到这条消息。紧接着主节点宕机运维把从节点顶上去最新那条消息就找不回来了。你在业务日志里看到Producer已经收到成功回包但消息就是不在了。这种“假成功”是异步复制方案最经典的暗坑。所以“多副本”只是第一步关键是副本之间的同步要满足什么条件才算完成。如果规定“主节点和从节点都写入成功才给Producer返回成功”那主节点宕机时数据大概率还在从节点上。但“几个副本算多数”又成了新问题。这时候Raft给出的答案是过半确认。1.2 主从复制与“多数派确认”的本质区别Raft里的核心规则之一是Leader必须把日志复制到集群中超过半数的节点通常是自己加若干Follower之后才能提交这条日志并向客户端返回成功。这个“多数派确认”和传统的主从复制有本质区别。传统主从可以类比成公司里的“领导说了算”。领导拍板了下面的人记不记得不重要反正领导说了算。问题在于领导一旦失联下面的人谁都不知道领导拍过什么板。多数派确认则是“核心团队集体记录”。一个决策要生效必须至少一半以上的人都亲手记下来了。这样即使少数人失联剩余的大多数手里仍然握着完整记录业务可以无缝继续。放到消息队列里Raft的多数派确认意味着Producer收到成功回包就已经等于“这条消息至少有超过一半的副本都持久化了”。这个保证听起来简单但它是消息不丢的唯一可信依据。没有这个依据后面所有流处理语义都无从谈起。1.3 为什么偏偏是Raft而不是Paxos很多人问过消息队列里用的一致性协议为什么大家最终都选了RaftPaxos不是更早吗这事得从工程角度去看。Paxos的理论正确性毋庸置疑但它太难落地了光是搞清Multi-Paxos里的几个规则就劝退了大半工程师更别提做出一个可运维、可调试的实现。Raft厉害的地方在于它把共识问题拆成了三个相对独立的小问题Leader选举、日志复制、安全性保证。这种拆解让工程师可以照着论文直接写代码。而且Raft对日志同步的语义定义得非常清晰——所有节点按顺序追加日志日志项一旦提交就不可变更新Leader必定拥有已提交的全部日志。这个语义和消息队列的需求简直天作之合。消息队列的分区Leader本来就是一个“单点写入、多副本备份”的模型。Raft恰好提供了一套天然适配这种模型的机制选出一个Leader负责接收写入其余Follower同步日志Leader变更时通过选举选新主。所以Raft不是“碰巧”出现在消息队列里它根本就是为这类场景设计的。2. Raft在消息队列里的工作逻辑一条消息怎样才算“安全”理解了多数派确认的价值之后下一个问题是Raft在消息队列里到底是怎么运转的从Producer发出一条消息到它真正成为“已提交且不会丢”的消息中间发生了什么这一段我会把完整的链路拆开讲顺便解释一下为什么offset本质上就是一种Raft日志。2.1 从Leader写入到多数派提交的完整链路消息队列的每个分区或者队列会选取一个副本作为Leader其他副本作为Follower。Producer只能向Leader发送消息这是Raft协议的核心约束——一切写请求都走Leader不能走Follower。这样设计是为了让日志追加的顺序在所有副本上保持一致。当Leader收到一条消息执行的操作大致是这样的先把消息追加到本地日志也就是写进自己的存储引擎然后再把这条消息发给所有Follower等Leader自己收到多数派包括Leader自身的写入确认之后才把这条消息标记为“已提交”并向Producer返回成功。这个过程和Raft论文里日志复制的流程几乎一模一样。你可能会问Leader把消息追加到本地日志为什么也算一个确认因为Raft的“日志”和“状态机”是分开的。写入日志不代表立即对外可见但它已经是多数派视图的一部分了。对消息队列来说“日志已提交”就等于“消息安全”。这也是为什么Kafka在底层做了大量日志顺序写的优化——因为它本质上就是在一套Raft风格的日志复制协议上堆硬件优化。2.2 offset到底是一种什么样的日志消息队列里最容易被忽略但又最重要的概念是offset。它表示一条消息在分区日志中的位置是一个单调递增的序号。很多流处理框架依赖offset做位点管理比如从某个offset开始消费、把消费进度提交到某个offset。但offset本身也是日志的一部分。这意味着什么在Raft的视角里offset不是一个“存储状态”而是一条条日志项的自然编号。消息写入到日志的第N位对应的offset就是N。所有副本的日志都必须保持一致所以同一分区内不同副本看到的offset一定是一致的。只要日志被多数派提交这个offset上的消息就永久存在。我见过一些刚接触流处理的人把offset当成“商品上的标签”觉得它只是给消息编个号。其实offset更像是“卷宗的页码”。Raft保证的是同一卷宗在多个档案馆里的页码完全一致任何人翻开第N页看到的都是同一条记录。这种一致性是后续所有“断点续传”“重放”操作的基础。2.3 消费进度为什么也需要一致消息本身的一致性解决了还有一块头疼的问题是消费进度Consumer Offset怎么办消费者每次处理完一批消息要把当前消费到的offset提交一下否则下次重启就会从头读。这一块如果各个节点各记各的流处理状态就乱了。很多现代消息队列把消费进度保存在一个内部的特殊主题里这个主题的写入和读取同样需要一致性。如果不同消费者实例看到的消费进度不一致就会出现“这个消息你处理了但另一个实例不知道”的情况。在流处理里这会造成状态错乱、重复计算、窗口数据对不齐等一连串问题。Raft在这里的用武之地是为保存消费进度的元数据层提供一致的日志复制。比如Kafka中的内部消费进度主题本质上是Kafka自己管理的一组分区分区的多副本一致性由ISR机制保证而在KRaft模式下元数据本身已经被Raft管起来了。总之消息队列里任何“状态”只要被多副本复制就需要一致协议来兜底。3. 大数据流处理场景中的Raft顺序、并行与恰好一次流处理引擎依赖消息队列不只是因为它能存消息而是因为消息队列能提供两样东西顺序性和可重放性。顺序性保证同一分区的消息按生产顺序到达可重放性保证出错了可以回溯。这两样恰好都是Raft模型带来的。3.1 流处理最怕乱序Raft如何守住顺序做过流处理的人都知道乱序是最难搞的问题之一。窗口聚合、状态更新、事件时间对齐全都依赖数据顺序。你写了一个每分钟滚动窗口的计数任务如果窗口边界前后的两条消息顺序颠倒了统计结果就是错的。Raft对日志顺序的保证是硬性的所有节点追加日志的顺序完全一致已提交日志的次序不可更改。映射到消息队列中就是同一个分区内的消息严格按照写入顺序排列消费者从Leader读取时看到的就是这个顺序。更重要的是Raft的Leader是一个全局唯一的存在写入都经过它天然串行化了消息热点问题也少。网上有人杠过Raft保证的是“日志顺序”不保证“业务时间顺序”。这话没毛病但要区分清楚。Raft保证的是消息在队列中的物理顺序也就是Producer发送的顺序。如果你本身是乱序发送的任何队列都救不了你。所以流处理任务通常要求业务侧按Key分区同一Key的数据进同一分区Raft才能保住它的相对顺序。3.2 分区机制与Raft副本组的组合拳大数据场景下单分区的吞吐远远不够所以消息队列搞出了分区机制一个主题拆成几十上百个分区每个分区独立读写并行消费。这个模型跟Raft的组组合拳是绝配。每个分区都可以看作一个独立的“Raft组”有自己的Leader和Follower。主题的N个分区就是N个独立的Raft复制组。它们各自持有日志、各自选主、各自保证一致性。Producer按Key把消息路由到不同分区同一个Key的消息永远去同一个分区于是并行和有序同时成立。我拿实际数据算过一笔账假设一个主题有64个分区每个分区有3个副本那么我们就有64个Raft组在同时工作。每组的Leader分摊写入压力Follower在后台同步。这样单主题就能打到百万级TPS而每一条消息都享受多数派确认带来的可靠性。如果哪个分区所在的Broker宕了只有这个分区的Leader会切换其他分区纹丝不动。故障半径被控制在单个分区级别这是大数据架构能持续稳定的重要原因。3.3 exactly-once的真相Raft只是其中一块拼图“恰好一次”Exactly-Once是流处理领域被讲得最烂也最容易误解的词。很多人以为消息队列用了Raft就一定保证消息不重不丢。这是极大的误解。Raft提供的保证是已提交的消息不丢并且所有副本看到的消息一致。但“不重复”这件事Raft管不着。真相是流处理系统宣称的“恰好一次”通常由多个层面共同完成。消息队列的Raft层负责“不丢”生产端开启幂等保证同一条消息重复发送不会产生重复数据消费端利用外部事务或状态存储实现“处理消息和提交offset”作为一个原子操作。这三个层面缺一不可。我们平时说的“至少一次”At-Least-Once和“至多一次”At-Most-Once差别就在于中间发生了多少次“重试”。Raft把重试变成可能因为日志还在offset可以回退重放但真正把重试“消化”掉的是下游的幂等逻辑。Raft是拼图的底座但绝不是整幅拼图。4. 实践中才能遇到的坑选主抖动、刷盘延迟与重复消费理论说得再好落地时总要踩几个坑。这一部分把我在维护消息集群过程中真正遇过的三个问题摆出来每一个都跟Raft机制直接相关。这些问题你不实操很难发现都藏在配置和运行细节里。4.1 一场扩容引发的Leader切回从“慢副本”说起前年给一套Kafka集群做在线扩容加了三个新Broker然后把一个热点的分区数从12调整到24。加完之后监控上出现了一个诡异现象部分分区的Leader频繁切换到其他Broker过几分钟又切回来就像钟摆一样来回跳。消费者消费时快时慢Producer写超时也开始出现。排查链路是这样的先看分区状态发现新增的Follower长期处于“不同步”状态再看网络指标新节点的入带宽被其他数据流占满日志复制一直跟不上最后排查副本追赶进度发现UnderReplicatedPartitions指标一直大于0。根因在于我把副本因子设成了3但新加入的节点同时承担了太多分区的同步任务网络成了瓶颈导致它们成为“慢副本”。Raft与Kafka的ISR机制有一点相同如果Follower落后太多会被移出同步副本集合此时Leader可以继续提交但可用副本数下降了容错能力变弱。如果配置不当比如min.insync.replicas保持为2网络抖动严重时会直接拒绝对外提供服务。这个坑的教训是扩容时一定要评估副本追赶的带宽消耗不要一次性把大量分区均衡到新节点上而且要提前调大副本同步的几个超时参数。4.2 fsync到底刷多勤Raft在性能和可靠之间的取舍Raft提交日志的第一步是Leader把日志写到本地存储。这个“写”到底写到哪里差别大了去了。操作系统会把数据先放进页缓存再由后台线程慢慢刷到磁盘。如果进程在这个时候崩溃页缓存里的数据可能就丢了。所以Raft类的系统几乎都要围绕fsync做文章。消息队列给用户提供了“刷盘策略”的开关。你可以配置每条消息都同步刷盘log.flush.interval.messages1这样可靠性最高但性能会掉一截你也可以配置成每隔一定时间或攒够一定消息数再刷盘换取吞吐。这个选择没有绝对的对错只有对业务场景的适配。但有一个关键点很多人不知道Raft多数派确认和fsync之间是有联动的。Leader给Follower发日志时如果Follower还没有真正把日志刷到磁盘就回复成功那么实际上的持久化承诺被“借”走了。这种“用页缓存假装刷盘”的做法可以换来更低的延迟但代价是一条日志只在内存里机器重启就没了。我个人的习惯是生产环境必须满足“至少一个副本真正刷盘、再配合多数派确认”这个底线。性能差一点换来的是不可丢失的承诺。4.3 消息重复消费的根因以及消息队列“不重复”的理想与现实“消息队列重复消费”是所有人迟早要面对的问题。网上关于这个问题的高频讨论我看了不少十有八九都在问“为什么我明明设置了ack消息还是重复消费”这个问题的根因它不在Raft而在消费确认的语义上。Kafka默认的消费模式是“至少一次”消费者拉取一批消息处理完业务逻辑再提交offset。如果消费者在处理完毕之后、提交offset之前崩溃了这条消息就会被重复拉取。Raft在这里能做的只是保证这条消息还在日志里可以再次被读到但它无法阻止下游再处理一次。所以“重复消费”这个热词背后真正需要解决的是消费端的幂等设计。解决重复消费的标准做法有三种第一种是业务侧做幂等比如用消息里的唯一ID去数据库查重第二种是使用事务API让“消息处理”和“offset提交”绑在一起提交第三种是消费端引入去重表处理过的消息ID直接跳过。这三种方法我都实践过最简单粗暴的是第一种但要求业务逻辑本身具备幂等性最稳妥的是第三种但会多一次存储查询。5. 从Raft看消息队列的演进ZooKeeper时代与新一代自管模式如果你只盯着“消息写入需要多数派确认”这个点觉得Raft就该这么用那你可能还漏了它更深刻的一面。Raft正在悄悄改变消息队列的系统架构。过去需要额外部署一套协调服务现在Broker自己就能管自己。5.1 用ZooKeeper协调Broker的旧玩法老一代的Kafka架构依赖ZooKeeper来选主、存元数据、做集群协调。Broker启动时要向ZooKeeper注册控制器由ZooKeeper选举产生分区的Leader变化也要写ZooKeeper。这带来了一个额外组件也带来了额外故障域ZooKeeper集群出问题Kafka就跟着出问题ZooKeeper的写入延迟变高Kafka的元数据操作就变慢。Raft并不是第一次出现在这个场景里。ZooKeeper本身用的就是ZAB协议和Raft高度同源。但ZooKeeper的设计目标是通用协调消息队列还得额外对接它。多一个组件就多一分运维复杂度我最头疼的运维问题里很大一部分不是Kafka自身而是ZooKeeper的JVM调参、会话超时、磁盘IO。5.2 KRaft与DLedgerRaft直接内嵌进Broker新一代的消息队列架构把Raft直接内嵌到了Broker进程里。Kafka的KRaft模式用自定义实现的Raft协议来存储元数据、管理控制器选举完全去掉了ZooKeeper。RocketMQ 5.x也引入了DLedger让Broker用自己的Raft组来做日志存储和主从切换不再依赖外部协调者。这个演进的好处是显而易见的部署组件变少了故障域缩小了Leader切换的路径变短了。以Kafka KRaft为例元数据变更由内部Raft组完成整个集群的控制面和数据面都统一了协议语义。RocketMQ DLedger模式下一个Broker组内的写请求先经过Raft日志复制提交后才能真正对外可见和Kafka的ISR逻辑殊途同归。配置层面有个例子RocketMQ DLedger的Broker组大概是这样的brokerClusterNameDefaultCluster brokerNameRaftBroker brokerId-1 enableDLegerCommitLogtrue dLedgerGroupNameRaftNode dLedgerPeersn0-127.0.0.1:40911;n1-127.0.0.1:40912;n2-127.0.0.1:40913 dLedgerSelfIdn0三个节点组一个Raft组多数派确认日志故障时自动选主。这种模式的好处是你不再需要为每个队列单独指定谁是主Raft自己决定。5.3 传统Windows/MSMQ队列与Raft消息队列的差距既然话题里有“Windows消息队列”“MSMQ”这些热词顺便提一嘴。MSMQ是微软早期提供的系统级消息队列主要面向单机和Windows域环境依赖本地系统服务跨平台能力弱扩展性也很有限。它的可靠性主要靠事务性队列和本地磁盘存储并没有Raft这种分布式多数派机制。拿MSMQ和现代基于Raft的消息队列对比差距可以用一张表格说清楚项目MSMQ基于Raft的现代消息队列高可用方式依赖Windows集群或异地转发多副本多数派确认故障切换手动或依赖系统群集自动选主秒级切换跨平台基本绑定Windows跨平台部署扩展模型单机队列为主分区并行扩展大数据流处理适配弱强这个对比想表达的是大数据流处理面对的是海量数据、跨地域部署、弹性扩容它的基础设施选型从一开始就不会考虑操作系统绑定的队列方案。Raft类协议消息队列能成为今天大数据流处理的事实标准本质上是把“分布式环境下的可靠性”从系统能力变成了可编程的、可验证的协议能力。做了这么多年消息队列和流处理基础设施我自己的体会是别急着啃Raft论文先把消息队列里acks、min.insync、副本同步超时这些参数调明白再回头去看Raft源码思路会顺很多。Raft不是高高在上的理论它就是Kafka、RocketMQ这些系统里你我每天都在用的那套可靠性逻辑。理解了它看消息队列的很多“奇怪现象”就不再奇怪了。
返回列表