多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

Zookeeper分布式队列原理与实现:从顺序节点到FIFO队列的工程实践

Zookeeper分布式队列原理与实现:从顺序节点到FIFO队列的工程实践 做大数据的同学接触Zookeeper的时间点通常很有意思——装Hadoop集群你会先把它起起来配置HBase的时候你得确认它是健康的跑Kafka集群的时候还会单独给它准备三台机器。可如果你追问这些框架到底拿ZK做了什么很多人只会含糊接一句用来做协调的。我再往下问一句ZK自己能不能实现一个分布式队列它和Kafka的队列有什么区别能把这个问题讲透的人确实不多。这篇文章就围绕“Zookeeper分布式队列”展开。我会先把ZK的核心机制用大白话讲透再带你用原生API从零手写一个能跑的FIFO分布式队列然后拆解入队出队背后那些容易翻车的边界细节最后把ZK队列放回大数据生态里看它适合站在哪里、不该出现在哪里。适合正在学Zookeeper入门、想动手写分布式组件但没写过的人也适合准备大数据开发面试、想理解ZK原理细节的同学。1. 先搞清楚Zookeeper在大数据生态里的真实定位1.1 它不是一个存储库而是一个“分布式协调记事本”第一次打开zkCli执行ls /看到只有zookeeper一个节点时大部分人的第一反应是这不就是一个简陋的文件系统吗能建节点、能写数据、能删节点好像没啥特别的。这个直觉没有错但要理解ZK必须记住一点它不是为存储业务数据而生的它的核心能力是让分布式系统里的一群机器“对某件事达成一致”。打个比方。你们部门要确定今天的发布顺序如果只有三个人群里吼一声就行。但如果一百多人分散在三个城市每个人都吼一嗓子现场会乱成一锅粥。ZK提供的就是一个“带公告栏的会议室”任何人想发布一条信息必须在公告栏上写清楚任何人想知道当前状态直接到公告栏读谁先写、谁后写由会议管理员统一编号。这个管理员就是ZK的Leader编号规则就是ZAB协议。这句话是大数据场景理解ZK的钥匙ZK的每一个写操作都会被赋予全局唯一的序号而且这个序号在所有副本上完全一致。分布式队列能成立靠的就是这个特性——入队的本质是“让公告栏按顺序产生编号”出队的本质是“所有消费者看到同一份编号清单按序号依次处理”。所以你看ZK做队列不是用它存了多少数据而是用它把“先后顺序”这个最难分布式化的东西给确定下来了。1.2 大数据框架到底让ZK干了什么把ZK放进整个大数据生态里看它的工作基本可以归成三类元数据协调HBase把meta表所在的RegionServer位置告诉ZK客户端启动时先问ZK再决定去哪个RegionServer拿数据Kafka早期版本把broker的注册信息、topic的分区分配信息都放在ZK上生产者和消费者通过ZK感知集群变化。Leader选举HDFS的NameNode Active/Standby切换靠ZK选主HBase的HMaster选主也走ZK。谁拿到临时节点谁就是老大节点一删就触发新一轮选举。状态感知与发布Flink的Standalone集群高可用方案里JobManager通过ZK发布当前运行状态Spark也有类似机制。ZK在这里扮演的角色就是一个所有节点都能访问的“状态公告栏”谁挂了、谁活了大家都看得见。想清楚这些再看ZK队列它本质上就是“状态公告栏”的一种具体用法。队列里不是堆着海量消息而是堆着一串“任务顺序编号”。节点上的少量数据可能只是任务描述、元信息或者ID真正的大payload不放在ZK里这是ZK队列和消息队列在定位上最根本的区别。2. 队列能搭起来全靠这套节点模型2.1 znode既是文件又是目录ZK的数据模型是一棵树每个节点叫znode。它像文件系统里的“文件加目录二合一”可以像目录一样有子节点也可以像文件一样存一段数据。你可以把每个znode理解成一个带版本号的独立对象每次修改版本号都会递增配合-1这种特殊版本号就能实现CAS式的乐观锁操作。有一点必须养成习惯ZK里的路径是标准Unix风格比如/task-queue/task-0000000001而且每个节点只能通过完整路径访问。节点上能存的数据大小有上限默认1MB左右。这个限制决定了你不可能把消息体本身往ZK里塞只能塞元数据或者小状态标识。2.2 四种节点类型重点记住“临时”和“顺序”理解了znode再看节点类型就轻松了。ZK提供四种创建模式节点类型创建方式核心特点典型用途持久节点PERSISTENT不主动删除就永远存在队列根路径、固定配置信息临时节点EPHEMERAL会话断开自动删除在线状态、Leader标识持久顺序节点PERSISTENT_SEQUENTIAL带全局唯一递增序号FIFO队列入队临时顺序节点EPHEMERAL_SEQUENTIAL带序号且会话断开自动清理分布式锁、公平型队列做分布式队列时“顺序”二字是关键中的关键。PERSISTENT_SEQUENTIAL模式下你创建子节点时不用自己算序号ZK会在你指定的路径后缀自动追加一个10位数字的递增序号比如task-0000000001。这个序号由ZK内部的ZAB协议保证全局唯一、严格递增不会因为并发创建而错乱。至于用临时还是持久取决于你的业务语义。如果任务不允许因为客户端掉线就消失用持久顺序节点如果希望客户端掉了以后任务自动清理用临时顺序节点。这个选择没有绝对对错但要先想清楚。2.3 watch监听机制一次性的订阅通知第四个核心机制是watch。客户端可以给某个节点或者某个父节点的子节点列表挂上一个监听器当节点数据变化、子节点增删、节点被删除时ZK会向客户端推送一个事件通知。这个“通知”不包含具体数据只告诉你“变化发生了”你还得自己再调用一次getData或getChildren去拿最新内容。再补一句容易忽略的watch是一次性的。你处理完一件事件后如果不重新注册后续变化就再也收不到了。这个特性直接影响了队列消费者的写法后面讲代码的时候你会看到它怎么逼你把轮询和监听组合起来。3. 基于原生API手写FIFO分布式队列3.1 队列的目录结构设计动手写之前先画清楚结构。假设队列根路径是/task-queue它必须是持久节点否则根节点一旦消失挂在上面的所有任务子节点全都没了。/task-queue 持久节点队列根 ├── task-0000000001 任务A ├── task-0000000002 任务B └── task-0000000003 任务C入队就是创建带顺序号的子节点出队就是找出序号最小的子节点、读取数据、删除节点。整个设计几十行代码就能落地完全不依赖第三方框架只用一个ZooKeeper客户端对象。3.2 入队核心就是一行create入队逻辑最简单public void enqueue(byte[] taskData) throws KeeperException, InterruptedException { String childPath queueRoot /task-; String createdPath zk.create(childPath, taskData, ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT_SEQUENTIAL); System.out.println(workerName enqueued: createdPath); }看到/task-这个路径后缀了吗你传的路径不需要带完整序号ZK会自动生成并返回带序号的完整路径。这里我只用了两个参数需要注意ACL用了最宽松的OPEN_ACL_UNSAFE生产环境要按集群配置收紧节点模式选了PERSISTENT_SEQUENTIAL这样即使入队客户端崩了任务节点也不会消失。3.3 出队争抢最小号学会处理竞争出队要复杂一点因为多个消费者可能同时看到同一个“最小序号节点”。我用的方法是标准的最小序号竞争模式public byte[] dequeue() throws KeeperException, InterruptedException { while (true) { ListString children zk.getChildren(queueRoot, false); if (children.isEmpty()) { return null; // 队列为空 } Collections.sort(children); // 序号升序 String smallestPath queueRoot / children.get(0); try { byte[] data zk.getData(smallestPath, false, null); zk.delete(smallestPath, -1); return data; } catch (KeeperException.NoNodeException e) { // 别人抢先删了循环重来 } } }流程拆开看是这样的调用getChildren(queueRoot, false)拿所有子节点名。因为子节点名天生带10位序号直接按字符串排序就能得到正确顺序。取第一个尝试getData后再delete。如果delete抛出NoNodeException说明另一个消费者已经把节点抢走删掉了本次竞争失败回到循环重新取。这个while(true)加异常重试的模式看着简单实际就是ZK分布式锁里最常见的自旋方式。它保证同一时刻最多只有一个消费者能成功删除最小序号节点删除成功的那个人就是当次任务的消费者。不需要额外的锁因为“删除节点”这个操作本身在ZK里就是原子的、排他的。3.4 跑一个完整Demo一个生产者两个消费者有了上面两个方法串一个主流程public class QueueDemo { public static void main(String[] args) throws Exception { DistributedQueue producer new DistributedQueue(localhost:2181, /task-queue, producer); producer.initRoot(); producer.enqueue(任务A.getBytes()); producer.enqueue(任务B.getBytes()); producer.enqueue(任务C.getBytes()); DistributedQueue consumer1 new DistributedQueue(localhost:2181, /task-queue, consumer-1); DistributedQueue consumer2 new DistributedQueue(localhost:2181, /task-queue, consumer-2); consumer1.initRoot(); consumer2.initRoot(); byte[] task1 consumer1.dequeue(); byte[] task2 consumer2.dequeue(); byte[] task3 consumer1.dequeue(); System.out.println(consumer-1 got: new String(task1)); System.out.println(consumer-2 got: new String(task2)); System.out.println(consumer-1 got: new String(task3)); } }跑起来以后的效果通常是consumer-1拿到任务Aconsumer-2拿到任务Bconsumer-1再拿任务C。顺序完全按入队序号来谁抢到完全看运气但整体先进先出不会乱。这就是ZK分布式队列最朴素的形态能跑、能演示、能帮你理解原理但它离生产还差好几步下面这些坑我一个个说。4. 队列实现里最容易翻车的四个细节4.1 惊群效应不要每个人都盯着同一个父节点上面那版dequeue是纯轮询逻辑简单但效率低。你要是把getChildren改成带watch的写法让所有消费者都监听/task-queue的子节点变化就会出现典型的惊群效应每入队一个新任务ZK给所有消费者各推一条通知所有消费者同时醒来、同时调用getChildren、同时去抢那一个最小节点结果99%的消费者白醒一场。消费者多、任务频繁时ZK会被这种无意义通知打爆。优化思路是这样的每个消费者不要监听整个父节点而是只监听“比当前最小号更小的前一个节点”。比如你是消费者A上次处理到task-0000000002你就只watchtask-0000000001的变化它被删了才说明轮到你出手。这个做法的本质是把通知范围从“全体广播”压缩到“精确单播”和ZK分布式锁里规避羊群效应的思路完全一致。如果只是想快速实现一个没问题的小工具轮询加短暂sleep也能凑合但你要意识到这是拿CPU换简单。生产环境我会建议先想清楚“任务量到底多大”量小轮询完全够用量大了直接上Curator封装别自己造轮子。4.2 会话过期与节点丢失临时节点是一把双刃剑用EPHEMERAL_SEQUENTIAL做队列会有一个隐藏风险消费者处理一个任务时如果客户端和ZK的会话超时了它创建的临时节点会被ZK自动清理任务呢你以为它丢了其实它在被别人持有的一瞬间就已经从队列里删除了没丢只是“没人知道这个任务到底处理到哪一步了”。这个问题真正麻烦的地方在于ZK的会话超时是一个window不是立刻判死。客户端可能在十几秒内实际上还活着但ZK已经认为会话死亡并把所有临时节点清掉。对于任务队列来说这意味着一批节点可能被“误删”。所以我自己的倾向是队列任务节点用持久顺序节点业务层自己维护任务状态不要把任务存活的命运寄托在ZK会话上。保证数据可用性再谈一致性。4.3 消费确认删节点不等于任务成功上面dequeue里“删除节点即消费完成”的方式非常粗暴。如果消费者拿到数据、删掉节点之后还没来得及执行任务就宕机了这个任务永远回不到队列里。这在消息系统里叫“至少一次”语义都没满足严格来说连at-most-once都算不上因为连“至少执行”的保证都没做。要解决就得引入任务状态字段。我比较常用的做法是给节点加一段状态信息task-0000000001节点数据里写任务描述初始状态为PENDING消费者读到PENDING后先尝试把它更新为PROCESSING更新成功才算抢到任务执行完再把节点删除或者把状态改成DONE如果消费者中途崩了由另一个巡检线程找出长时间停留在PROCESSING的节点重新置回PENDING或直接转移到死信队列这样队列就从“一个节点一份数据”升级成“一个节点一份任务状态机”。这也是生产级任务分发系统的通用设计和MQ里的消费位点、ack机制做的是同一件事。ZK能给的只是操作的原子性但“任务到底成功没有”这种业务判断只能自己写。4.4 1MB数据限制别把队列当仓库最后一个常被忽略的点znode默认数据上限约1MB可通过jute.maxbuffer调大但不建议。很多人写个Demo觉得挺好用就把任务内容整个塞进去任务一多、数据一大写操作直接在Leader上被打回。我的习惯是节点里只放任务ID、任务类型、优先级这类小字段真正的任务体放HDFS、对象存储或者MySQL消费者拿到ID后再去外部拉。毕竟ZK每次写都要走一遍ZAB协议达成多数派共识本质上是“慢而可靠”的路径不是给大流量数据准备的。5. 大数据场景下ZK队列的真实用武之地5.1 别拿ZK队列和Kafka比它们不是一个物种聊到这儿你应该能感受到ZK队列的吞吐天花板。每个写操作都要Leader广播给Follower过半确认才算成功顺序节点虽然能全局排序但排序本身也是成本watch通知一次一配消费者多了还会惊群。这套机制决定了ZK队列的QPS上限低得可怜和Kafka动辄十万百万级的写入吞吐完全不在一个量级。维度ZK队列Kafka等MQ底层模型树形节点全局序号分区顺序日志吞吐量受ZAB协议制约适合低频分区并行高吞吐数据容量节点上限约1MB不适合大消息海量消息持久化可靠性依赖会话、watch机制副本同步、消费位点、ack机制典型场景任务分发、分布式锁、节点协调业务消息、日志管道、事件流选型上不需要纠结业务数据流、实时日志、离线和在线系统解耦用Kafka或Pulsar分布式系统内部的低频协调动作、需要全局顺序的任务分发才考虑ZK。拿ZK队列替代Kafka属于典型的用错工具。5.2 真正值得用ZK队列的场景低频任务分发和批量调度那ZK队列在大数据领域到底能干什么我实际用过的场景是任务分发在自研的数据血缘采集系统里每天凌晨调度平台会生成几百个采集任务需要分发给多个worker节点去执行。这些任务一次生成、执行时间长、总量不大但对“不能重复执行”要求很高。用ZK队列后每个采集任务一个顺序节点worker启动后抢最小号节点抢到就执行执行完删节点。天然具备全局唯一性、顺序性和故障感知不需要另外引入MQ运维成本少一大截。再比如批量指令下发给一批边缘节点下发同一份配置变更先把指令按节点ID写入ZK队列各节点按自己的顺序号来拉取执行执行进度也能实时写入对应节点管理端一看就知道谁没完成。这种场景如果用Kafka消费位点和任务状态全要自己写反而麻烦。还有一个角度值得说出来ZK队列本质上就是一个公平锁的变体。你仔细看dequeue逻辑它和ZK分布式锁里“创建临时顺序节点、监听前一个节点”的思路一模一样。理解了队列就理解了大半个ZK分布式锁。这也是我建议大数据初学者必须先动手实现一次队列的原因。5.3 生产建议直接用Curator的DistributedQueue自己用原生API写队列适合学习和面试生产环境我建议直接用Apache Curator它是ZK官方推荐的客户端封装库内部提供了DistributedQueue和DistributedDelayQueue。Curator不仅把watcher重复注册、连接重试、会话重建这些脏活全做了还在DistributedQueue里实现了多消费者下避免惊群的内部逻辑并且支持设置节点数据上限、任务消费方式等参数。用Curator大概长这样CuratorFramework client CuratorFrameworkFactory.newClient( connectString, new ExponentialBackoffRetry(1000, 3)); client.start(); QueueBuilderString builder QueueBuilder.builder( client, new QueueConsumerString() { Override public void consumeMessage(String message) throws Exception { System.out.println(consume: message); } Override public void stateChanged(CuratorFramework client, ConnectionState newState) { // 处理连接状态变化 } }, new QueueSerializerString() { Override public byte[] serialize(String item) { return item.getBytes(StandardCharsets.UTF_8); } Override public String deserialize(byte[] bytes) { return new String(bytes, StandardCharsets.UTF_8); } }, /task-queue); DistributedQueueString queue builder.build(); for (int i 0; i 10; i) { queue.put(task- i); }Curator把create、getChildren、delete这些底层调用全部封装起来你只需要关心序列化和消费回调。它内部对watcher的重置、连接断开的保护都处理好了是我见过最不容易写错的ZK队列入口。如果你只是想在项目里快速用一个可靠的任务队列直接走这条路就好没必要自己处理那些边角细节。6. 我在项目里用ZK队列的实际体会最后说点个人感受。我做过一个数据采集调度模块一开始方案准备直接用数据库表做任务队列用SELECT ... FOR UPDATE锁行来保证不重复消费。结果任务一多行锁竞争、死锁排查搞得人很头疼。后来换成ZK队列代码量反而更少根节点一个、任务一个子节点、消费者抢最小号、处理完删除整套逻辑清晰得多出问题也容易在ZK上直接观察节点状态。那次经历让我意识到一件事分布式系统里很多让人抓狂的问题根源都是“多个节点对同一件事的先后顺序无法达成共识”。ZK队列解决的就是这个问题它不追求高吞吐但追求“顺序被所有人认可”。你把它用明白以后再看HDFS的NameNode选举、Kafka早期版本的controller选举、HBase的RegionServer上下线通知会发现核心思想都能对上——无非是节点、监听、序号、会话这套组合拳。如果你正在入门大数据我的建议很直接不要只停留在看文档和搭集群动手用原生API写一个分布式队列然后故意杀掉一个消费者、故意让会话超时观察节点状态怎么变化。这个过程跑一遍比你刷十篇教程都有用。以后面试被问到Zookeeper能做什么、分布式锁怎么实现、顺序节点有什么含义你都能用自己的代码经验答出来而不是背概念。
返回列表