十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

Storm容错机制全解析:节点故障后如何保证数据零丢失

Storm容错机制全解析:节点故障后如何保证数据零丢失 直接从一个真实的夜里说起吧。当时我们线上的 Storm 集群跑着一条实时订单风控流上游 Kafka 里积压着几百万条消息下游 Bolt 做规则匹配和用户画像关联。突然一台 Supervisor 节点因为物理机内存故障宕了监控大屏上一片飘红我当时第一反应是“完了这波数据要丢”。结果等集群完成故障转移再去核对 Kafka 的 offset 和 Storm 实际处理数一条没少。那一刻我才真正意识到Storm 的容错机制不是文档里写的“理论设计”而是真的能在生产环境里扛事的东西。这篇文章不打算复述官方文档而是想从一个长期维护 Storm 集群的人的角度把“节点故障”到“数据零丢失”这条链路完整拆开。你会看到 Nimbus、Supervisor、ZooKeeper 是如何分工的Acker 如何用异或运算追踪每条消息以及当 Worker 进程被杀、网络抖动、任务重新调度时背后到底发生了什么。我也会把我在生产环境踩过的坑、验证过的排查方法一并写出来希望能给正在用 Storm、或者正准备做实时计算选型的朋友一些真正能落地的参考。1. 容错机制的整体架构Storm 凭什么敢说“不丢”1.1 从“进程会挂”这个现实出发在学习 Storm 容错之前需要先接受一个分布式系统的基本前提任何节点在任何时刻都可能挂掉。物理机可能宕机、虚拟机可能被回收、进程可能被 OOM Killer 干掉甚至机房网络可能抖动几秒。容错机制不是用来“避免故障”的而是用来“在故障发生后让系统还能继续正确工作”。Storm 对这个问题的回答分成了两个层面。第一个层面是“集群层面的故障转移”也就是当某台机器挂了之后运行在这台机器上的 Worker 进程怎么办、任务怎么重新分配、上下游怎么感知。第二个层面是“消息层面的可靠性保障”也就是从 Spout 发射的每一条 tuple如何确保它被完整地处理完哪怕处理过程中某个 Bolt 挂了也能通过重放机制补上。这两个层面缺一不可集群恢复得再快如果消息本身会丢那对业务来说依然是灾难。1.2 三个角色一台戏Nimbus、Supervisor、ZooKeeper 各管什么Storm 集群的容错之所以做得优雅核心在于它的架构设计刻意做了“去状态化”。我先用一个生活化的类比来帮你建立直觉——可以把整个集群想象成一家物流公司。Nimbus 是公司的调度中心只负责“派单”。它把拓扑拆分成一个个 task然后决定哪些 task 跑在哪台机器上。它只做这一件事而且做完就忘不保存任何中间状态。这就是所谓的“无状态主节点”也是 Storm 容错的一个关键设计决策。Supervisor 是各分部的现场主管。它监听调度中心的指令在本机启动或停止对应的 Worker 进程。每个 Worker 进程是真正干活的 JVM里面跑着一组 Executor 线程每个 Executor 再负责执行一个或多个 Task 实例。ZooKeeper 则是整家公司的“公告板”。调度中心把任务分配信息写到公告板上现场主管从公告板上读取指令Worker 进程的心跳也往公告板上写。任何一方挂了再起来只要读公告板就能知道自己该干什么。这个设计的精妙之处在于因为所有协调信息都持久化在 ZooKeeper 上所以 Nimbus 挂了集群的 Worker 进程依旧能正常处理数据流。Supervisor 挂了Nimbus 发现心跳超时后会重新调度它负责的 Worker 到其他机器。从头到尾没有一个单点是“非活不可”的。1.3 为什么“无状态”是容错的最优解很多刚接触 Storm 的人会不理解为什么 Nimbus 不直接和 Worker 保持长连接、实时下发指令那样调度延迟不是更低吗但答案是长连接方案在故障面前非常脆弱。一旦 Nimbus 和 Worker 之间的网络出现短暂分区两边都不知道对方是否存活很容易出现“双主”或者“误杀”的混乱局面。Storm 选择了一条更稳妥的路所有状态都丢给 ZooKeeperNimbus 和 Supervisor 只做“状态变更的触发者”。Nimbus 挂了ZooKeeper 里已有的分配信息不会消失Supervisor 继续按旧指令运行 Worker。Nimbus 重新拉起后只需要重新读取 ZooKeeper 中的元数据就能恢复对整个集群的认知然后继续做故障检测和新任务的调度。实际上这背后是一种“外部化状态”的思路。真正的系统状态不在进程内存里而在 ZooKeeper 里。所有组件进程都变成了可以随时杀掉、随时重建的“无状态计算单元”。这也是为什么 Storm 的守护进程可以放心地用storm kill强杀重启后不需要做任何数据恢复操作。2. 节点故障的检测与恢复全流程2.1 心跳机制与超时阈值一台机器宕机后的 15 秒Storm 的故障检测核心是“心跳 超时”。每个 Supervisor 和 Worker 都会定期往 ZooKeeper 上报自己的心跳信息表示“我还活着”。Nimbus 则像一个考勤员定期检查这些心跳。这里有一个生产环境必须注意的参数nimbus.supervisor.timeout.secs默认值是 60 秒。也就是说如果 Nimbus 超过 60 秒没收到某个 Supervisor 的心跳它就判定这台机器挂了。一旦判定Nimbus 会把该 Supervisor 下所有 Worker 的任务重新分配。我在生产环境会把这个值调小到 15 到 30 秒因为我们的业务容忍不了太长故障窗口。不过要提醒一句调得太小也有风险——网络瞬间抖动导致心跳延迟可能引发不必要的任务重排反而造成更多颠簸。2.2 任务重新分配一个 Worker 的重生过程当 Nimbus 判定某台机器故障后具体的恢复流程是这样的首先Nimbus 会读取 ZooKeeper 中存储的“当前拓扑分配信息”找出故障节点上原本运行了哪些 Worker、每个 Worker 负责哪些 Executor。Nimbus 会把这些 Executor 标记为“待分配”状态。接着Nimbus 调用调度算法在剩余健康的 Supervisor 中挑选合适的目标机器。这里的“合适”包括CPU 核数是否充足、内存余量是否够用、是否已经有该拓扑的其他 Worker尽量做到负载均衡。然后Nimbus 更新 ZooKeeper 中的分配信息目标机器的 Supervisor 监听到了变化就会启动新的 Worker 进程。新 Worker 启动后会从 ZooKeeper 里拉取 Topology 的 DAG 结构定义恢复自己的 Executor 执行上下文。最后还需要注意一个细节Storm 的新 Worker 启动后上游的 Worker 并不会自动知道“下游换了地址”。它们是通过 ZooKeeper 中的“任务到主机映射”信息来发现对端的。新 Worker 注册成功后上游会通过网络连接池重建连接通信自动恢复。整个过程对外部数据源来说是透明的。2.3 Nimbus 自己挂了怎么办没人要求 Nimbus 必须永不宕机。默认配置下Nimbus 是单点运行的所以如果 Nimbus 进程所在机器崩溃集群会暂时失去“调度大脑”。但注意仅仅是“暂时失去调度能力”而不是“集群不可用”。因为 Worker 还在继续处理数据拓扑的运行并不依赖 Nimbus 实时参与。Nimbus 挂了数据流不会断只有发生新的故障需要重新调度时才会因为没有 Nimbus 而无法转移任务。恢复方式很简单使用storm命令重新启动 Nimbus 进程。它启动后会自动从 ZooKeeper 中读取集群元数据和拓扑状态恢复对整个集群的管理权。如果担心 Nimbus 单点可以部署 Supervisor 监控脚本在 Nimbus 进程退出时自动重启。我的经验是Nimbus 的可用性其实没那么关键因为它不参与数据通路。你更应该操心的是 ZooKeeper 集群的稳定性。如果 ZooKeeper 整体挂了那 Storm 的所有心跳和元数据读写都会失败整个集群的可靠性会大打折扣。2.4 实操模拟节点故障的验证步骤理论讲再多不如亲手验证一遍。下面是我在测试环境常用的故障演练步骤也是建议所有使用 Storm 的团队在投产前做一遍的操作。在集群里部署一个简单的 WordCount 拓扑设置 3 个 Supervisor 节点每个节点 2 个 Worker。启动拓扑后用storm list确认任务分布记录每个 Worker 所在的主机。选择一个 Worker 运行中的主机先jps查看进程找到对应的 Supervisor 进程 PID。直接kill -9该 Supervisor 进程模拟物理机崩溃因为普通 kill 会被 JVM 捕获可能触发优雅停机不够“暴力”。观察 Nimbus 日志目标是在 15 到 60 秒内看到类似Failed to get nimbus heartbeat或者Supervisor is not alive的日志。继续用storm list和 Storm UI 观察拓扑执行状态会发现 Executor 总数不变但部分 Executor 的 Host 字段变成了其他机器。检查消息处理是否中断以及 Kafka 的 offset 与 Storm 的处理记录是否一致。如果第 7 步发现问题说明你的拓扑没有正确开启消息可靠性保障或者某个 Bolt 在 emit 后没有做 anchor这时就需要回到第三大部分讲到的 ack 机制去排查。3. 数据零丢失的链路保证Spout、Bolt 与 Acker 的三角协奏3.1 “零丢失”到底是什么意思很多初学 Storm 的人会对“数据零丢失”产生误解以为是“exactly-once”。实际上Storm 默认提供的是“at-least-once”语义每条消息至少被处理一次但可能因为故障重放而被处理多次。真正意义上的“幂等去重”需要配合 Trident 或者下游存储自己实现。那为什么我们不直接追求“exactly-once”因为分布式系统的 CAP 限制下“不丢”和“不重”是矛盾的。要保证不丢就必须在有故障时重放消息而重放天然会带来重复。Storm 选择把“不丢”作为底线把“重复”的解决交给业务层。这是一个非常务实的取舍——对大多数实时计算场景风控、监控、流量统计来说数据不能丢偶尔重复则可以接受或者通过业务幂等来化解。3.2 Spout 的 ack/fail 与消息重放机制消息可靠性的起点在 Spout。Spout 是数据流的源头它从 Kafka、RocketMQ、数据库 binlog 或者别的数据源拉取数据然后发射为 tuple。要让 Storm 帮你保证不丢Spout 发射 tuple 后必须记录一个“未完成”状态。Storm 的 BaseRichSpout 中有一个核心方法nextTuple()以及配套的ack(Object msgId)和fail(Object msgId)回调。这里的msgId就是 Spout 用来关联原始消息和处理结果的“凭证”。具体流程是这样的Spout 从 Kafka 拉取一条消息增加到待确认队列将 msgId 和消息内容对应好然后调用collector.emit(tuple, msgId)将 tuple 发送出去。之后如果整条 tuple 树都处理完毕Spout 会收到 ack 回调这时就可以把消息从待确认队列中删除。如果某个环节失败或超时Spout 会收到 fail 回调这时 Spout 需要重新从 Kafka 拉取该条消息再次发射。这里有一个关键点Spout 必须在fail()回调里实现“重新拉取并发送”。Storm 只负责通知你“失败了”重放动作需要你自己写。通常做法是把 msgId 对应的原始消息缓存起来fail 时再次 emit。如果用的是 KafkaSpout框架已经自动实现了这套重放逻辑但如果你自己写 Spout就要特别小心这个环节。3.3 Bolt 的 anchoring 与 tuple 树的形成消息从 Spout 出去之后会在 DAG 里经历多个 Bolt。这里最关键的概念是anchoring——锚定。当一个 Bolt 处理一条输入 tuple 时它可能会派生出一条或多条新的 tuple并通过collector.emit(inputTuple, newTuple)发射出去。这个操作就把新 tuple “锚定”到了原来的 tuple 上Storm 会据此建立起一棵tuple 树。tuple 树的根是 Spout 发射的原始 tuple叶子节点是最终被处理完的输出。只有当整棵树的所有节点都被 ack 了根 tuple 才算处理成功。任何一个节点 fail或者超时未 ack根 tuple 就会整树重放。很多新手在写 Bolt 时最容易犯的错误是在 Bolt 里用collector.emit(new Tuple(...))直接发新 tuple而没有把输入 tuple 作为第一个参数传进去。这样新 tuple 就没有被锚定Storm 无法追踪它和原始 tuple 的关系。结果是即使这个下游 Bolt 处理失败上游 Spout 根本不知道数据照样丢。这也是我每次 code review 时必查的一个点。3.4 Acker 如何用异或运算追踪整棵树接下来是整个容错机制最精妙的部分——Acker。Acker 是一种特殊的 Bolt它不处理业务数据只负责追踪 tuple 树的状态。每个拓扑可以配置一个或多个 Acker数量由topology.acker.executors参数决定。Acker 的数量不是越多越好因为它需要维护大量 tuple 的追踪状态数量太多会带来额外内存消耗通常 1 到 3 个就够高并发场景可按 Worker 数量的四分之一左右估算。Acker 的追踪算法核心是一只“异或计数器”。当 Spout 发射一条 tuple 时它会为这条 tuple 生成一个随机 64 位整数作为初始校验值并将其发送到某个 Acker。Acker 会记录这个 tuple 对应的“预期校验值”。当 Bolt 处理完一条输入 tuple 并 ack 它时Acker 会把该 tuple 的锚定关系做一个异或运算更新。当 Acker 发现某个 tuple 的校验值归零时就表示整棵树处理完毕Acker 通知 Spout 执行 ack 回调。这个异或方案的好处是Acker 不需要保存整棵树的结构只需要保存一个 64 位整数和每条 tuple 的发送方、接收方信息。即使 tuple 树有几十万个节点追踪成本也几乎不变。3.5 超时机制防止“永久悬挂”的最后一道防线Acker 能追踪成功和失败但无法追踪“既没成功也没失败”的悬挂 tuple。比如某个 Bolt 进程被 kill 之前刚收到 tuple 还没来得及处理这个 tuple 就永远不会 ack 了。如果没有超时机制Spout 会一直等下去。Storm 的topology.message.timeout.secs参数就是为这个场景设计的默认值是 30 秒。如果一条 tuple 从发射到 ack 超过这个时间Acker 会判定失败Spout 会收到 fail 回调并触发重放。在生产中这个值需要根据业务处理耗时合理设置。设置的太短会导致慢处理被频繁重放造成重复数据爆炸设置的太长故障恢复的延迟会变大。3.6 一个完整链路示例订单风控中的一次成功 ack为了方便理解用一个纯业务例子从头到尾串一遍。假设我们的 Spout 从 Kafka 读取一条“用户支付事件”消息内容包含用户 ID、支付金额、设备指纹等字段。Spout 调用emit(tuple, msgIdorderId-123)此时 msgId 为orderId-123。第一个 Bolt 是“用户标签关联 Bolt”它读取用户 ID从 Redis 中获取该用户的 VIP 等级然后 emit 一个新 tuple并 anchor 到输入的 tuple。第二个 Bolt 是“规则匹配 Bolt”它根据前面的标签和金额判断这笔订单是否命中风险规则。如果命中写入预警表无论是否命中都会 ack 输入 tuple。当第二个 Bolt ack 时Acker 更新校验值。假设第一个 Bolt 的 emit 和第二个 Bolt 的 ack 都执行正确那么 Acker 最终算出的异或结果为 0Acker 通知 Spout ackorderId-123。Spout 收到 ack 后从待确认队列中删除这条记录并向 Kafka 提交 offset。如果第 3 步的 Bolt 因为内存溢出崩溃来不及 ack超时后 Acker 会通知 Spout fail 这个 msgIdSpout 从 Kafka 重新拉取该消息整个流程重新执行。这样从源头上保证了数据不丢。4. 真正用好的关键配置调优与生产避坑实录4.1 与容错相关的核心配置项容错机制能不能在生产环境发挥作用很大程度上取决于你如何配置。下面是几个最重要的参数以及我自己的调优经验。topology.acker.executorsAcker 的数量。这个值直接决定了系统能够同时追踪多少 tuple 树。如果为 0则完全关闭可靠性追踪Spout 发射的消息不再有 ack/fail 回调但同时系统吞吐会更高。是否设置为 0取决于业务是否需要“不丢”。建议生产环境至少设置为 1。topology.message.timeout.secs单条 tuple 树上所有处理的最长等待时间。需要根据业务链路的实际耗时去设置。我曾经见过有人把 KafkaSpout 的 timeout 设置的比端到端处理时间还短导致大量消息被误判失败反复重放Kafka 的消费位点一直原地打转。合理的做法是先压测出 P99 的端到端耗时再用它乘以 2 到 3 作为 timeout。nimbus.supervisor.timeout.secs前面提到过的节点心跳超时时间默认 60 秒。生产环境建议根据故障恢复窗口需求调小但不要低于 10 秒否则网络抖动就会导致误判。topology.max.spout.pending这个参数控制 Spout 最多允许同时有多少条 tuple 未被 ack。它本质上是一个“滑动窗口”限制数据源发射速率。调大可以提升吞吐但会增加内存占用和重放时的重复量调小则能降低故障恢复时的重复范围。我的习惯是结合下游处理能力逐步上调。topology.sleep.spout.wait.strategy.timeout.msSpout 在没有消息可拉取时的休眠时间。这个参数虽然不直接属于容错但会影响到空转时的 CPU 消耗过短会导致空转频繁过大则影响消费延迟。4.2 生产环境最容易踩的 5 个“丢数据”陷阱我在排查各种“离奇丢数据”问题时发现大多数都不是 Storm 机制本身的问题而是使用姿势的问题。这里列几个高频坑位。第一个坑Spout 没有用 msgId。有些同学写 Spout 时直接调用emit(tuple)不传 msgId。这样 Storm 无法追踪 tuple 的完成状态ack 和 fail 都无从谈起。实际上emit(obj, msgId)的 msgId 参数就是可靠性追踪的前提没有它一切免谈。第二个坑Bolt 里 emit 后忘掉 ack 或 anchor。如果 Bolt 在处理输入 tuple 后调用了collector.emit(newTuple)但没有把输入 tuple 作为锚定参数传入或者处理完毕后忘了调用collector.ack(inputTuple)那么这条 tuple 永远无法完成直到超时并触发重放。这会导致“消息无限重放”对下游造成巨大冲击。第三个坑重放导致下游重复写入。当消息重放时你已经写入到数据库里的结果可能会被重复写入一次。如果下游是 Redis 的 set 操作重复写无影响但如果下游是 MySQL 自增主键的插入操作就会产生脏数据。解决方案是下游存储使用业务唯一键或者配合 Trident 做幂等。第四个坑Kafka offset 提交过早。有些团队用的是自定义 KafkaSpout为了提升吞吐在消息刚 poll 出来就手动提交 offset而不是在 Spout 的 ack 回调后再提交。这样一旦后续处理失败消息已经提交了 offset重放时从已消费位置开始前面的数据就真的丢了。第五个坑忽略 ZooKeeper 的目录权限。ZooKeeper 上 Storm 的根目录如果配置不当可能导致多个环境共用一套 ZooKeeper 时互相干扰。Storm 的启动参数-Dstorm.zookeeper.root可以在不同环境间做隔离建议生产环境务必设置成不同路径。4.3 监控指标怎么判断容错是否生效就算配置全部正确如果不在线上持续监控故障发生的那一刻你可能还是两眼一抹黑。Storm UI 本身就是最好的监控工具它提供了几个关键指标日常巡检要养成看的习惯。第一个是Complete Latency也就是从 Spout 发射 tuple 到整棵树处理完成的平均耗时。如果这个值接近 timeout 设置值说明系统压力很大需要扩容或者优化 Bolt。第二个是Failed 数量。Storm UI 每个 Spout/Bolt 的统计页里都有 Failed 列它表示处理失败的 tuple 数量。如果 Failed 数量持续增长说明有 bug 导致某些 tuple 永远无法成功需要及时排查。第三个是Kafka 消费位点。如果你用的是 KafkaSpout务必在监控面板上画出未消费消息数Lag。Lag 持续上涨说明消费能力跟不上生产速度一旦积压太久消息过期或者重放可能造成灾难性后果。还有一个容易忽略的点Storm 的日志。生产环境一定要配置 log4j 的日志级别和滚动策略尤其要保留 Nimbus 和 Supervisor 的日志。节点故障发生时这些日志是最直接的排查线索。Nimbus 日志中常见的关键词包括Failed to get nimbus heartbeat、Task is not alive、Connection refused。4.4 一次线上“数据丢失”谜案的真实排查过程最后分享一个我经手的真实案例帮助你把前面这些知识串起来。某天凌晨业务方反映实时大屏上的成交金额骤降但 Kafka 的生产端数据是正常的。我先去 Storm UI 上看拓扑运行状态发现 Spout 的 Failed 指标为 0Complete Latency 也正常看起来一切“正常”。但大屏数字确实少了于是我开始怀疑是不是压根就没消费到数据。我马上看 Kafka 消费组的 Lag。结果发现 Lag 在持续增长而且消费组的成员列表里显示当前活跃的消费者数量为 0。也就是说Storm 集群的 KafkaSpout 根本没有在消费。接着看 Nimbus 和 Supervisor 的日志发现两个小时前某台 Supervisor 因为磁盘写满触发了 OOM导致 Worker 进程全部被杀。但诡异的是Storm UI 上显示拓扑还是 RUNNING 状态。后来才明白因为我们当时把topology.acker.executors设成了 0消息可靠性追踪被完全关闭Worker 挂了之后Spout 的 tuple 不再有 ack/fail 回调对外表现就是“一切都正常”实际上数据早就停住了。这个案例给我的教训非常深刻。当你不关心 ack 机制时Storm 依然不会帮你检测“消息没被处理”这件事。只有开启 acker让每条 tuple 都有明确的完成信号才能真正感知到故障。这也是为什么我会在文章里反复强调ack 机制不仅是数据零丢失的保障更是系统“自检”能力的来源。5. 容错之外的扩展思考与 Kafka、Trident 的配合选型5.1 Kafka Storm 的标准容错组合在实际生产里Kafka 几乎成了 Storm 的标配数据源。Kafka 的 offset 机制 Storm 的 ack 机制刚好形成一套完整的“端到端至少一次”保障。这里要特别提一下 KafkaSpout 的容错行为细节。KafkaSpout 在 emit 每条消息时都会把 msgId 设置为“分区 offset”的组合。当 Spout 收到 ack 回调时KafkaSpout 内部会维护每个分区的“已提交 offset”并在合适时机向 Kafka 提交。如果消息处理失败fail 回调会触发该消息的重放。这套机制结合的非常紧密因此我强烈建议除非你有非常特殊的自定义需求否则不要自己写 Kafka Spout直接用官方提供的即可。在调优上KafkaSpout 有一个值得关注的参数max.poll.records它控制每次从 Kafka poll 出来的最大消息数。这个值如果设置太大会导致 Spout 内存里积压大量未 ack 消息设置太小则可能降低消费吞吐。我的经验是结合topology.max.spout.pending一起调先确定整个拓扑的“在途消息量”再倒推每次 poll 的大小。5.2 什么时候该用 Trident 的 exactly-once如果业务对“重复”完全无法容忍比如金融交易流水、余额变动通知那么即使 Storm 能保证“不丢”at-least-once 的重复问题依然会让下游数据错乱。这个时候就要考虑 Trident。Trident 是 Storm 的高层抽象它通过“事务性状态存储 批次去重”实现了 exactly-once 语义。它把流式数据切分成一个个 batch为每个 batch 分配唯一的 transaction id然后保证每个 batch 要么完全处理成功要么完全失败重放。配合 Redis、HBase 等支持事务的外部存储可以实现“精确一次”的效果。但这里必须说一句实在话Trident 虽好但它的吞吐量比原生的 Storm API 低不少而且在 Debug 时非常痛苦。我的建议是如果业务能通过“业务唯一键 幂等写入”解决重复问题就不要轻易上 Trident只有当你发现幂等操作无法覆盖所有场景时再考虑事务性拓扑。5.3 从 Storm 迁移到 Flink一次架构选型的思考写到这里我想聊一点可能有些敏感但很有价值的话题。最近几年Storm 的使用者确实在减少很多人转向了 Flink。Flink 在容错上的最大卖点是“精确一次的状态一致性”它基于 Chandy-Lamport 分布式快照算法做 checkpoint在保证状态一致性的同时吞吐能力还很强。但我不认为这意味着 Storm 就“过时了”。Storm 的架构设计依然是理解分布式计算系统的最佳教材之一它的 ack 机制、去状态化设计、ZooKeeper 协调方式至今仍是很多系统的设计蓝本。如果你的团队对实时计算的延迟要求极高并且业务能够接受 at-least-onceStorm 依然是一个足够稳定且轻量的选择。反过来说如果你需要大规模状态存储、窗口计算和精确一次保障那么 Flink 会是更合适的工具。架构选型没有绝对的对错只有“适不适合当前场景”。我这篇文章写的是 Storm 的容错但希望你学到的不只是 StORM 本身更是那种“从故障角度看系统设计”的思路。我在实际运维中最大的体会是容错机制的最优状态是“隐形”的——平时你感觉不到它的存在直到故障发生的那一刻你才意识到它在默默兜底。但这份“隐形”是有代价的它要求你在配置上做足功课在监控上保持敬畏在生产环境的每一次变更前先想清楚“如果这台机器现在就挂了会发生什么”。如果你正在设计自己的实时计算体系或者准备深入排查一个诡异的数据丢失问题希望这篇文章能成为你的第一份排查手册。先从开启 acker 和正确使用 msgId 开始再逐步模拟节点故障去验证你的集群你会发现Storm 的容错机制并没有那么玄妙它只是把每一件小事都做对了而已。
返回列表