ARTICLE DETAIL

资讯详情

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

基于ZooKeeper顺序节点实现轻量级FIFO队列的原理与实战

基于ZooKeeper顺序节点实现轻量级FIFO队列的原理与实战 前阵子有个任务调度模块把我折腾得够呛一批离线数据处理任务必须严格按提交顺序执行前一个任务不结束后一个就不能开始。业务方最开始用数据库表加状态位轮询查询越频繁锁冲突越多后来试过 Redis List又担心积压太多数据时内存和持久化扛不住。最后我回头看了一眼机房那套挂了很久的 ZooKeeper用顺序节点写了一个轻量级 FIFO 队列不到两百行代码问题迎刃而解。这篇文章就把这套方案从头到尾拆开讲顺带聊聊那些只有踩过坑才会懂的设计细节。ZooKeeper 在很多人印象里就是个分布式协调组件用处无非是配置中心、服务注册、分布式锁。但它最容易被忽视的一个能力就是利用顺序节点天然实现一个严格先来后到的 FIFO 队列。顺序节点那个自动追加单调递增序号的行为本质上就是给每个任务发了一个排队号码牌。理解了这个机制你不仅能自己手写一个队列还能彻底看懂 Curator、HBase、Kafka 老版本里那些基于 ZooKeeper 的队列和选主逻辑。1. 队列的刚需与 ZooKeeper 的独特位置1.1 哪些场景真的需要严格先来后到先别上来就写代码得先搞清楚一个问题你的业务真的需要 FIFO 吗我做过的项目里真正需要严格 FIFO 的场景通常长这样。一是任务调度比如用户批量提交了一堆数据清洗任务要求按提交顺序依次执行后面的任务可能依赖前面的产出。二是订单状态流转同一笔订单的创建、支付、发货、完成事件必须按发生顺序处理顺序乱了整个状态机就崩了。三是某些消息通知场景运营后台发了一串配置变更指令每条指令的执行结果会影响下一条所以必须排队。这类场景的共同特点是对顺序敏感任务量不算特别巨大但每条任务又需要可靠落盘。你当然可以用数据库实现给任务表加一个自增 ID然后ORDER BY id LIMIT 1轮询取任务。听起来没毛病但数据库轮询在高并发下会带来锁竞争而且取任务、改状态、释放任务这三个动作要做到原子性代码会越来越复杂。这时 ZooKeeper 的顺序节点方案反而更优雅。它把生成顺序号持久化数据感知任务变化这三件事都做了你只需要专注业务逻辑。1.2 对比 Redis 和 MQ我为什么选了 ZooKeeper如果你去技术群里问FIFO 队列用什么十个人里有八个会告诉你可以用 Redis List或者干脆上个 RocketMQ/Kafka。确实每种方案都有它的适用边界但 ZooKeeper 在特定场景下有自己的独特价值。Redis List 的RPUSH BLPOP天然就是 FIFO性能还极高。但 Redis 是内存数据库虽然有 RDB/AOF 持久化宕机时还是存在丢失窗口而且在大数据量积压时内存占用会非常恐怖。真正让我放弃它的原因是我们当时不想为一个每天几万条的小任务队列单独托管一套 Redis 集群运维成本不划算。RocketMQ 和 Kafka 当然是正规军功能强大、吞吐惊人。但它们的部署和维护比 ZooKeeper 重得多尤其 Kafka 的全局有序依赖单分区单分区吞吐有限实现严格全局限序其实并不方便。而且如果只是需要一个轻量级任务队列为了它引入一套完整消息中间件对很多中小团队来说是杀鸡用牛刀。ZooKeeper 的定位恰好卡在中间。它本身是强一致性的可以保证每个节点数据在所有节点间一致可见持久节点能把数据可靠写到磁盘顺序节点保证号码牌不重不漏Watcher 机制还能实现事件驱动消费。如果你所在的环境本来就有 ZooKeeper 在跑边际成本几乎是零。这也是很多团队在 Hadoop/HBase/Hive 生态里顺手用它做协调的原因。方案FIFO 支持持久化运维成本适合场景数据库轮询可以靠自增ID强低小流量、逻辑简单Redis List天然支持一般中高吞吐、内存足够消息队列 MQ分区内有序强高海量消息、复杂路由ZooKeeper 顺序节点天然支持强中轻量级分布式协调、任务调度1.3 为什么 ZooKeeper 在大数据生态里这么常见学 ZooKeeper 的时候你会发现它在 Hadoop 生态里无处不在。HDFS 的 NameNode 高可用用它做 Active/Standby 切换HBase 的 RegionServer 用它做元数据管理和故障发现老版本的 Kafka 用它的节点管理 Broker 和消费者组。甚至 HiveServer2 也会把配置和服务地址注册到 ZooKeeper 上客户端启动时再去读取。如果你在排查问题时见过类似unable to read hiveserver2 configs from zookeeper的报错本质就是客户端在 ZooKeeper 的某个路径下找不到配置节点要么是 ZK 连接不通要么是路径写错了。这个现象背后有一个共性ZooKeeper 提供的分布式协调原语很多都能用节点 监听的方式建模。队列只是其中一个应用但它几乎涵盖了 ZooKeeper 所有核心概念——节点类型、顺序序号、Watcher、一致性、会话管理。把队列搞透了你再看其他协调方案会轻松很多。2. 顺序节点队列的基石2.1 先认识四类节点ZooKeeper 的命名空间是一个树形结构每个节点叫做 znode。znode 可以存数据也可以有子节点。它一共有四种类型这是理解一切的基础。持久节点PERSISTENT是默认类型创建后一直存在直到显式调用 delete 删除。持久顺序节点PERSISTENT_SEQUENTIAL和持久节点的区别在于创建时指定这个类型ZooKeeper 会在你给的路径末尾自动追加一个 10 位数字的递增序号。临时节点EPHEMERAL的生命周期绑定创建它的客户端会话会话断开节点就自动消失。临时顺序节点EPHEMERAL_SEQUENTIAL则是临时 序号的组合。用命令行操作的话这四种类型分别对应create /queue data 创建持久节点 create -s /queue/msg- data 创建持久顺序节点 create -e /queue/lock data 创建临时节点 create -e -s /queue/consumer- data 创建临时顺序节点顺序节点创建成功后你实际拿到的路径可能长这样/queue/msg-0000000001。这个 10 位数字就是实现 FIFO 的关键。因为它是左补零存储的所以按照字符串字典序排序的结果和按照数字大小排序的结果完全一致。也就是说你拿到子节点列表后直接Collections.sort()排在最前面的就是序号最小的那个节点。2.2 顺序节点那个 10 位递增序号服务端是怎么分配的很多教学文章讲到这里就停了用-s创建顺序节点然后排序取最小。但为什么序号不会重复这个问题很少有人说清楚。在 ZooKeeper 服务端每个父节点维护着一个子节点版本号cversion。每创建一个顺序节点服务端会在父节点的 cversion 基础上递增并把得到的值格式化为 10 位数字追加到节点名后面。因为所有写请求在 ZooKeeper 集群内部都会经过 Leader 节点协调同一时刻只会有一个请求在执行创建操作所以序号是全局唯一且严格单调递增的。你连续创建 100 个顺序节点无论请求来自多少个客户端这 100 个序号都不会重复。不过这里有一个非常容易踩坑的细节cversion 不区分节点的类型。如果你在一个父节点下既创建顺序节点又创建普通节点、临时节点、锁节点那普通节点的创建也会推动 cversion 往前走。结果就是你创建的顺序节点序号会出现跳号。比如第一个顺序节点拿到 0000000001此时你创建了一个普通节点cversion 变成 2再创建顺序节点时拿到的就是 0000000003。很多人在测试时看到序号不连续就以为系统出了问题其实这完全正常。ZooKeeper 官方只承诺序号单调递增从来没承诺连续不断。2.3 临时顺序节点还是持久顺序节点关键看消息可靠性设计一个队列首先要决定用哪种节点存消息。这个选择直接决定了消息的可靠性语义。生产者的职责是往队列里丢一条消息消费者取走并处理。如果消息节点用的是临时顺序节点那么一旦生产者客户端会话断开消息就会自动从队列里消失。这很明显不靠谱因为你无法保证生产者发送完消息后不会断连。所以消息本身应该用持久顺序节点这样即使客户端和 ZooKeeper 之间网络抖动消息也已经可靠地躺在服务端磁盘上等消费者来取。那临时节点在队列里就完全没用了吗不是的。消费者注册、分布式锁占位这种场景临时节点才是首选。比如多个消费者竞争同一批消息时可以在取消息之前先创建一个临时节点作为我正在处理的标记会话断开这个标记就会自动消失别的消费者就能接手这条消息。这利用的正是临时节点生命周期随会话的特性。所以结论是队列消息用持久顺序节点协调和锁标记用临时节点两者配合使用。3. 从零手写一个 FIFO 队列3.1 环境准备与工程依赖开始之前先把环境准备好。本地开发建议直接跑一个 standalone 模式的 ZooKeeper下载解压后改一下配置就能启动。生产环境至少 3 台组成集群避免单点故障。学习阶段没必要折腾集群单机足够跑通代码。Java 工程只需要引入一个 ZK 客户端依赖。ZooKeeper 原生的客户端类org.apache.zookeeper.ZooKeeper就够用不需要额外装别的框架这样才能把底层逻辑看透。用 Maven 的话加这一个依赖dependency groupIdorg.apache.zookeeper/groupId artifactIdzookeeper/artifactId version3.7.1/version /dependency连接 ZooKeeper 时注意几个参数连接串写成127.0.0.1:2181多个节点用逗号分隔会话超时时间设置要根据业务处理时长来定太短会导致处理慢的时候会话被误判过期太长又会让客户端在故障后长时间处于僵尸状态。我习惯设成 30 秒。另外要设置一个默认 Watcher即使暂时不用也先放着用event - {}占位即可。3.2 生产者一条消息就是一个顺序子节点生产者端的逻辑非常简单一句话总结往队列根节点下面创建一个持久顺序节点节点数据就是消息内容。创建时指定的前缀名是什么无所谓关键是CreateMode.PERSISTENT_SEQUENTIAL。public class FifoQueueProducer { private static final String QUEUE_PATH /fifoQueue; public static void main(String[] args) throws Exception { ZooKeeper zk new ZooKeeper(127.0.0.1:2181, 30000, event - { }); // 确保队列根节点存在 if (zk.exists(QUEUE_PATH, false) null) { zk.create(QUEUE_PATH, null, ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); } // 生产 10 条顺序消息 for (int i 1; i 10; i) { String payload message- i; String nodePath zk.create( QUEUE_PATH /msg-, payload.getBytes(StandardCharsets.UTF_8), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT_SEQUENTIAL ); System.out.println(produced: nodePath - payload); } zk.close(); } }这段代码跑完之后你用ls /fifoQueue看会得到[msg-0000000000, msg-0000000001, ..., msg-0000000009]。每个节点里都存着一条消息内容。注意我第一次创建的时候是从 0 开始的因为父节点的 cversion 初始是 0。如果之前创建过别的节点序号就会从后面接着排。创建节点的 API 里有三个参数要说明一下。第一个参数是路径顺序节点实际创建时会在路径后面追加上序号。第二个参数是 ACLOPEN_ACL_UNSAFE表示完全开放开发环境用它最省事。生产环境建议用 digest 认证或者 IP 白名单只允许特定客户端写队列。第三个参数是 CreateMode就是选节点类型。这里用 PERSISTENT_SEQUENTIAL理由刚才已经讲过。3.3 消费者取编号最小的子节点消费者端的核心逻辑是获取队列根节点的所有子节点排序取第一个读取数据处理删除节点。这一步如果写成最简单的轮询版本大概是这个样子public class FifoQueueConsumer { private static final String QUEUE_PATH /fifoQueue; public static void main(String[] args) throws Exception { ZooKeeper zk new ZooKeeper(127.0.0.1:2181, 30000, event - { }); while (true) { ListString children zk.getChildren(QUEUE_PATH, false); if (children.isEmpty()) { Thread.sleep(500); continue; } Collections.sort(children); String headNode QUEUE_PATH / children.get(0); // 读取队头消息 Stat stat new Stat(); byte[] data zk.getData(headNode, false, stat); System.out.println(consuming: headNode - new String(data, StandardCharsets.UTF_8)); // 模拟业务处理耗时 Thread.sleep(200); // 消费完成删除节点 zk.delete(headNode, -1); } } }这里有个关键点要反复强调ZooKeeper 顺序节点的序号是 10 位左补零格式所以Collections.sort()的字符串排序结果和数字排序结果完全一致。你不需要自己写 Comparator 去解析数字直接用默认排序即可。如果序号没有补齐位数比如msg-1、msg-10、msg-2字符串排序就会变成 1、10、2队列就乱套了。所以使用 ZooKeeper 自己生成的顺序节点别手动拼序号。delete的版本参数我用的是 -1表示忽略版本号直接删除。这样最简单但存在覆盖别人改动的小风险。对于消息队列场景某个节点只会被消费一次业务上本来就该是排他的所以直接用 -1 也没问题。3.4 并发消费翻车实录同样一条消息两个消费者都拿到了单消费者跑通之后你可能会想多开几个消费者实例是不是能提高消费速度这里有个大坑。假设队列里有消息节点msg-0000000002此时两个消费者 A 和 B 同时执行了getChildren拿到的子节点列表一模一样排序后都认为msg-0000000002是队头。A 先执行getData拿到消息内容开始处理B 紧接着也执行getData同样拿到了消息内容也开始处理。A 处理完成后执行delete节点被删除B 处理完后执行delete必然抛KeeperException.NoNodeException。消息已经被两个消费者各处理了一遍这在业务上等于重复消费。你可以想象成窗口只有一个但是两个人都挤到了窗口前各自都拿到了叫号单结果业务员把排队号发给了两个人。这个问题的根源就是取队头和删除队头这两个动作没有合并成一个原子操作。要解决这个问题最简单的方式是使用分布式锁。用 ZooKeeper 自己实现一把取队头锁也非常简单消费者在消费某个消息节点之前先尝试创建一个临时节点作为锁比如/fifoQueue/lock/msg-0000000002。创建成功的消费者才有资格处理这条消息另一个消费者创建同名锁节点时ZooKeeper 会抛出NodeExistsException说明这条消息已经被别人抢走了它只能跳过重新取队头。临时节点的好处是如果持有锁的消费者突然宕机会话断开后锁节点自动消失消息节点还在其他消费者可以继续接手不会出现锁永远不释放的死锁问题。不过更优雅的解法是设计一个消费者排队的方案让消费者之间也按顺序节点排队从根上避免竞争。这个思路放到下一节讲。4. 让队列真正可用Watcher 驱动与惊群治理4.1 轮询的问题延迟与空转上面那版消费者代码虽然能跑但用起来会很难受。当队列里没有消息时消费者会一直Thread.sleep(500)然后再次getChildren。这就是典型的轮询。轮询的毛病很明显。第一是响应延迟生产者 00:00:00 时刻入队了一条消息消费者可能最多要等 500 毫秒才能感知到。如果你把 sleep 时间缩短到 100 毫秒延迟是降下来了但消费者的空转次数也变多了。第二是无意义的 ZK 请求会占满网络连接因为每个消费者都每秒钟发好几个getChildren请求集群的请求压力会白白增大。生产环境里如果队列根节点下的子节点很多每次getChildren还会返回全量子节点列表网络传输开销更大。一个靠谱的队列不该用轮询应该用回调。ZooKeeper 的 Watcher 机制就是为这种场景设计的你告诉服务端帮我盯着某个节点有变化就叫我服务端会在节点发生变化时推送一个事件给你。这样消费者可以一直阻塞等待不需要空转。4.2 Watcher 一次性触发实现有新消息再干活使用 Watcher 的核心思想是消费者先注册监听然后获取子节点列表并处理已有消息处理完后不退出继续等待下一个事件。写代码之前必须强调一个特性ZooKeeper 的 Watcher 是一次性的。它触发一次之后就会失效如果你想继续监听必须在回调里重新注册。很多初学者在这里栽了跟头调了一次回调之后后续消息再也不来了排查半天发现是 Watcher 没续上。一个可以工作的套路是这样的public class FifoQueueConsumerWithWatch { private static final String QUEUE_PATH /fifoQueue; private static final Object LOCK new Object(); public static void main(String[] args) throws Exception { ZooKeeper zk new ZooKeeper(127.0.0.1:2181, 30000, event - { // 任何事件到来时唤醒主线程重新处理 synchronized (LOCK) { LOCK.notifyAll(); } }); while (true) { // 注册 Watcher监听队列子节点变化 ListString children zk.getChildren(QUEUE_PATH, true); if (children.isEmpty()) { // 队列为空阻塞等待事件 synchronized (LOCK) { LOCK.wait(); } } else { // 队列非空取队头处理 Collections.sort(children); String headNode QUEUE_PATH / children.get(0); byte[] data zk.getData(headNode, false, null); System.out.println(consume: new String(data, StandardCharsets.UTF_8)); // 业务处理完成后删除节点 zk.delete(headNode, -1); // 继续循环重新 getChildren 并注册新的 Watcher } } } }这个模式虽然能用但严格来说还是有缺陷删除节点时触发的事件和你新注册的 Watcher 之间可能存在竞争关系极端情况下会丢失事件。更严谨的写法是删除操作也带上一个Watcher确保每次状态变化都被捕获。不过工程实践中上面这个版本的思路已经足够清晰能帮你理解注册 Watcher - 处理事件 - 重新注册这个循环。4.3 从惊群唤醒到链式唤醒如果你有多个消费者同时监听同一个队列根节点那么任何一条新消息入队所有消费者都会被唤醒然后竞争同一个队头。这就是典型的惊群效应。虽然分布式锁解决了重复消费问题但每次消息到达所有消费者都被打扰一次大多数消费者抢锁失败回去继续等待白白浪费了资源。更好的方案是让消费者之间也排成一个队。每个消费者启动时在/fifoQueue/consumers下创建一个临时顺序节点比如consumer-0000000001。所有消费者节点按序号排队序号最小的消费者拥有当前处理权。没有处理权的消费者只需要监听自己前一个消费者节点的状态前一个消费者处理完消息或退出后删除自己的节点后一个消费者收到NodeDeleted事件成为新的最小序号消费者开始干活。这个机制的本质是把所有消费者抢一个队头变成了消费者之间按顺序依次上岗。谁先注册谁先消费不会出现两人争抢也不需要分布式锁。当某个消费者实例宕机时它的临时节点会自动消失后一个消费者立即就会被事件唤醒并顶替上来。这种逐级唤醒的方式每次只唤醒一个消费者没有惊群问题。用 ZooKeeper 实现这个逻辑核心代码大致是这样的String myNode zk.create(CONSUMER_PATH /consumer-, sessionId, ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL_SEQUENTIAL); ListString consumers zk.getChildren(CONSUMER_PATH, false); Collections.sort(consumers); int myIndex consumers.indexOf(myNode.substring(myNode.lastIndexOf(/) 1)); if (myIndex 0) { // 我排在最前面可以消费 consumeHead(zk); } else { // 我排后面监听我前面的消费者节点 String prevConsumer CONSUMER_PATH / consumers.get(myIndex - 1); Stat stat zk.exists(prevConsumer, event - { if (event.getType() Event.EventType.NodeDeleted) { // 前一个消费者下线了重新竞争 } }); }这个链式唤醒看起来很漂亮但真正的生产级队列很少会自己手写这么一套复杂逻辑。因为 Curator 已经把这些设计模式封装好了。Curator 提供的DistributedQueue、DistributedPriorityQueue、DistributedDelayQueue等组件内部已经处理了分布式锁、会话重连、Watcher 重新注册等脏活。能用成熟组件的时候优先用组件。自己手写的意义在于理解原理而不是重复造轮子。5. 生产环境避坑指南5.1 序号不连续不是故障我在前面已经提过顺序节点的序号来自父节点的 cversion而 cversion 会因为该父节点下任意类型子节点的创建而递增。所以你的队列如果混入了锁节点、消费者节点序号跳号是正常现象。就算没有普通节点如果某条消息创建后又被删除cversion 已经往前走了一步下一条消息的序号也不会接着上一个已删除节点的序号。所以跳号完全不影响 FIFO 的正确性只要序号单调递增排序结果就是确定的。真正要小心的反而是序号位数问题。ZooKeeper 使用 10 位数字序列像msg-0000000001这样。只要创建的是 ZooKeeper 顺序节点补零是框架保证的。但如果你为了实现别的功能自己拼节点名比如把序号拿去当 ID 存到数据库取回来再拼路径一定要自己补齐位数否则字符串排序就会出错。5.2 Watch 回调后别忘重新注册这一点值得用一整节强调因为它太容易忽略了。ZooKeeper 的 Watcher 触发一次就会失效无论是getData、getChildren还是exists注册的 Watcher都是如此。如果你在回调里只做了处理消息但没有重新调用getChildren(path, watcher)注册新的 Watcher那么第二次消息入队时你的消费者是感知不到的。生产环境里还有更隐蔽的情况客户端和 ZooKeeper 之间的会话会因为网络抖动而重连。重连期间注册的 Watcher 可能会丢失。所以严谨的做法是每次处理完事件后立即重新注册并且在重连回调里也要重新注册一次全局 Watcher。如果你用的是 Curator这个问题框架会帮你处理但用原生客户端就必须自己扛。我自己吃过这个亏线上队列没有新消息触发排查了半天最后发现是重连后 Watcher 丢了。5.3 临时节点与幽灵消费者消费者节点如果使用临时节点注册在客户端异常退出时ZooKeeper 会很快清理掉这个节点。但很快不是立刻。客户端的会话过期检查需要一个超时时间在超时之前这个临时节点依然是存在的。换句话说一个已经宕机的消费者它的临时节点可能还要存活 10 到 30 秒。这段时间里如果这个宕机的消费者恰好排在最前面队列的消费就会停滞直到 ZooKeeper 判定会话过期并删除节点。为了避免这个窗口引发长时间停滞有两点经验一是合理设置 session timeout不能太长一般来说 15 到 30 秒比较合适二是消费逻辑要做到可重入因为节点删除后后续消费者可能会重新处理之前消费者已经处理到一半的消息。另外要特别注意持久节点不会自动清理。如果你用持久顺序节点作为消息节点消费者正常删除还好但如果消费者消费失败、代码 bug 导致删除逻辑没执行消息就会一直残留在 ZooKeeper 里。下线清理时别只盯着地自己机器的服务还要记得清除队列里积压的脏节点。5.4 子节点数量膨胀队列积压的隐患ZooKeeper 作为一个 CP 系统存储和性能是有边界的。虽然官方没有强限制但经验法则是单个父节点下的子节点数量最好不要超过十万级。如果你的队列积压了大量未消费消息每次getChildren都会返回庞大的子节点列表网络传输和客户端排序的开销都会急剧上升。一个实用的治理手段是分片。比如把队列根节点按时间分桶/fifoQueue/2024-01-01、/fifoQueue/2024-01-02消息创建时落到对应日期分片消费者优先消费最早的分片。这样每个分片下的节点数被控制在一个合理范围内积压再严重也不会让某个分片变成超大节点列表。还有一个通用建议队列里放的消息数据本身要小ZooKeeper 不是给大对象设计的单节点数据建议控制在 1MB 以内最好只有几 KB。任务详情放数据库或对象存储ZooKeeper 节点里只保存任务 ID 和必要元数据消费时再回源查详情。5.5 Curator生产上更推荐的封装如果你手写了一遍上面的代码对 ZooKeeper 的机制已经有了足够的体感那我可以很明确地建议生产环境直接用 Apache Curator它是一个更高层的 ZooKeeper 客户端封装把分布式锁、队列、Leader 选举都做好了。用 Curator 实现一个分布式队列非常简洁。声明一个QueueBuilder传入队列根路径、序列化器和消费者回调然后调用buildQueue()就能拿到一个DistributedQueue。它的内部实现核心就是顺序节点加分布式锁所有并发竞争、会话重连、Watcher 重新注册都由框架处理你只需要关注业务消费逻辑。CuratorFramework client CuratorFrameworkFactory.newClient( 127.0.0.1:2181, new ExponentialBackoffRetry(1000, 3)); client.start(); QueueBuilderString builder QueueBuilder.builder(client, createConsumer(), 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); } }, /curatorQueue); DistributedQueueString queue builder.buildQueue(); queue.put(task-1);我自己在线上使用 Curator 的DistributedQueue时会额外注意它的消费线程模型默认消费是在回调线程里执行如果业务处理较慢要考虑线程池配置和背压控制。同时Curator 对于连接中断也有自己的重试策略建议配置成指数退避避免断连时所有请求瞬间重试压垮 ZK 集群。6. 进阶扩展从队列到分布式协调的更多可能6.1 顺序节点还能做什么理解了顺序节点之后你会发现它能做的事情远不止一个队列。分布式 ID 生成器就是最直接的一个应用利用顺序节点创建后的返回路径提取末尾的序号可以做出一个全局单调递增的 ID。虽然它的吞吐量不如 Snowflake 算法但在某些需要强一致、按顺序排号的场景下非常好用。分布式锁也是最典型的应用之一甚至可以说是 ZooKeeper 的看家本领。用临时顺序节点实现公平锁每个客户端创建一个临时顺序节点序号最小的客户端获得锁没拿到锁的客户端监听自己前一个节点的删除事件被唤醒后重新判断是否能拿锁。这套逻辑和前面讲的消费者排队几乎一模一样只是场景换成了抢锁。还有一个小技巧利用节点顺序号可以判断两个事件的先后顺序。因为序号由 ZooKeeper 统一分配所以谁先创建谁序号小。这在某些需要判定谁先谁后的分布式场景里很有价值比如确定版本、确定 leader lease 的归属。6.2 结合业务场景的设计建议最后给出一些我实际踩坑后总结出来的设计建议希望能帮你少走弯路。一是消息数据尽量轻量化ZooKeeper 节点里只放必要信息详情数据放外部存储。二是消费逻辑必须支持幂等即使有锁和顺序保证网络分区、会话过期这些极端情况仍会导致重复消费业务侧要做同一任务重复执行不产生副作用的防御。三是监控一定要做好队列积压是最直观的告警指标可以用 ZK 节点数量统计也可以把入队和出队次数埋点上报到监控系统。四是优先使用 Curator 等成熟封装手写用于学习和场景定制生产环境用框架更稳。这套方案能不能扛住千万级消息说实话不能。ZooKeeper 不是消息中间件它的设计目标是协调一致性而不是高吞吐消息分发。如果业务量真到了那个量级还是老老实实上专业的消息队列。但像任务调度、轻量级协调、事件排序这类场景ZooKeeper FIFO 队列简单、可靠、够用而且能沉淀下来一整套分布式协调的思维方式。这种思维方式才是比队列本身更值钱的东西。我个人的体会是ZooKeeper 的 API 看起来简单真正难的是理解它背后的几个模型节点模型、会话模型、Watcher 模型、一致性模型。这四个模型反复出现在各种分布式系统里无论是 Kafka 的协调器、HBase 的 RegionServer 管理还是各种分布式锁和队列底层都是这一套东西的组合。把顺序节点这个点彻底吃透你再看其他组件的源码很多当时看不懂的代码会一下子豁然开朗。
返回列表