{T}

Apache Storm:最早的开源流计算框架

0. 引言

进入模块四,开始用前面建立的核心概念(DAG、队列、线程、流数据状态、流信息状态、反向压力)来审视开源流计算框架。Apache Storm 是 Twitter 开源的大规模分布式流计算平台,是最早被广泛接受的流计算框架。本文从系统架构、流的描述、流的处理、流的状态、消息处理可靠性五个方面分析。

1. 分析框架的统一思路

回顾自建框架的不足,可归纳为五类考察维度,这也是后续所有框架的分析框架:

维度考察内容
系统架构主从结构、组件划分、资源管理
流的描述DAG 拓扑与编程 API 的抽象程度
流的处理处理 API、是否支持反向压力
流的状态流数据状态与流信息状态的支持
消息处理可靠性at most once / at least once / exactly once

2. 系统架构:Nimbus + Supervisor + ZooKeeper

图表渲染中…
  • Master(Nimbus):代码分发、任务分配、状态监控;
  • Worker(Supervisor + Worker 进程):Supervisor 管理 Worker 生命周期,Worker 创建 Executor 线程执行 Task;
  • ZooKeeper:在 Nimbus 与 Supervisor 之间共享作业状态、协调调度。

理解主从架构后,会发现大数据系统与微服务架构(Kubernetes 等)整体同构,可触类旁通。

3. 流的描述:Topology、Tuple、Stream、Spout、Bolt

概念对应抽象说明
TopologyDAG完整描述流计算执行过程,部署后持续运行直到显式停止
Tuple消息一条消息
Stream队列(消息流)无边界 Tuple 序列
Spout输入源从外部数据源读取数据并发送到 Stream
Bolt线程(处理单元)消息的过滤、运算、聚类、关联、数据库访问等逻辑

Topology 的节点对应 Spout/Bolt,边对应 Stream——抽象概念在 Storm 中一一落实。

4. 流的处理:Stream API

早期版本用 TopologyBuilder 构建应用,API 较底层。Storm 2.0.0 起提供更现代化的 Stream API(目前仍标记为实验性),建议新开发直接使用。

4.1 输入:Spout

java
public class DemoWordSpout extends BaseRichSpout {
    public void nextTuple() {
        Utils.sleep(100L);
        String word = words[rand.nextInt(words.length)];
        this._collector.emit(new Values(word));  // 逐条读取并发射 Tuple
    }
}
StreamBuilder builder = new StreamBuilder();
Stream<String> words = builder.newStream(new DemoWordSpout(), new ValueMapper<String>(0));

Spout 的 nextTuple 逐条读取消息;ack/fail 与消息可靠性相关(成功后确认、失败后重发或处理)。

4.2 处理与输出

Stream API 提供三类操作:常用流式处理(filtermapreduceaggregate)、流数据状态操作(windowjoincogroup)、流信息状态操作(updateStateByKeystateQuery):

java
wordCounts = words
    .mapToPair(w -> Pair.of(w, 1))
    .countByKey();          // 单词计数

wordCounts.forEach(new Print2FileConsumer());  // 输出:print/peek/forEach/to

输出操作:print(stdout)、peek(原样中继并检测各阶段状况)、forEach(任意逻辑)、to(复用既有 Bolt 输出)。

4.3 反向压力

  • 早期机制:acker + max.spout.pending 参数——Spout 发出但未确认的消息数超阈值即暂停发送。缺陷:静态配置难达最优、只在源头限速导致处理速度抖动;
  • 改进机制:除监控未确认消息数外,还监控每级 Bolt 接收队列的消息数,超阈值时通过 ZooKeeper 通知 Spout 暂停——实现各阶段动态反向压力。

5. 流的状态

5.1 流数据状态(早期方案)

方案说明
Trident将流切分为元组块(Tuple Batch)分发处理,提供 State API 处理状态与事务一致性;支持 Transactional(强一致)/ Opaque Transactional(弱一致)/ No-Transactional(无保证)三级
窗口Bolt 按滑动窗口(Sliding)与滚动窗口(Tumbling)处理
自定义批处理内置 tick tuple 定时机制,可自行实现窗口、Watermark 等

5.2 流信息状态

Trident 状态接口耦合 Trident 机制,非 Trident 的 Topology 无法使用;Stream API 引入更通用的 updateStateByKeystateQuery,对普通 Topology 同样适用。

6. 消息处理可靠性

Storm 提供尽力而为(best effort)、至少一次(at least once)、限于 Trident 的精确一次(exactly once)三级保证。核心机制是追踪消息是否被完全处理:一条消息(及其子元组、孙元组……)全部成功处理才算完成,任一失败即视为失败。使用时需配合两件事:生成新元组时 ack 告知系统、处理完成时通过 ack/fail 上报。

业务选型建议:按需选择可靠性级别即可——越严格的可靠性,实现越复杂、性能损耗越大,很多场景 at least once 已经足够。

7. 小结

  • 架构:Nimbus + Supervisor + ZooKeeper 的标准主从结构;
  • 描述:Topology(DAG)/ Tuple(消息)/ Stream(消息流)/ Spout(输入)/ Bolt(处理);
  • 处理:早期 TopologyBuilder 底层,2.0+ Stream API 现代化;反向压力从源头限速演进为逐级动态监控;
  • 状态:Trident(批式、耦合重)→ Stream API(窗口/join/cogroup + updateStateByKey/stateQuery);
  • 可靠性:追踪式确认机制,best effort / at least once / exactly once(仅 Trident)。

Storm 曾风靡一时,但因早期设计与状态接口的不足,且后继者 Flink 在各方面都能取代它,逐渐没落。入门流计算建议直接学习 Flink

下一章讲解 Spark Streaming:从批处理走向流处理。