凌晨1点17分,我收到一条告警短信:ods_orders_inc这个增量任务在过去10分钟内的产出量比正常值高了一倍。
一开始我以为是上游业务发生了大促,点开数据面板后才发现,同一批订单数据被两个Worker同时拉了两遍,落到下游表里生成了重复记录。更让人恼火的是,这个问题的根子不在数据同步插件,而在调度框架本身:备调度器通过数据库心跳判断主调度器“挂了”,结果顶上来以后,主调度器其实还在正常运行,于是两个主节点同时派发任务,下游自然就乱了。
那之后我花了很长时间重构了整个ETL调度链路,核心就是引入了ZooKeeper来承担分布式协调的工作。在这个过程里我意识到一件事:很多做数据平台的同学对ZooKeeper的印象还停留在“Hadoop里NameNode做HA用的组件”,但它在自研/开源大数据ETL工具里的角色远比想象中重要——任务分片防重、主节点选举、执行器注册、动态节流、检查点保存,几乎每一层都离不开它。
这篇文章我就从实际落地经验出发,把ZooKeeper在ETL工具中的几个典型应用场景拆开讲清楚,包括每类场景的架构思路、代码写法、参数选择和真实踩坑。
1. 为什么离线数仓的“调度控制链条”需要ZooKeeper
1.1 单靠关系型数据库做任务协调的局限
很多初版ETL调度平台的设计都很简单:数据库里放一张task_instance表,多个Worker通过SQL去Update这张表来抢任务。比如一个任务要被某个Worker执行时,Worker先执行UPDATE task_instance SET owner=?, status='running' WHERE id=? AND status='pending',如果影响行数为1,就认为抢到了。
这个方案在任务量小的时候没毛病,一旦任务规模和Worker数量上来,麻烦就一个接一个地出现。
第一,数据库的锁机制在高频抢任务下容易拖慢整个调度链路。几十个Worker定期轮询并Update同一张表,行锁竞争和死锁日志会变得非常频繁。第二,数据库记录任务状态只是“某个时刻的拍照”,它无法判断一个Worker到底是真死了,还是只是网络抖动。基于最后心跳时间判定Worker是否可用,在主备切换时特别容易误判,我开头说的“双跑事故”就是这么来的。第三,关系型数据库没有“临时数据”和“监听变更”的天然语义,如果某些任务执行到一半Worker宕机,数据库里的状态最终要靠另一个清理线程去扫,响应速度不够实时。
1.2 ZooKeeper真正提供的是哪几项能力
用ZooKeeper做ETL协调,很多人第一反应是“又一个分布式锁组件”。这也没错,但如果只把它当锁用,就太亏了。ZooKeeper在ETL场景里真正值钱的,是下面这几个原语:
- 临时节点(Ephemeral Node):节点和客户端会话绑定,客户端会话超时后节点自动删除。这个特性天然适合做“活体检查”——Worker还活着,节点就在;Worker失联,节点就消失。
- 监听通知(Watcher):客户端可以监听某个节点的变化,一旦节点被创建、修改、删除,服务端会推送事件。这比轮询数据库高效得多。
- 顺序节点(Sequential Node):在同一个父节点下创建带序号的子节点,能保证全局的创建顺序。分布式锁、Leader选举、公平任务队列都依赖这个特性。
- 数据版本号(Version):每次数据更新都会带版本号,配合
版本不一致则更新失败的机制,可以做成乐观锁,避免配置信息被并发覆盖。
1.3 一句话:ZooKeeper在ETL里是“控制通道”
我习惯把整个数据集成平台分成两条链路:数据流和控制流。数据流是从源端数据库、消息队列、日志系统把业务数据搬到数仓,走的是DataX、Canal、Kafka Connect、Flink CDC这些组件;控制流负责“谁来做、什么时候做、做哪些数据、做到什么进度”,走的是调度器、任务队列、状态判断和故障转移。
ZooKeeper不参与数据流的搬运,它只是控制通道上的一个协调中心。
有了这个定位,后面就好设计了:所有需要多台机器对同一件事达成共识的地方,才放ZooKeeper;所有高频、大批量的数据临时状态,不往ZooKeeper里塞。
2. 任务分片抢锁:一台机器只能有一个Worker执行同一份分区数据
2.1 原始场景:12个分片被多个Worker并发领取
我自己维护的同步服务有一个很典型的分区并行模型:每天凌晨把MySQL的业务大表按主键范围切成多个分片,每个分片是一个子任务,写入待执行列表,由多个Worker并发拉取。比如ods_orders_inc这次同步,一共被切成12个分片,每片大概处理80万条增量数据。
问题是:多个Worker同时扫描待执行列表时,很容易抢到同一个分片。
如果把“分片是否被领取”放在关系型数据库里做标记,就回到第一部分的轮询和锁竞争问题。如果直接依赖Redis的SETNX来加锁,又要额外处理“锁超时后任务还在跑,另一个Worker把同一个任务抢走”的经典问题。
ZooKeeper解决这个场景的方式很优雅:每个分片在领取前先尝试创建一个临时顺序节点,哪个Worker创建的序号最小,哪个Worker就获得这个分片的处理权。
2.2 临时顺序节点的抢锁流程
先用一段示意代码讲原理。假设每个分片对应一个锁路径,比如/etl/jobs/ods_orders_inc/shards/0/lock,Worker想处理分片0时,在这个路径下创建一个EPHEMERAL_SEQUENTIAL节点:
String participantPath = "/etl/jobs/ods_orders_inc/shards/0/woker-"; String createdPath = zookeeper.create(participantPath, null, ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL_SEQUENTIAL); // 例如返回 /etl/jobs/ods_orders_inc/shards/0/woker-0000000012 List<String> children = zookeeper.getChildren("/etl/jobs/ods_orders_inc/shards/0", false); Collections.sort(children); // 判断自己创建的节点是不是所有子节点中序号最小的如果自己的序号是最小的,说明锁拿到了。如果序号不是最小,就监听比它序号小一个的那个节点,等待前一个节点被删除(也就是前一个Worker处理完或宕机),然后再次尝试。
这里最关键的细节是:因为节点是临时节点,持有锁的Worker如果突然宕机,ZooKeeper会话超时后会自动删除这个节点,锁立刻释放。这彻底解决了“持锁进程死了但锁没人释放”的问题,比Redis锁需要手工设置过期时间要可靠得多。
2.3 为什么我不建议自己写锁实现
上面这段代码只是演示原理。真实生产环境里,我强烈建议直接使用Apache Curator封装好的InterProcessMutex,而不是自己维护锁逻辑。
原因是,原生实现有几个非常容易出错的边界状态:
- 创建节点成功了,但还没走到
getChildren这一步,客户端和ZooKeeper之间的连接恰好断开了。 ConnectionLoss异常发生时,你无法确认刚才那个临时节点到底创建成功了没有。- 监听前一个节点时,如果前一个节点在你注册监听前已经删除,你可能会一直等待,直到超时。
这些都是分布式协调里最常见的“未知状态”问题。自己处理一两个还能应付,全套做完的成本不亚于重写一个协调库。用Curator的写法就简单很多:
InterProcessMutex lock = new InterProcessMutex(client, "/etl/locks/ods_orders_inc_shard_0"); if (lock.acquire(30, TimeUnit.SECONDS)) { try { // 真正执行分片抽取逻辑 } finally { lock.release(); } }2.4 Watcher挂载时的两个工程坑
用ZooKeeper做任务抢锁,有两个工程问题特别容易踩。
第一个坑是Watcher一次性的。ZooKeeper的Watcher在触发一次事件后就会失效,如果代码里没有在回调中重新注册Watcher,下一次节点变化时客户端就收不到通知了。尤其在任务分片很多、锁竞争很频繁的场景里,一次事件丢失可能导致一批Worker都在等待超时,直到靠锁的acquire超时机制兜底才恢复。用Curator的NodeCache或PathChildrenCache可以自动重新注册,减少这类低级问题。
第二个坑是“惊群效应”。如果每个分片被多个Worker同时抢,最朴素的方案是所有抢锁失败者都监听同一个父节点,只要锁一释放,所有等待的Worker都被唤醒。但这样ZooKeeper会在同一时刻向大量客户端推送事件,节点越多,压力越大。正确的做法是每个等待者只监听它前一个节点,形成一条通知链,这也是InterProcessMutex底层的默认做法。
注意:在使用ZooKeeper做分片锁时,一定不要把锁的持有时间和业务执行时间混为一谈。业务执行时间可能因为下游数据库慢查询而无限拉长,而临时节点只和ZooKeeper会话绑定,一旦网络波动导致会话超时,锁释放了,但业务线程还在执行,此时另一个Worker可能会拿到锁跑同一个分片。所以ETL任务本身要设计成可重复执行的幂等逻辑,至少在写入阶段加上主键冲突更新或去重。
3. 主调度器选主与路由切换:双活之后没再出现过双跑
3.1 从双主事故到Leader选举
回到开头说的那次事故。后来我复盘时发现,问题的本质是两台调度器之间无法确认“谁是真正的当前主节点”。
当时用的是数据库心跳方案:主调度器每5秒更新一次心跳表,备调度器发现心跳超过15秒没更新,就自动接管成为新主。结果主调度器只是因为GC停顿导致心跳写入延迟,并没有真正宕机,备调度器却已经接管了任务,两台机器同时派单,下游就重复了。
把主节点的身份放到ZooKeeper之后,这个问题就从根上解决了。做法不复杂:每台调度器启动时都在同一个路径/etl/leader/active下创建一个临时节点,谁能成功创建,谁就是当前主调度器;其余调度器监听这个节点,一旦节点消失(主调度器宕机或会话超时),它们立刻尝试重新创建节点。
为了避免多台备机同时抢建同一个节点引发的竞争,更稳定的方式是让所有调度器在/etl/leader下创建EPHEMERAL_SEQUENTIAL节点,序号最小的为主,其余节点监听前一个节点。Curator的LeaderLatch直接封装了这层逻辑:
LeaderLatch leaderLatch = new LeaderLatch(client, "/etl/leader", instanceId); leaderLatch.start(); leaderLatch.await(); // 阻塞直到当前实例获得领导权 if (leaderLatch.hasLeadership()) { // 本实例成为主调度器 startScheduler(); }这套机制最大的价值在于:主节点失联后,ZooKeeper会话超时会触发临时节点删除,备用调度器通过Watcher感知到删除事件后马上发起选主,整个过程通常是秒级的,不再依赖数据库5秒一轮的心跳轮询。
3.2 “当前模型调度”是怎么运转的
把Master选出来只是第一步。真正让我觉得ZooKeeper好用的,是它能让“当前模型调度”这套机制跑起来。
什么叫“当前模型调度”?我举一个实际例子。我们的同步平台里,一个ETL数据模型会对应一组任务规则,包括从哪张源表读取、按什么字段切分、每批拉多少行、落到哪个目标表、任务之间的上游依赖是什么。这个规则本身不是静态的,业务方可能随时调整同步模型,比如新增一个关联字段,或者把全量同步改成增量同步。
模型有变化时,如果靠人工去每台调度器上改配置,很容易漏改、改错,而且无法保证所有调度器在同一时刻切换到同一套规则。
我用ZooKeeper把它变成了这样一个结构:
/etl/models/order_sync/config # 当前模型配置快照 /etl/models/order_sync/schedulers # 由Master生成的多个调度器注册信息 /etl/models/order_sync/pending # 待执行分片队列 /etl/models/order_sync/workers # 当前正在处理该模型的Worker实例Master调度器启动后,会读取模型配置,并根据配置生成多个调度器实例,这些调度器实例负责监控数据源的最新状态和待处理分片。每个调度器实例会把自己的调度状态写到对应路径下,例如它当前盯到哪一批数据、上一次派发的时间、下一次计划派发的时间。其他组件通过查看这些节点,就能知道当前模型下有没有调度器存活、调度器有没有正常在推进。
如果某个调度器实例所在的节点异常退出,ZooKeeper上对应的临时节点消失,Master会重新生成一个新的调度器实例来接管。整个过程中,用户只会观察到任务队列的推进发生了一次短暂停顿,不会出现两个调度器同时对一个任务派单的情况。
3.3 网络分区和“脑裂”是怎么避免的
在主备切换场景里,“脑裂”是两个节点同时认为自己是主的最危险情况。我在做数据库心跳方案时最担心的就是这个问题:主调度器所在的网络分区和备调度器所在的网络分区完全隔离,两边都无法连到对方的数据库,于是两边都认为自己才是活着的那个,双双开始派发任务。
ZooKeeper解决这个问题的核心是多数派机制。一个ZooKeeper集群正常对外提供服务需要多数派节点存活。当调度器集群发生网络分区时,只有能和ZooKeeper多数派节点建立会话的那一侧,才能成功创建临时节点;另一侧的调度器即使认为自己还活着,也无法更新ZooKeeper中的主节点状态。也就是说,无论网络上分成几个区,在同一个ZooKeeper集群里能持有主节点锁的,永远只有一边。
不过这里要提醒一句:如果网络分区发生得非常快,旧主调度器持有的临时节点可能还没有立刻消失,新主已经通过选主逻辑产生,这时旧主上的部分任务可能仍在继续执行。所以依赖ZooKeeper选主能让“双主”的概率降到极低,但不能完全消灭极短时间窗口内的重复派单。真正兜底的,仍然是每个ETL任务自身的幂等性设计。主从切换期间宁可让任务重复执行一次,也不能让任务被丢。
3.4 切换路由的粒度不需要太细
有人可能会把ZooKeeper理解成一个“万能任务路由器”,所有任务的执行路径都要经过它动态路由。实际情况不是这样。
在调度平台里,我的实践原则是:路由信息只放“不经常变化但需要全局一致”的数据。比如当前主调度器是谁、当前模型使用哪一版同步配置、当前模型由哪组Worker实例处理。而具体的、频繁变化的数据分片列表,不应该全部塞到ZooKeeper上,否则节点数量和写入压力都会失控。
换句话说,ZooKeeper负责的是“路由规则的共识”,而不是“每条数据的搬运路径”。数据分片列表我一般放在消息队列或任务DB里,ZooKeeper里只保留任务队列的游标指针和状态。
4. 节流阀、检查点和全局配置:用ZK给ETL装一个“控制面”
4.1 动态调整全局抽取速率
ETL任务跑起来以后,另一个绕不开的问题是“限流”。我们平台同时有几十个同步任务在跑,如果所有Worker都用最大速率去源端拉数据,业务高峰期的数据库CPU会立刻被打满;但如果在凌晨低谷期也限制得很死,数据延迟又会越积越严重。
最早我们让运维逐台机器改Worker的配置文件,然后重启进程。这种做法效率极低,一台一台重启的间隙里,有的任务在跑,有的任务停了,进度完全不齐。
后来我把“全局流量阈值”放到了ZooKeeper的一个持久节点里,路径设计为/etl/config/throttle/global_rate,节点值存一个JSON字符串:
{ "maxRecordsPerSecond": 10000, "windowSeconds": 10, "updatedBy": "ops", "updatedAt": "2024-01-15 23:00:00" }数据平台管理员需要调整全局速率时,只需要在管控端修改这个节点的值。每个Worker进程内部启动一个后台线程,用Curator的NodeCache监听这个节点,一旦发现值变化,就把最新的maxRecordsPerSecond刷进本地的一个AtomicLong。业务线程每次从源端拉数据前,检查一下本地限流令牌是否允许继续拉取。
这套方案让我特别满意的是,调整限流参数完全不需要重启任何进程,几秒内整个集群的Worker都会感知到变化。而且ZooKeeper只负责“通知值变了”,实际的速率控制逻辑全部在Worker本地执行,不会因为高频拉取而把ZooKeeper的写请求打满。
提示:不要让业务线程直接去读ZooKeeper节点来决定“这1秒能不能拉数据”。ZooKeeper的Watcher事件推送是有延迟的,不适合当成毫秒级的实时配置中心来用。正确做法是Watcher只负责更新本地缓存,业务线程永远只读本地缓存。
4.2 控制ETL任务并发数的分布式信号量
除了限流,ETL平台还经常需要控制“同一时刻最多只允许N个同步任务在跑”。比如下游Hive表同时写入的分区不能太多,否则HDFS的写入压力会很大。
ZooKeeper的分布式信号量可以很好地做这件事。为每个任务在/etl/semaphore/job下注册一个临时顺序节点,同时约定最多允许5个节点存在,一旦节点数达到5,新的任务就只能等待前面的节点释放。
我自己用Curator的InterProcessSemaphoreV2封装了这套逻辑:
InterProcessSemaphoreV2 semaphore = new InterProcessSemaphoreV2(client, "/etl/semaphore/job", 5); Lease lease = semaphore.acquire(); try { // 执行同步任务 } finally { semaphore.returnLease(lease); }这个场景和第一部分的分布式锁不太一样:分布式锁追求的是“某个资源同一时刻只能被一个任务持有”,而信号量允许一个资源被最多N个任务并发持有。如果直接在ZooKeeper里写一个计数器然后靠数据库更新,多个Worker之间的可见性很难保证,用临时节点作为租约则简单且自动容错。
4.3 增量任务的检查点(Checkpoint)保存
增量同步任务最怕的是什么?是跑到一半挂了,重启后不知道上次同步到哪个位置。
以前我们用一张DB表保存每个任务的同步水位,比如source_table的最后同步主键ID或者binlog位置。这样做的问题是:如果任务所在的Worker因为网络分区失联,ZooKeeper可能已经把该Worker的临时节点清理掉了,但任务进程本身还在运行,它最后提交的检查点可能还没写入DB。新Worker接管后读到的是旧检查点,就会从更早的位置重新同步一遍,产生重复数据。
后来我把“每个任务的检查点”也挪到了ZooKeeper路径下:
/etl/checkpoints/ods_orders_inc/last_offset节点值存任务最近一次确实处理完的位点。每次任务完成一批数据处理后,才更新这个节点。因为是持久节点,即使当前Worker挂了,只要ZooKeeper集群还活着,新Worker就能读到上一次提交的真实进度。
这个方案比直接用DB保存检查点好在哪?主要是ZooKeeper节点天然带版本号。当两个任务实例因为网络原因同时尝试更新同一个检查点时,版本号冲突会让其中一方更新失败,从而避免检查点被旧进度覆盖。如果想要完全可靠的“只提交不丢”,还可以开启ZooKeeper的sync操作,在提交后强制同步到Leader,但这会牺牲一点延迟,通常只在任务进度非常重要的时候才用。
4.4 哪些数据千万别往ZooKeeper里写
和我一样一开始接触ZooKeeper的人,很容易犯一个错误:什么东西都觉得可以往ZooKeeper里放,最后ZooKeeper变成了一个“慢速数据库”,写入性能被拖垮。
ZooKeeper每个节点的默认大小上限是1MB,而且每次写请求都要走Leader并通过ZAB协议达成多数派确认,写入性能远比不上普通数据库。我见过有人把同步任务的明细日志、每次心跳的详细负载、甚至Kafka消费位点的每个批次都写进ZooKeeper,结果集群经常出现事务日志堆积和节点响应变慢。
一条简单经验:ZooKeeper里只保存“最终状态”和“变更事件”,不保存“过程明细”。检查点可以放,因为每个任务只有一条进度记录;每次心跳的精确耗时、每个分片处理了多少行这种细粒度数据,不应该放。
5. Worker实例注册与感知:ETL执行器如何被调度器发现和剔除
5.1 启动自注册,宕机自动摘除
ETL平台把任务分给Worker时,调度器必须知道当前有哪些Worker可用。如果Worker列表靠人工维护在配置文件里,那么新增一台Worker就得改所有调度器的配置,Worker宕机也要等人工去改,非常不灵活。
用ZooKeeper做Worker实例注册是标准解法。每台Worker启动时,在/etl/workers/路径下创建一个临时节点,节点名带上实例ID,节点值写上该实例的主机地址、端口和当前负载信息:
public void registerWorker(String workerId, String address) throws Exception { String path = "/etl/workers/worker-" + workerId; zookeeper.create(path, address.getBytes(StandardCharsets.UTF_8), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL); }因为这是临时节点,Worker正常运行时节点一直在;Worker宕机或网络断开超过会话超时时间后,节点被自动删除。调度器通过监听/etl/workers的子节点变化,就能实时感知Worker的上下线。
这套机制在扩容时特别省心。新Worker只要在启动时配置好ZooKeeper地址,启动流程里创建一个临时节点,调度器立刻能看到它,不需要改任何已有节点的配置。
5.2 Worker的负载信息怎么上报
临时节点解决的是“存活状态”问题,但调度器做任务分配时,还需要知道每个Worker当前空闲不空闲。如果每个Worker都把自己当前的任务数写在一个固定节点上,那么调度器只要读这个子节点列表,就能把新任务优先分给负载最低的Worker。
实现上,可以让每个Worker在自注册节点的数据里写入当前的负载信息,并在任务开始和结束时更新一次。节点值类似:
{ "host": "10.10.0.31", "runningTasks": 3, "maxConcurrentTasks": 10, "lastHeartbeat": "2024-06-01 12:00:00" }调度器在分配任务时,先拉取/etl/workers下所有子节点的值,过滤掉runningTasks >= maxConcurrentTasks的Worker,再从剩余实例里选择负载最低的。
需要注意,更新的频率要控制住。如果Worker每处理完一小批数据就更新一次节点,写入频率可能过高。我的做法是任务开始和结束时各更新一次即可,调度器允许短时间内的负载信息滞后,只要不把任务分配给一个已经满载的Worker就行。
5.3 和Kubernetes的服务发现有什么关系
有人会问:既然已经有了Kubernetes,它的Service和Endpoint也能做服务发现,为什么还要用ZooKeeper?
我的理解是,两者解决的是不同级别的问题。Kubernetes解决的是“容器实例的IP和端口发现”,它保证的是网络层的连通性。而ZooKeeper在ETL场景里额外提供了强一致性的“分布式共识”能力,比如抢锁、选主、信号量,这些是Kubernetes原生Service机制不直接提供的能力。
对于纯容器化部署且没有复杂分布式协调需求的ETL工具,直接用Kubernetes的服务发现就够了;但对于需要精细控制任务不重复执行、支持多活调度器切换的场景,ZooKeeper仍然是更顺手的选择。
6. ZK集群部署与参数配置:在ETL规模下应注意什么
6.1 集群规模和磁盘规划
我见过不少团队在测试环境用单机ZooKeeper,这没问题,但生产环境一定要至少3台组成集群,否则一旦单点宕机,整个调度控制面就瘫痪了。
ZooKeeper节点数量最好用奇数,3台或5台最常见。ETL调度平台场景下,除非任务量和Worker数量特别巨大,否则3台已经足够。5台带来的额外成本不是机器本身,而是ZAB协议在写请求上需要达成多数派确认,节点越多,写延迟略高。
磁盘方面,ZooKeeper对事务日志的写入延迟非常敏感。它每个写请求都要落事务日志,如果事务日志磁盘和HDFS的数据盘共用一块普通机械硬盘,大量监听事件和节点变更会让IO延迟飙升,进而影响整个集群的会话稳定性。有条件的话给ZooKeeper单独挂一块SSD,哪怕容量不大也行。
6.2 关键参数设置
我经常碰到的几个生产参数,整理成下面这张表,方便直接抄作业:
| 参数 | 配置位置 | 典型值 | 说明 |
|---|---|---|---|
tickTime | zoo.cfg | 2000 | ZooKeeper基础时间单位,单位毫秒 |
initLimit | zoo.cfg | 10 | 初始化连接时,允许的tick数量 |
syncLimit | zoo.cfg | 5 | Leader与Follower之间心跳超时的tick数量 |
sessionTimeout | 客户端代码 | 10000~30000 | 客户端会话超时时间 |
autopurge.purgeInterval | zoo.cfg | 1 | 事务日志/快照自动清理间隔(小时) |
autopurge.snapRetainCount | zoo.cfg | 3 | 保留的快照数量 |
在ETL平台里,sessionTimeout这个参数的设计最重要。如果设置得太短,比如3秒,那么一次Java Full GC造成几秒停顿就可能让所有Worker的临时节点全部失效,调度器认为Worker全部宕机,引发大规模任务重跑。如果设置得太长,比如60秒,那么Worker真正宕机后,调度器要等60秒才能感知,任务恢复时间被拉长。
我个人的习惯是业务端设置30秒,同时配合监控Worker的JVM GC情况。如果要追求更敏感的故障发现,可以缩短到15秒,但需要保证任务线程不会出现超过15秒的长时间停顿。
6.3 连接池和客户端的统一管理
不要让每个业务任务都自己创建ZooKeeper连接。ZooKeeper客户端和集群之间需要维持长期会话,每新建一个连接都会增加一次握手开销,而且会话数量太多会占用服务端内存。正确做法是在平台进程里维护一个全局的ZooKeeper连接单例,所有模块都复用这个连接。
Curator自带的CuratorFramework本身就管理了连接重试和会话恢复逻辑,建议直接用,不要用原生ZooKeeper客户端裸写业务。当然,业务代码中要显式捕获ConnectionLossException、SessionExpiredException这类异常,连接丢失后Curator会自动重连,但会话过期后Curator可能无法恢复,需要重新创建客户端。
7. 真实故障复盘:三件让我改变ZK使用习惯的事
7.1 故障一:Watcher事件丢失导致任务队列“堵死”
有段时间,我们线上一个分片同步任务偶尔会挂在“等待执行”状态很久,查看ZooKeeper节点,发现锁节点已经被删除了,但Worker一直没有去抢下一个分片。
排查后定位到问题:这些Worker用的是原生ZooKeeper API,在Watcher回调里收到锁释放事件后,忘记重新注册对下一个节点的监听。事件是一次性的,错过以后就再也没有下一次通知了,任务只能靠超时重试机制被拖起来。
这次之后,我把所有和监听相关的代码都换成了Curator的Cache组件,并且代码审查时规定:禁止直接使用原生getChildren+exists+Watcher的组合来做事件感知,除非能明确证明不会丢失事件。
7.2 故障二:Node进程活着但线程池卡死,任务无法被接管
另一个印象深刻的案例:某台Worker处理一个超大分片时,业务线程池全部阻塞在了一个外部API调用上,进程本身没死,JVM也没发生异常退出。但从任务角度看,这个Worker已经在超过30秒的时间里没有任何推进。
此时ZooKeeper里注册的临时节点仍然存在,因为会话还活着。调度器认为这台Worker是健康的,不会重新分配它的任务。结果就是被卡住的任务一直卡到外部调用超时,才被自身的重试逻辑救回来。
那次之后我才明白,ZooKeeper的临时节点只能检测进程级的存活,检测不了业务级的卡死。如果你需要调度平台能在Worker“假死”时自动把任务转移,就必须额外在任务节点上维护一个“任务心跳”,任务每处理一小批数据就更新一次lastProgress时间。调度器在分配任务时不仅要看Worker实例是否注册,还要看它当前执行的任务是否有进度。
7.3 故障三:写入频率过高导致事务日志疯长
还有一次是ZooKeeper集群磁盘报警。检查后发现,是某个新上线的同步任务在检查点保存时,每个分片每处理完1000条数据就更新一次ZooKeeper节点。任务的实时性要求高,所以更新频率也高,几十个分片一起跑,ZooKeeper每秒钟要处理上千次写请求。
虽然单次写请求的数据量很小,但ZooKeeper每次写都要落事务日志,高频写请求很快就把磁盘IO打满了。后来我把检查点更新策略改成了“每处理完一个分片才更新一次”,并且用本地内存暂存分片内的中间进度,写入频率一下就降了几百倍,ZooKeeper集群的压力完全恢复。
这次故障让我给自己定了一条规矩:任何写入ZooKeeper的操作,先估算最大QPS,超过每秒几十次就要考虑设计是不是有问题。
接触ZooKeeper几年下来,我自己最大的体会是,它并不是一个“你所有分布式问题都能靠它解决”的银弹。它的定位更像是ETL工具集里的“共识开关”,放对了地方,能让复杂的任务调度变得有条理;放错了地方,比如拿来当任务明细存储,只会引入新的瓶颈。
如果你正在做一个自研的ETL调度平台,或者你维护的同步工具在多活切换、任务防重上反复踩坑,可以优先把这几件事交给ZooKeeper:任务分片锁、主调度器选主、Worker实例注册、全局限流配置、增量检查点。在这些场景里,它提供的临时节点、Watcher、顺序节点机制,刚好补上了关系型数据库和Redis解决不了的“分布式共识”这一环。
最后再分享一个经验:ZooKeeper集群不只是搭起来就完事的,它在生产环境中的稳定性,很大程度上取决于你对客户端连接、会话超时、写入频率和监听事件的管理。每次新加一种使用方式前,我会先问自己三个问题:这个数据需要全局一致吗?需要多台机器一起感知变化吗?单点故障会造成什么后果?如果前两个答案是“是”而第三个答案是“严重后果”,那就该用ZooKeeper;否则,还是让它待在控制面里当一个安静的协调者最稳妥。