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
| 概念 | 对应抽象 | 说明 |
|---|---|---|
| Topology | DAG | 完整描述流计算执行过程,部署后持续运行直到显式停止 |
| 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
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 提供三类操作:常用流式处理(filter、map、reduce、aggregate)、流数据状态操作(window、join、cogroup)、流信息状态操作(updateStateByKey、stateQuery):
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 引入更通用的 updateStateByKey 与 stateQuery,对普通 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:从批处理走向流处理。