Apache Storm 守护进程容错机制全解析:Worker 重启、Nimbus 故障恢复与集群高可用实践
Apache Storm 守护进程容错机制全解析:Worker 重启、Nimbus 故障恢复与集群高可用实践
本篇文章围绕 Apache Storm 的守护进程容错(Daemon Fault Tolerance)展开,系统讲解 Nimbus、Supervisor、Logviewer、UI 等守护进程在 Worker 死亡、节点宕机、自身崩溃三种故障场景下的恢复流程,并结合本仓库中 storm-core 的源码实现与 conf/defaults.yaml 配置项,说明容错背后的心跳检测、无状态重启与任务重调度机制。读完本文,你将掌握 Storm 集群守护进程的故障处理原理、关键超时参数的调优方法,以及"Nimbus 是否为单点故障(SPOF)"这一经典问题的确切答案。
一、Storm 的守护进程全景
在深入故障处理之前,先厘清 Storm 集群中几类核心守护进程的角色分工。正如 Daemon-Fault-Tolerance.md 所述,Storm 的守护进程包括:
- Nimbus:负责调度 Worker(把任务分配到各机器上),是集群的"主控"角色。对应源码为 nimbus.clj。
- Supervisor:负责在本机启动和杀死 Worker 进程。对应源码为 supervisor.clj。
- Logviewer:提供日志访问能力。对应源码为 logviewer.clj。
- UI:提供浏览器访问的集群与拓扑状态诊断页面。
此外还有 DRPC、Acker 等角色(见 drpc.clj、acker.clj),它们共同构成集群的运行骨架。集群的整体结构与守护进程之间的协作关系可以用下图概括:
从图中可以看到,Zookeeper 承担 Nimbus 与 Supervisor 之间的协调职责(Storm 不通过 Zookeeper 传递消息数据,因此 ZK 负载很低)。这正是后续容错机制得以成立的基石:所有关键状态都存在 Zookeeper 或本地磁盘上,而不是驻留在守护进程的内存里。
二、Worker 死亡时会发生什么
Worker 进程是真正执行 Spout/Bolt 任务的进程。当某个 Worker 死亡时,容错链条分为两层:
-
Supervisor 立即重启 Worker:每个 Supervisor 守护进程在本地持续监控 Worker 进程的运行状态。一旦发现 Worker 进程消失(进程退出、被 kill 等),Supervisor 会重新拉起该 Worker。从源码结构看,supervisor.clj 中有独立的线程负责"启动事件以在必要时重启死掉的进程"(注释原文为 another thread launches events to restart any dead processes if necessary),并提供了
shutdown-worker、launch-worker等关键函数,分别负责回收与启动 Worker。 -
启动持续失败时由 Nimbus 重调度:如果 Worker 在启动阶段持续失败,无法向 Nimbus 发送心跳(heartbeat),Nimbus 就会判定该 Worker 已失联,并将其包含的任务重新调度到其他机器上。
这里的心跳机制值得展开。在 nimbus.clj 中,Nimbus 维护着一个 :heartbeats-cache(心跳缓存),通过 update-all-heartbeats! 与 update-heartbeats! 函数在每次监控周期内刷新各执行器(executor)的心跳,并通过 update-heartbeat-cache 结合超时阈值判断哪些执行器已超时(is-timed-out)。注释中还特别说明:Storm 不假设各机器时钟同步——心跳只用于让 Nimbus 知道"收到了新心跳",所有计时判断均由 Nimbus 基于心跳缓存完成。
2.1 与 Worker 容错相关的关键配置
以下参数定义于 Config.java,默认值可在 conf/defaults.yaml 中核对:
| 配置项 | 默认值 | 作用 |
|---|---|---|
supervisor.worker.timeout.secs |
5 | Supervisor 判定 Worker 心跳超时的秒数,超过则视为该 Worker 死亡并重启 |
supervisor.worker.start.timeout.secs |
120 | 等待 Worker 启动(完成初始化并上报心跳)的最长时间 |
supervisor.worker.shutdown.sleep.secs |
1 | 关闭 Worker 时的等待时间 |
nimbus.task.timeout.secs |
30 | Nimbus 判定执行器任务心跳超时的秒数,超时即触发任务重调度 |
nimbus.task.launch.secs |
120 | 等待任务启动完成的窗口 |
nimbus.monitor.freq.secs |
10 | Nimbus 心跳监控与调度检查的执行频率 |
nimbus.supervisor.timeout.secs |
60 | Nimbus 判定 Supervisor 心跳超时的秒数 |
这些参数共同决定了故障检测的敏感度:调小超时可以让故障更快被发现和恢复,但过小可能在正常负载波动时造成误判;生产环境需要结合机器性能与网络状况权衡。
三、节点(机器)死亡时会发生什么
当一个物理节点宕机时,该机器上所有 Worker 进程都会随之停止。此时:
- 分配给该机器的所有任务(tasks)会进入心跳超时状态;
- Nimbus 在监控周期内检测到这些执行器心跳超时后,会将这些任务重新分配给其他存活的机器,从而保证拓扑继续运行。
这一过程正是"无状态 + 集中调度"设计的结果:Nimbus 掌握全部拓扑的任务分配信息(assignment),节点故障只是触发一次重新调度(reassign)。需要注意的是,nimbus.reassign(默认 true,见 conf/defaults.yaml)是重调度功能的开关,生产环境不应关闭它。
四、Nimbus 或 Supervisor 守护进程死亡时会发生什么
这是 Storm 容错设计中最具特色的部分。Nimbus 与 Supervisor 被刻意设计为:
- fail-fast(快速失败):进程在遇到任何非预期状况时直接自杀退出,而不是尝试在异常状态下继续运行;
- 无状态(stateless):所有状态都保存在 Zookeeper 或本地磁盘(
storm.local.dir)中,进程内存中不留持久性状态。
因此,当 Nimbus 或 Supervisor 死亡时,只要通过监督工具将其重启,它们就能"像什么都没发生过一样"恢复工作。这正是 Setting-up-a-Storm-cluster.md 反复强调"必须以监督方式运行守护进程"的原因:
- 在 master 机器上以
bin/storm nimbus启动 Nimbus; - 在每台 worker 机器上以
bin/storm supervisor启动 Supervisor; - 以
bin/storm ui启动 UI(浏览器访问http://{ui host}:8080)。
推荐的监督工具包括 daemontools、monit 等进程守护工具(Zookeeper 也应同样置于监督之下,因为 ZK 同样是 fail-fast 设计)。由于 fail-fast 与无状态,守护进程可以"安全地在任意时刻停机,并在重启后正确恢复"——这与 Hadoop 形成鲜明对比:Hadoop 中若 JobTracker 死亡,所有正在运行的作业都会丢失;而 Storm 中 Nimbus 或 Supervisor 死亡不会影响任何正在运行的 Worker 进程。
从源码看,fail-fast 的定位在 supervisor.clj 中也有体现:Supervisor 对 Worker 进程的启停(shutdown-worker、launch-worker、kill-process-with-sig-term 等)都通过独立的监控线程驱动,自身崩溃后由外部监督工具拉起即可继续履行这些职责。
五、Nimbus 是单点故障(SPOF)吗
原文给出的结论是"某种程度上是,但实际影响不大":
- 失去 Nimbus 节点后,正在运行的 Worker 依然正常工作;
- Supervisor 依然会在 Worker 死亡时就地重启 Worker;
- 但是,没有 Nimbus 时,Worker 不会被重新分配到其他机器(例如丢失一台 Worker 机器时,其任务无法迁移)。
也就是说,Nimbus 死亡不会引发灾难性后果,它只是暂时丧失了"调度"能力。因此实践中无需为 Nimbus 短暂下线过度紧张。需要说明的是:本仓库版本中 Nimbus 尚无内置的高可用(HA)实现,文档明确标注"未来有计划让 Nimbus 高可用"。若要在多台机器上部署候选 master,则需在 nimbus.seeds 中列出所有运行 Nimbus 的机器 FQDN(详见 Setting-up-a-Storm-cluster.md 中的 nimbus.seeds 配置说明),Worker 通过该配置获取拓扑 jar 与配置的下载来源。
六、节点与消息丢失时,Storm 如何保证数据处理
守护进程容错只是集群层面的保障,数据层面的保障由 Storm 的可靠性机制(reliability API)提供。关于"节点死亡或消息丢失时如何保证数据处理",Guaranteeing-message-processing.md 给出了完整说明,其核心思想如下:
-
tuple tree(元组树)追踪:Spout 发出的一个元组会触发下游成千上万个元组,形成一棵"元组树"。Storm 只有当整棵树的元组全部处理完成,才认为该 Spout 元组"完全处理(fully processed)";如果在配置的超时时间内未能完全处理,则该元组判定失败。
-
超时可配置:超时时间由
Config.TOPOLOGY_MESSAGE_TIMEOUT_SECS(即topology.message.timeout.secs,默认 30 秒,见 Config.java)控制。 -
ack / fail 回调:完全处理时,Storm 调用 Spout 的
ack(msgId);超时失败时调用fail(msgId),且一定由创建该元组的同一个 Spout task 完成回调。Spout(如 KestrelSpout)收到fail后可把消息重新放回队列实现重放。 -
anchoring(锚定)与 ack/fail API:Bolt 通过
collector.emit(tuple, new Values(...))锚定新元组,并在处理完输入元组后调用collector.ack(tuple);未锚定的元组失败不会触发 Spout 元组重放。BaseBasicBolt自动完成锚定与 ack,适合过滤/简单函数类 Bolt。 -
Acker 任务的异或追踪算法:拓扑中有一组专门的 acker 任务,通过 64 位随机元组 id 的 XOR(异或) 累加值(ack val)跟踪整棵树的状态;当 ack val 归零即代表树处理完成。该算法每个 Spout 元组仅需约 20 字节固定内存,即使元组树有数万个节点也不会撑爆 acker 内存。ack val 偶然归零的概率极低——按每秒 10K 次 ack 计算,约 5000 万年才可能出现一次误判。
-
失败场景覆盖:
- 任务(task)死亡导致元组未 ack:根元组超时,从 Spout 重放;
- Acker 任务死亡:其跟踪的所有 Spout 元组超时并重放;
- Spout 任务死亡:由 Spout 对接的消息源负责重放(如 Kestrel、RabbitMQ 会在客户端断开时把 pending 消息重新放回队列)。
-
关闭可靠性的三种方式:设置
topology.ackers为 0(元组发射后立即 ack);发射元组时省略 message id(按消息关闭追踪);以未锚定方式发射下游元组(按子树关闭追踪)。
以上机制与守护进程容错相互配合,构成 Storm"节点可死、数据不丢(至少一次语义)"的完整容错体系。若需要精确一次(exactly once)语义,则应使用 Trident API。
七、守护进程容错的运维实践
结合上文,在真实集群中落实守护进程容错,需要注意以下几点:
-
所有守护进程必须置于监督之下运行:Nimbus、Supervisor、UI 乃至 Zookeeper 都要用 daemontools、monit 等工具托管,使其"死亡即被拉起"。同时 Zookeeper 建议配置 cron 定期压缩数据与事务日志,避免磁盘耗尽。
-
按需调优超时参数:故障检测的灵敏度取决于
nimbus.monitor.freq.secs、nimbus.task.timeout.secs、nimbus.supervisor.timeout.secs、supervisor.worker.timeout.secs等配置。调优前应观察集群负载,避免误判与频繁重调度。 -
利用节点健康检查脚本:Setting-up-a-Storm-cluster.md 提到,Storm 支持管理员在
storm.health.check.dir(默认healthchecks)目录放置健康检查脚本,脚本若输出以ERROR开头的行则表明节点不健康,Supervisor 会关停本机 Worker 并退出;可通过/bin/storm node-health-check在监督框架中判断是否应启动 Supervisor。脚本需有执行权限,执行超时由storm.health.check.timeout.ms(默认 5000ms)控制。这套机制让"节点级故障"可以在真正宕机前被提前发现并优雅处理。 -
通过 UI 观察集群健康:UI 页面展示拓扑与集群状态,Worker/Acker 的存活与吞吐情况可以在浏览器中直接观测,是验证容错配置是否生效的第一现场。
八、总结
Storm 的守护进程容错建立在三条设计原则之上:fail-fast 进程模型、无状态(状态全部外置到 Zookeeper 与磁盘)、集中式心跳监控与任务重调度。Worker 死亡由 Supervisor 就地重启,启动失败或节点宕机则由 Nimbus 重新调度;Nimbus/Supervisor 自身死亡后由外部监督工具无痕拉起,且不影响任何运行中的 Worker——这是与 Hadoop JobTracker 故障模型最本质的区别。Nimbus 虽有"准 SPOF"之名,但因其死亡不产生灾难性后果,实践中只需做好监督托管即可。若要进一步深挖数据层面的保障,请继续阅读 Guaranteeing-message-processing.md;完整的集群搭建与监督部署步骤见 Setting-up-a-Storm-cluster.md。
