
第一次在生产环境搭 Apache Storm 集群我盯着官方架构图看了半天第一反应是这玩意儿怎么和 Hadoop 这么像都有一个主控节点都有一堆从节点都依赖 ZooKeeper。但真正把一个 Topology 跑起来仔细跟踪一条数据从进入集群到输出结果的全过程之后我才反应过来Storm 的架构并不是 Hadoop 的简单翻版它把实时计算这件事的底层假设彻底换了一套。这篇文章不打算做概念堆砌而是想把 Storm 的核心组件、它们各自的角色以及组件之间到底怎么协同工作这件事讲清楚。我尽量按自己实际用过、踩过坑之后的认知来写适合两类人一类是刚接触 Storm、被一堆名词搞晕的新人另一类是已经在写 Topology但对底层机制理解得比较模糊、遇到问题只能靠猜的老选手。1. 实时计算为什么需要专属架构批处理思维的局限与 Storm 的解法1.1 批处理计算模型的三个不适合实时的点先回到最根本的问题为什么实时计算不能直接拿 Hadoop 那套批处理框架来扛批处理的核心假设是数据可以等。把一天产生的日志落盘第二天凌晨跑一次 MapReduce扫一遍全部数据输出报表。这个模式下数据是一批、一批地处理批次之间可以离线等待调度系统追求的指标是吞吐量——单位时间能处理多少条记录而不是这条记录从到达系统到产出结果中间隔了多少毫秒。流式计算的核心假设完全反过来数据不停在来你没法等也没那么大的耐心把无穷无尽的数据攒下来。你的目标是在数据到达的瞬间尽快完成处理把结果送到下游追求的是延迟指标。为什么不能折中一下每隔几十秒跑一次 MapReduce也能做但问题很多。MapReduce 每个作业都要经历任务提交、调度、启动 JVM、数据混洗这些阶段固定开销太大频繁跑批意味着大量的任务启动和销毁在烧资源。更关键的是基于文件系统的批处理天然面向已经存在的数据而流式数据的输入是持续无界的攒几十秒再算在实时性要求高的场景——比如风控、秒级监控告警——仍然不可接受。1.2 Storm 的设计目标常驻进程、有向无环图、数据即流Storm 的解决思路很直接既然频繁启停任务不行那就让任务一直跑着。一个 Topology 提交到集群之后就是一个持续运行的常驻计算任务除非你手动 kill否则它永远在跑。数据不是分批发给它而是像水流一样持续灌入每个处理节点收到一条数据就立刻处理处理完立刻往下游发。这里引出了 Storm 架构中最核心的抽象——Topology。它本质上是一个有向无环图DAG节点是计算逻辑边是数据流向。数据从源头节点Spout不断流出经过一个个处理节点Bolt最终到达外部存储或者完成一次聚合输出。这个模型其实和工厂里的流水线非常像原材料的入口只有一个中间经过多个工位加工不同工位干不同的事最后输出成品。理解了这一点再往后看 Storm 的各种组件思路就清晰很多。物理层面它要解决谁来调度、谁来执行、状态存哪逻辑层面它要解决数据从哪来、在哪儿处理、怎么分配而组件之间的协同机制则负责把这些拼起来。2. 物理层面的角色分工Nimbus、Supervisor、Worker 和 ZooKeeper 各管什么2.1 四个角色的核心职责熟悉 Hadoop 的人对这四个名词肯定不陌生但 Storm 里的分工细节完全不一样。我用一句话先概括Nimbus 负责调度和分配任务ZooKeeper 负责记住集群里所有的元数据状态Supervisor 负责管理本节点上的执行进程Worker 负责真正干活。逐个拆开来看。Nimbus 是集群的主控节点也是大脑。它接收你提交的 Topology 代码把这份代码打包分发出去同时把 Topology 的结构转换成一份任务分配清单——哪个 Supervisor 的哪个端口上该跑哪个任务的哪个部分。但值得注意的是Nimbus 本身不参与任何数据流的处理它只是一个调度器。你提交完 Topology 之后Nimbus 的日常工作就很轻松了只需要盯着集群的健康状态和任务分配状态。ZooKeeper 是整个集群的记忆体。Nimbus、Supervisor、Worker 之间的所有协调信息包括任务分配、心跳状态、Topology 运行进度、元数据快照全部存在 ZooKeeper 里。Storm 选择 ZooKeeper 承担这个角色最关键的原因是让 Nimbus 和 Supervisor 都变成无状态组件。所谓无状态就是它们挂掉之后可以从 ZooKeeper 中读回全部状态信息重新恢复工作不需要在本地保留任何持久化数据。Supervisor 是跑在从节点上的守护进程。它在 ZooKeeper 上监听有没有分配给自己节点的任务一旦发现有新的分配就根据分配清单启动或者停止本地的 Worker 进程。同时 Supervisor 会监控这些 Worker 的心跳某个 Worker 挂掉之后Supervisor 负责把它重新拉起来。Worker 是最底层的执行容器就是一个 JVM 进程。一个 Worker 可以运行多个 Executor 线程每个 Executor 线程运行一个或多个 Task即 Spout 或 Bolt 的实例。多个 Worker 之间通过网络通信交换数据所以 Worker 数量实际上决定了 Topology 能利用的网络带宽和 CPU 资源。下面这张表是我自己整理的角色对照方便你记组件类比核心职责是否处理数据流Nimbus生产调度员任务分配、作业提交、故障转移否ZooKeeper工厂公告栏和资料室保存集群元数据与协调状态否Supervisor车间主任管理本节点 Worker 进程的启停否Worker车间/工段运行 Executor 线程承载具体计算是2.2 为什么 Nimbus 设计成无状态这个设计在早期版本的 Storm 里尤其有代表性。Hadoop 1.x 的 JobTracker 之所以备受诟病就是因为单点状态太重所有作业的状态都保存在内存里挂掉之后恢复困难成了整个集群最脆弱的环节。Storm 在设计上做了一个非常果断的取舍Nimbus 不保存任何运行期状态所有状态都放 ZooKeeper。于是 Nimbus 挂了正在运行的 Topology 完全不受影响——因为 Worker 和 Supervisor 之间的数据交换根本不经过 Nimbus数据流是 Worker 对 Worker 直连的。只是新的作业提交和故障转移会暂时不可用等 Nimbus 重新起来之后它再从 ZooKeeper 里读回状态继续履行调度职责。不过这里要提醒一句早期 Storm 版本的 Nimbus 是单点生产环境确实存在Nimbus 挂了集群暂时没法调度新任务的问题。从 1.0 开始官方支持 Nimbus HA可以同时部署多个 Nimbus 实例通过 ZooKeeper 做选主。如果你在维护老版本集群这点务必注意。3. 逻辑层面的核心组件Topology、Spout、Bolt、Stream 与 Grouping 策略3.1 Topology 的构成和提交逻辑物理层解决的是任务跑在哪逻辑层解决的是这个任务到底算什么。Storm 的逻辑单元是一个 Topology它由三部分组成Spout、Bolt、Stream Grouping。Spout 是数据源Topology 里所有数据的入口。它负责从外部系统读取数据比如从 Kafka 读日志、从消息队列读订单事件、从传感器网关读设备数据。Spout 的核心工作是发射emitTuple 到数据流中。它通常有两种状态一种是主动推送数据比如不停轮询 Kafka 拉消息一种是等待外部系统回调。Bolt 是数据处理节点接收上游发来的 Tuple执行自己的处理逻辑再发射新的 Tuple 给下游。一个 Bolt 可以干很多事过滤、聚合、连接、与外部存储交互、计算指标。理论上你可以把整个处理逻辑写在一个 Bolt 里但通常做法是拆成多个 Bolt每个螺栓只干一件事便于维护也便于并行扩展。Stream 和 Tuple 是数据流动的基本载体。Stream 是一组无界的 Tuple 序列可以理解成流水线的一段传送带Tuple 是传送带上的一个零件也就是 Storm 中的基本数据单元。Tuple 本质上是一个有序的字段列表每个字段有一个名字和对应的值格式上和数据库表的一行记录非常相似。一个简单的 Topology 定义代码长这样TopologyBuilder builder new TopologyBuilder(); // 数据源从 Kafka 读取订单事件 builder.setSpout(order-spout, new KafkaSpout(kafkaSpoutConfig)); // 处理节点1按逗号切分订单事件4个并行实例 builder.setBolt(split-bolt, new SplitOrderBolt(), 4) .shuffleGrouping(order-spout); // 处理节点2按订单ID聚合统计2个并行实例 builder.setBolt(count-bolt, new OrderCountBolt(), 2) .fieldsGrouping(split-bolt, new Fields(orderId));这里的关键信息有两个每个 Bolt 后面跟的并行度数字以及每个 setBolt 后面的 Grouping 声明。前者决定逻辑组件在物理上开多少个实例后者决定数据流怎么分配给这些实例——也就是下面要说的数据分发策略。3.2 六种数据分发策略怎么选Stream Grouping 定义了上游 Stream 中的 Tuple 如何被分发到下游 Bolt 的不同 Task 上。Storm 内置的 Grouping 策略有七种按使用频率我将最常用的六种整理成一张对比表Grouping 策略分配逻辑典型用途Shuffle Grouping随机均匀分发无状态过滤、切分、简单转换Fields Grouping按指定字段哈希分发同字段值进同一 Task按 key 聚合、状态化处理All Grouping广播给全部 Task全局配置下发、窄依赖广播Global Grouping全发给编号最小的 Task强制全局有序处理Direct Grouping由上游显式指定目标 Task特殊路由场景只能配合 emitDirectLocal or Shuffle Grouping优先在本地 Worker 内随机分发跨 Worker 传输开销敏感的场景这个选择其实反映了一个本质问题数据要不要按 key 分组处理。如果处理逻辑是无状态的比如只要把日志里的字段解析出来往下传用 Shuffle Grouping 就行因为谁处理都一样随机分发最均匀负载最平衡。如果处理逻辑有状态比如按订单 ID 做累加统计那么同一个订单 ID 的所有数据必须进入同一个 Task才能保证状态不冲突这时候必须用 Fields Grouping。我第一次写 Topology 的时候把状态维护相关的 Bolt 配成了 Shuffle Grouping结果统计结果一塌糊涂。原因很简单不同 Task 各自维护了独立的计数状态同一个订单 ID 被随机分到了不同的 Task各自的计数都是残缺的。换成 Fields Grouping 按订单 ID 哈希之后问题立刻消失。这个坑特别经典新手几乎必踩。另外说一下没有进表的 None Grouping。它的语义和 Shuffle Grouping 类似但声明时不关心分配目标官方文档表示它不保证任何分组行为实际使用中性能相比 Shuffle Grouping 提升并不明显所以我一般不建议优先使用。选型时把上面六种搞清楚就够了。3.3 Executor 和 Task 的关系这里补一个容易混淆的概念。Task 是 Spout 或 Bolt 的实际实例而 Executor 是承载 Task 运行的线程。默认情况下一个 Executor 运行一个 Task。但你可以通过配置让一个 Executor 运行多个 Task——多个 Task 共享同一个线程也就是把多个 Task 串行跑在一个线程里。这种配置在 TopologyBuilder 里可以通过 setBolt 的重载参数控制分别指定 executor 数量和 task 数量。实际生产里我基本都保持默认的 1:1 关系。一个线程跑多个 Task 在多数场景下收益不大反而会让单个线程变成一个热点瓶颈。如果发现某个 Bolt 实例的负载太高优先做的是增加并行度而不是让一个 Executor 里挤更多的 Task。4. 数据流在集群中的完整旅程从 Spout 发起到 Bolt 处理落地的每一步4.1 数据跨节点传递的物理路径理解了逻辑组件接下来要看数据到底怎么从上游流到下游。这里我把一条 Tuple 从 Spout 发射出来到 Bolt 处理完成的全过程拆开按物理路径一步步走。假设一个 Topology 有两个组件A 是 Spout运行在节点 1 的 Worker 1 上B 是 Bolt运行在节点 2 的 Worker 3 上。数据从 A 到 B需要经历六个环节第一Spout 发射 Tuple。Spout 调用 emit 方法后Tuple 不会直接飞到网络里去而是先进入当前 Worker 内部的发送队列。这个队列是 LMAX Disruptor 设计的无锁环形队列Storm 用它替代传统的 LinkedBlockingQueue就是为了在高吞吐场景下减少锁竞争和 GC 压力。Tuple 会在这里做序列化——把内存里的 Java 对象变成二进制字节流。第二Worker 1 的发送管理器把序列化后的消息通过网络发送给 Worker 3。Storm 默认使用的底层通信框架是 Netty早期版本用的则是 ZeroMQ。这里值得留意的是数据在 Worker 之间传的是序列化好的字节流而不是 Java 对象所以跨节点的数据传输天然要求自定义类型是可序列化的。第三Worker 3 收到数据后把字节流放入自己的接收队列。这个队列同样是 Disruptor 实现。接收队列的深度有上限满了之后会发生背压发送端会暂缓发送。第四Worker 3 内部有一个分发机制根据消息里的目标 Task 信息把数据从接收队列分发到对应 Executor 的接收队列。第五Executor 线程从自己的队列里取出一条 Tuple交给对应的 Bolt 实例执行 execute 方法。第六Bolt 处理完之后向 acker 发送确认消息并向它的下游发射新的 Tuple。如果这个 Bolt 已经是拓扑的末端处理过程就是一个完整的闭环。4.2 本地通信与跨节点通信的差异上面讲的场景是跨节点传输。实际上如果上游和下游落在同一个 Worker 进程里数据流的成本会低很多——不需要序列化不需要走网卡直接从发送端的内存块拷贝到接收端队列就行。Storm 的性能优化很多时候就是围绕这一点展开的。Local or Shuffle Grouping 存在的意义就在这里。它优先把 Tuple 分发给同一个 Worker 内的下游 Task只有本地找不到合适目标时才走网络。尤其是在数据量巨大、下游并行度又高的场景下启用这个策略能明显减少集群的带宽压力也能降低单条消息的端到端延迟。我自己做过一次小实验同样的 Topology下游 Bolt 的 Grouping 从 Shuffle 换成 Local or Shuffle在跨节点占比高的场景下整体吞吐大约提升了 15% 到 20%。这不是玄学而是省掉了序列化和网络传输这两个大头开销。所以如果你的业务允许尽量优先考虑这个策略。4.3 背压、超时和消息积压Storm 对数据流的一个默认假设是下游处理速度应该跟得上上游发射速度如果跟不上队列会逐渐填满。在没有背压机制的早期版本里队列满了之后数据只能丢弃或者大量堆积最后导致不稳定的延迟。从 1.0 开始Storm 引入了自动背压当接收队列超过高水位线时会逐级向上游反馈限速信号直到上游暂停发射等队列水位降下来再恢复。这个机制在真实场景里非常重要。我遇到过一档经典问题Kafka 里的消息短时间爆发Spout 疯狂拉取下游有个 Bolt 因为要做外部存储写操作速度上不去结果集群内存持续飙升。后来我做了两件事一是设置合理的 max spout pending限制 Spout 同时在途的未确认 Tuple 数量二是打开自动背压开关。两个措施配合才把内存压在一个可控范围。这些都是数据流机制层面的配套手段不理解队列和背压机制的人遇到这种问题往往只会盲目加资源效果却很差。另外Storm 还引入了 message timeout 概念。一条 Tuple 从 Spout 发射开始计时默认 30 秒内如果整棵树处理不完Spout 会认为失败触发重发。这个机制后面讲可靠性的时候还会再提但这里先记住一点超时时间不是越长越好因为它直接影响一次失败重试的周期也不宜设太短否则本身就慢的外部操作会频繁误报失败。5. 故障恢复与可靠性协同Worker 挂掉、Nimbus 挂掉和 Ack 机制都在做什么5.1 心跳、任务重新分配与三板斧恢复流程流式计算是常驻任务最怕的就是节点故障。Storm 的故障恢复机制是整个架构里最能体现协同两个字的部分因为一次故障的感知和处理涉及 Worker、Supervisor、Nimbus、ZooKeeper 四个角色共同配合。我把一次完整故障恢复拆成四个步骤第一步心跳丢失。每个 Worker 会定期通过 Supervisor 向 ZooKeeper 上报心跳Supervisor 自身也会向 ZooKeeper 上报心跳。这些心跳数据就是 ZooKeeper 节点上的临时会话信息。如果某个 Worker 进程异常退出它的临时节点会消失。第二步Supervisor 感知到本地 Worker 挂了之后会尝试在本节点上重新启动一个新的 Worker 进程这个新进程会接管原来的 Executor 和 Task。如果是整个节点宕机Supervisor 自己也挂掉了那 Nimbus 会通过 ZooKeeper 感知到这个节点失联。第三步Nimbus 把这个节点上所有分配的任务标记为失效并在存活节点上为它们重新寻找可用的 Worker 槽位生成新的任务分配信息写入 ZooKeeper。第四步其他节点上的 Supervisor 监听到新的分配信息后启动新的 Worker 接管任务。整个过程的协调中心就是 ZooKeeper。没有它Nimbus 无法准确感知集群全貌Supervisor 也无法知道该启动哪些进程。这也是我把 ZooKeeper 称作记忆体的原因。不过要说明一点故障恢复意味着任务状态会丢。Storm 默认语义是 at least once也就是说 Tuple 在故障发生后可能被重复处理。如果你在 Bolt 里维护了一些本地状态比如计数器、窗口累加值这些状态在 Worker 重启后会归零导致结果不准确。想要更严格的状态一致性要么依赖外部存储做状态管理要么上 Trident 这类更高层抽象要么干脆换 Flink 这种自带状态管理和 exactly-once 语义的引擎。这些选型问题我在后面详细讲。5.2 Ack 机制异或校验如何保证消息可达可靠性协同的另一块核心是 Ack 机制。它要解决一个很朴素的问题Spout 发出的一条消息经过多个 Bolt 加工之后到底有没有被完整处理如果中间某个 Bolt 挂掉导致某个 Tuple 丢失Spout 怎么知道并且怎么重新补救Storm 的解决办法是引入一个系统级的 Acker Bolt。每条从 Spout 发射出来的 Tuple 会生成一个 64 位的 rootIdAcker 维护这个 rootId 的校验状态。Bolt 每处理完一个 Tuple 并且完整继承父 Tuple 和子 Tuple 的关系之后会向 Acker 发送一个确认消息。Acker 收到确认后用异或操作更新状态最终当校验值归零Acker 通知 Spout 这个 Tuple 处理成功如果在 timeout 时间内没有归零则通知 Spout 失败并重发。这个 XOR 校验方式非常巧妙它不需要逐条追踪每棵处理树里每个节点的落单关系只需要在汇总层做异或聚合。但这套机制也有代价每个 Tuple 都要增加额外的序列化和网络传输开销。所以 Storm 允许你在 Spout 声明时关闭 acking也可以设置 topology.acker.executors 为 0 来完全关闭。取舍就是要可靠性就要付开销不要可靠性就省下来换吞吐。大部分生产场景我会开着 acking因为流式计算一旦丢数据下游的统计和风控结果就会悄悄出错排查起来的成本远比那点吞吐损失高。5.3 三种投递语义与生产中的取舍把 ack 机制和故障恢复放在一起就能看清楚 Storm 能做到什么语义级别At most once最多一次关闭 ack数据可能丢但延迟最低。At least once至少一次默认语义数据不会丢但可能重复。Exactly once恰好一次需要额外机制配合比如利用外部事务或 Trident 提供的语义。生产选型时常常有个误区觉得肯定要 exactly once但 exactly once 的代价很大而且很多时候业务本身能容忍重复。比如做指标监控一条重复数据最多让告警触发两次影响不大而如果下游是金融交易计数重复就不可接受。关键是先搞清楚业务对重复的容忍度再决定是否需要支付额外成本。Storm 原始 API 只保证 at least once如果你需要 exactly once建议优先考虑架构上的幂等设计例如下游写入 MySQL 时用唯一键去重而不是盲目上重型框架。6. 架构设计在生产部署中的经验教训从并行度规划到集群排障6.1 并行度怎么配才合理架构层面的理解最终要落到一个具体问题上并行度到底怎么设置。Storm 里的并行度分为三层Worker 数量、Executor 数量、Task 数量。Worker 是 JVM 进程决定进程级资源消耗Executor 是线程决定并发粒度Task 是 Spout/Bolt 实例决定逻辑上有没有多个副本。通常我们最关心的是前两个整个 Topology 开多少个 Worker每个组件开多少个 Executor。我的经验是先看上游数据量再看单实例处理能力最后反过来推并行度。比如你用单线程消费一个 Kafka partition发现一条消息的平均处理时间是 5 毫秒那单个 Executor 的吞吐上限大约是每秒 200 条。如果上游每秒来 1000 条这个 Bolt 就需要至少 5 个 Executor。当然这只是粗算实际还要考虑中间传输、反序列化、外部 IO 抖动的影响保守做法是乘一个 1.5 到 2 的冗余系数。Worker 数量一般建议不要开太多因为每个 Worker 都是独立 JVM进程数量多了内存开销和全集群协调的负载都会上来。常见做法是让一个 Supervisor 节点上跑 4 到 8 个 Worker每个 Worker 分配 2 到 4 个 Executor。具体数值要结合节点内存和 CPU 核数调整。还有一个容易被忽略的点并行度不是越大越好。加并行度意味着更多的线程切换、更多的队列交互和网络连接超过某个临界点之后继续加并行度反而会让性能下降。我见过有人把一个简单的过滤 Bolt 开了 64 个 Executor集群资源被吃满但吞吐几乎没有提升。合理定位瓶颈之后再针对性扩容才是正确姿势。6.2 常见瓶颈和一次真实排障过程集群跑久了最常见的问题类型就几类延迟飙高、内存溢出、数据倾斜、消费堆积。我印象最深的是一次延迟飙高的排查。现象是某个 Bolt 的平均处理时间从几毫秒涨到了几百毫秒整个 Topology 的完成时间大幅上升。第一反应是看这个 Bolt 的代码发现它在一个 for 循环里调用了外部 REST 接口接口响应偶尔超时。这里的问题不是代码逻辑不对而是把同步阻塞 IO 放进了处理链路的主线程——一个 Executor 线程在等网络响应的时候这个 Executor 的队列里堆积的所有消息都在等。解决方案是把这个外部调用改成异步或者把慢操作单独拆一个 Bolt 放到另一个线程池里处理。这也是架构设计层面的一个问题Bolt 里到底该不该做重量级 IO我的原则是Bolt 只做纯计算和轻量级状态更新所有需要访问外部系统的操作都尽量拆到独立组件里或者至少用异步化处理避免阻塞处理线程。另外一个高频坑是数据倾斜。Fields Grouping 按 key 哈希分发以后如果某个热门 key比如某个大卖家的订单量远高于其他卖家占了数据的大部分对应的 Task 就会成为热点其他 Task 闲着。这种倾斜本质上是分组策略的问题。有时靠改善 key 的粒度比如加盐拆分可以缓解有时需要重设计聚合层级——做两层聚合第一层按盐分散负载第二层再按自然 key 聚合还原结果。6.3 到底是选 Storm 还是 Flink最后聊一个很多人在架构选型时会纠结的问题都到现在的技术生态了还要不要学 Storm新项目要不要用 Storm客观讲Storm 在实时计算领域是一个里程碑式的存在它的架构设计影响了后来的 Flink、Heron 等很多系统。Flink 在状态管理、窗口计算、exactly-once 语义、流批一体方面都做得更完善社区也更活跃所以新项目建议优先考虑 Flink。但我个人的看法是理解 Storm 的架构依然很有价值。一方面很多遗留系统还在大规模使用 Storm会读、会调、会排障本身就是一项可以落地的技能另一方面Storm 的组件划分和协同机制——无状态调度器、流式 DAG、可靠性树——这些设计思想非常基础理解了它们再去看 Flink 的 JobManager、TaskManager、Checkpoint 机制会发现很多概念都是相通的。知识从来不会白学架构思维更是如此。最后分享一个我自己实际排障时的小技巧当 Topology 的延迟或吞吐出现问题先不要急着看代码先打开 Storm UI 的页面依次扫三个数据——Spout 的 emitted 数量和 acked 数量是否匹配、每个 Bolt 的 execute latency、各个 Executor 的队列水位。这三个数据基本能定位 80% 的问题。如果 emitted 远大于 acked大概率是下游处理慢或消息超时重发如果某个 Bolt 的 execute latency 明显高于其他 Bolt就该去查那个组件的代码和外部依赖了。这个思路其实就是把前面说的架构知识落到了排障动作里。