Apache Storm 从入门到实践:架构、可靠性、开发、部署与性能调优
Apache Storm 是一个开源的分布式实时计算系统,适合持续处理日志、消息、监控指标、交易事件等无界数据流。本文将系统介绍 Storm 的核心概念、集群架构、可靠性机制、并行度模型,并通过一个完整的 WordCount 案例演示拓扑开发、打包、运行、部署和调优。
1. Storm 是什么
Apache Storm 是一个分布式实时流处理框架。它可以持续接收数据,并在数据到达后立即完成过滤、转换、聚合、分组、持久化等操作。
Storm 处理的通常是一个没有明确结束时间的数据流,例如:
- Kafka 中持续产生的订单消息;
- 应用服务器不断输出的日志;
- 物联网设备持续上报的传感器数据;
- 用户在网站或 App 中产生的行为事件;
- 交易系统中的支付、退款和风控事件;
- 监控系统中的 CPU、内存、接口耗时等指标。
传统批处理系统通常需要先积累一批数据,再启动任务处理。Storm 则会让计算拓扑长期运行,只要上游有数据,就可以立即处理。
可以用下面的方式理解 Storm 与批处理系统的区别:
| 对比项 | 批处理系统 | Storm |
|---|---|---|
| 数据类型 | 有界数据集 | 无界数据流 |
| 运行方式 | 任务处理完成后退出 | 拓扑持续运行 |
| 处理时机 | 定时或数据积累后处理 | 数据到达后立即处理 |
| 延迟 | 秒级、分钟级或更高 | 通常为毫秒级或秒级 |
| 典型用途 | 离线报表、历史统计 | 实时监控、实时计算、事件处理 |
Storm 本身不是消息队列,也不是数据库。
在一个完整的数据架构中,它通常位于消息系统和存储系统之间:
业务系统
↓
Kafka、RabbitMQ 等消息系统
↓
Apache Storm 实时计算
↓
MySQL、Redis、Elasticsearch、HBase、ClickHouse 等存储
消息队列负责保存和传递消息,Storm 负责处理消息,数据库或缓存负责保存处理结果。
2. Storm 的应用场景
2.1 实时日志处理
应用服务器不断产生访问日志、异常日志和审计日志。Storm 可以实时完成:
- 日志格式解析;
- 敏感字段脱敏;
- 异常日志分类;
- 接口请求量统计;
- 错误率计算;
- 日志写入 Elasticsearch;
- 达到告警阈值后发送通知。
典型链路如下:
应用日志 → Kafka → Storm → Elasticsearch → Kibana
↓
告警系统
2.2 实时指标统计
电商、支付、广告等系统经常需要实时统计:
- 每分钟订单量;
- 每个商户的交易金额;
- 每个接口的成功率;
- 每个地区的访问人数;
- 广告点击量和转化率;
- 实时在线人数。
Storm 可以按照用户、商户、地区、商品等字段重新分组,使相同业务键的数据进入同一个处理任务。
2.3 实时风控
支付或交易事件进入 Storm 后,可以执行:
- 单位时间内交易次数统计;
- 异常 IP 检测;
- 异地登录检测;
- 高频退款检测;
- 设备指纹关联;
- 黑名单匹配;
- 风险规则计算;
- 风险评分。
需要注意:Storm 可以承担实时规则计算,但余额、支付状态和最终账务结果不能只保存在 Storm 内存中,必须由可靠的外部存储负责。
2.4 实时 ETL
Storm 可以从消息队列读取原始数据,然后执行:
- 字段校验;
- 数据清洗;
- 格式转换;
- 无效数据过滤;
- 数据补充;
- 多流关联;
- 写入目标数据库。
2.5 物联网数据处理
设备持续上报温度、湿度、电量、位置等信息,Storm 可以实时完成:
- 异常值过滤;
- 设备状态聚合;
- 阈值告警;
- 滑动窗口统计;
- 实时趋势计算。
2.6 不适合使用 Storm 的场景
以下场景通常不需要优先考虑 Storm:
- 每天只执行一次的离线报表;
- 主要处理有限历史数据的批任务;
- 强依赖复杂 SQL 分析的离线数仓任务;
- 需要大规模状态计算,但团队没有足够的状态管理经验;
- 简单消费 Kafka 后写入数据库,不需要分布式计算;
- 业务量很小,一个普通消费者程序就能完成处理。
如果任务只是“读取消息并调用一个接口”,直接编写 Kafka Consumer 往往更加简单。
3. Storm 的核心概念
Storm 最重要的概念包括:
| 概念 | 说明 |
|---|---|
| Topology | 拓扑,描述完整的实时计算流程 |
| Stream | 数据流,由连续不断的 Tuple 组成 |
| Tuple | Storm 中的一条数据记录 |
| Spout | 数据源,负责读取并发送数据 |
| Bolt | 数据处理节点,负责转换、统计、持久化等操作 |
| Stream Grouping | 决定数据如何分发给下游任务 |
| Worker | 执行拓扑的 JVM 进程 |
| Executor | Worker 中执行组件任务的线程 |
| Task | Spout 或 Bolt 的实际处理实例 |
| Acker | 跟踪 Tuple 是否被完整处理的系统任务 |
3.1 Topology
Topology 是 Storm 中最高层的抽象,表示一个长期运行的实时计算应用。
一个拓扑由 Spout、Bolt 和它们之间的连接关系组成:
KafkaSpout
↓
ParseBolt
↓
FilterBolt
↓
AggregateBolt
↓
DatabaseBolt
它与普通批处理任务最大的不同是:拓扑提交后通常不会自动结束,而是持续运行,直到被主动停止。
3.2 Stream
Stream 表示一个无界的数据序列,由连续产生的 Tuple 组成。
例如,一个订单数据流可能包含:
{orderId=10001, userId=20001, amount=99.90}
{orderId=10002, userId=20002, amount=199.00}
{orderId=10003, userId=20001, amount=59.90}
一个 Spout 或 Bolt 可以声明多个 Stream。每个 Stream 可以拥有自己的 Stream ID 和字段结构。
3.3 Tuple
Tuple 是 Storm 中传输和处理的基本数据单元,可以理解为一条具有字段名称的记录。
例如:
new Values("order-10001", "user-20001", 99.90)
对应的字段可以声明为:
declarer.declare(new Fields("orderId", "userId", "amount"));
下游 Bolt 可以按字段名称读取数据:
String orderId = input.getStringByField("orderId");
String userId = input.getStringByField("userId");
Double amount = input.getDoubleByField("amount");
生产环境应尽量使用简单、稳定的数据类型,并控制 Tuple 大小。对于自定义对象,应显式注册序列化器,不要依赖低效或不安全的 Java 原生序列化回退。
3.4 Spout
Spout 是拓扑的数据源,主要职责包括:
- 从 Kafka、RabbitMQ、文件或接口读取数据;
- 将数据转换成 Tuple;
- 为可靠消息设置 Message ID;
- 在处理成功时接收
ack; - 在处理失败时接收
fail; - 必要时重新发送失败消息。
Spout 的核心方法是 nextTuple()。
需要特别注意:nextTuple() 不应该长时间阻塞,因为 Storm 会在同一个线程中调用 Spout 的相关方法。
3.5 Bolt
Bolt 是数据处理节点,可以完成:
- 数据解析;
- 数据过滤;
- 字段转换;
- 聚合统计;
- 多流关联;
- 调用外部服务;
- 写入数据库;
- 发送新消息;
- 产生告警。
Bolt 的核心方法是 execute(Tuple input)。
一个复杂业务通常会拆分成多个 Bolt,而不是将所有逻辑都写进一个 Bolt。
4. Storm 集群架构
Storm 集群主要由 Nimbus、Supervisor、ZooKeeper、Worker、UI 和 Logviewer 等组件组成。
4.1 Nimbus
Nimbus 是 Storm 集群的控制节点,主要负责:
- 接收客户端提交的拓扑;
- 保存和分发拓扑代码;
- 计算任务分配;
- 将任务调度到 Supervisor;
- 监控拓扑运行状态;
- 在节点故障时重新分配任务。
Nimbus 的角色类似于集群调度中心,但它不直接处理业务 Tuple。
生产环境可以配置多个 Nimbus 节点,提高控制平面的可用性。
4.2 Supervisor
Supervisor 运行在工作节点上,主要负责:
- 接收集群中的任务分配;
- 启动和停止 Worker 进程;
- 监控本机 Worker 状态;
- 在 Worker 异常退出后重新启动它。
一个 Supervisor 可以管理多个 Worker Slot。每个 Slot 对应一个可运行 Worker 的端口。
4.3 ZooKeeper
ZooKeeper 用于 Storm 集群协调和元数据管理,例如:
- 保存集群状态;
- 保存任务分配信息;
- Nimbus 选举;
- Supervisor 心跳;
- Worker 心跳;
- 拓扑状态协调。
业务 Tuple 并不会通过 ZooKeeper 传输,因此 ZooKeeper 不是 Storm 的消息通道。
4.4 Worker
Worker 是一个 JVM 进程,只属于一个 Topology。
一个拓扑可以在多台服务器上运行多个 Worker,一个 Worker 内又可以运行多个 Executor 线程。
4.5 Storm UI
Storm UI 提供可视化管理页面,可以查看:
- 集群中的 Supervisor 数量;
- 可用和已使用的 Worker Slot;
- 正在运行的 Topology;
- Spout 和 Bolt 的吞吐量;
- Tuple 的 ACK 和失败数量;
- 处理延迟;
- Bolt Capacity;
- Worker 和 Executor 信息;
- 组件异常日志。
默认情况下,UI 常使用 8080 端口,但生产环境不应直接暴露到公网。
4.6 Logviewer
Logviewer 用来通过 Web 页面查看 Worker 日志。
由于日志中可能包含业务数据、异常堆栈和系统信息,生产环境必须配置访问控制。
5. Storm 的数据处理流程
以订单实时统计为例,整个处理流程可以是:
Kafka
↓
OrderSpout
↓
ParseOrderBolt
↓
ValidOrderBolt
↓
MerchantAmountBolt
↓
MySQL 或 Redis
一次完整处理过程如下:
OrderSpout从 Kafka 读取订单消息;- Spout 将消息封装成 Tuple;
ParseOrderBolt将 JSON 转换成订单字段;ValidOrderBolt校验订单是否合法;- 按
merchantId对消息重新分组; MerchantAmountBolt统计每个商户的金额;- 结果写入外部存储;
- Bolt 调用
ack()告诉 Storm 当前 Tuple 处理完成; - 整棵 Tuple Tree 完成后,Spout 收到成功通知;
- 如果处理超时或失败,Spout 可以重新发送消息。
Storm 的拓扑本质上是一个有向图:
- 节点是 Spout 或 Bolt;
- 边是 Stream;
- Grouping 决定 Stream 如何分区;
- Tuple 在这张图中持续流动。
6. Worker、Executor 与 Task
Worker、Executor 和 Task 很容易混淆。
它们的关系如下:
一台服务器
└── Supervisor
├── Worker JVM 1
│ ├── Executor Thread 1
│ │ ├── Task 1
│ │ └── Task 2
│ └── Executor Thread 2
│ └── Task 3
└── Worker JVM 2
└── Executor Thread 3
└── Task 4
6.1 Worker
Worker 是进程级别的执行单元:
Config config = new Config();
config.setNumWorkers(3);
这表示拓扑希望使用三个 Worker 进程。
6.2 Executor
Executor 是 Worker 中的线程。
下面的第三个参数是 Parallelism Hint,表示该组件初始使用的 Executor 数量:
builder.setBolt("parse-bolt", new ParseBolt(), 4);
这里的 4 表示初始创建四个 Executor 线程。
6.3 Task
Task 是 Spout 或 Bolt 的实际执行实例。
可以使用 setNumTasks() 指定 Task 数量:
builder.setBolt("parse-bolt", new ParseBolt(), 2)
.setNumTasks(4)
.shuffleGrouping("order-spout");
这里表示:
- Executor 数量为 2;
- Task 数量为 4;
- 平均每个 Executor 线程运行 2 个 Task。
如果没有调用 setNumTasks(),默认情况下 Task 数量通常与 Executor 数量相同。
6.4 为什么要区分 Executor 和 Task
Task 数量在拓扑生命周期内通常保持不变,而 Executor 数量可以通过 Rebalance 调整。
例如,预先创建 12 个 Task,但只使用 4 个 Executor:
builder.setBolt("count-bolt", new CountBolt(), 4)
.setNumTasks(12);
以后可以将 Executor 从 4 调整为 8,而不改变 Task 的逻辑分区数量:
storm rebalance order-topology -e count-bolt=8
需要满足:
Executor 数量 ≤ Task 数量
7. Stream Grouping 数据分组策略
Stream Grouping 决定上游 Tuple 应该发送给下游哪个 Task。
7.1 Shuffle Grouping
builder.setBolt("parse-bolt", new ParseBolt(), 4)
.shuffleGrouping("order-spout");
特点:
- 尽量均匀地把消息分配给下游 Task;
- 不保证相同业务键进入同一个 Task;
- 适合无状态解析、过滤和格式转换。
7.2 Fields Grouping
builder.setBolt("merchant-bolt", new MerchantBolt(), 4)
.fieldsGrouping("parse-bolt", new Fields("merchantId"));
特点:
- 按指定字段进行分区;
- 相同字段值会进入同一个下游 Task;
- 适合分组统计、聚合、状态计算。
例如,按 merchantId 分组后,同一商户的订单会由同一个 Task 处理。
需要注意数据倾斜。如果某个商户的数据量特别大,它对应的 Task 可能成为性能瓶颈。
7.3 Partial Key Grouping
Partial Key Grouping 也按照指定字段分区,但会在候选 Task 之间进行负载选择,适合缓解热点 Key 导致的数据倾斜。
它适用于:
- 数据按 Key 分组;
- Key 分布明显不均匀;
- 普通 Fields Grouping 出现单个 Task 过载。
是否能使用它,还要看业务状态是否允许同一个 Key 分散到多个处理实例。
7.4 All Grouping
builder.setBolt("rule-bolt", new RuleBolt(), 4)
.allGrouping("config-spout");
特点:
- 每条 Tuple 都发送给下游所有 Task;
- 类似广播;
- 适合规则、配置、黑名单更新;
- 数据量大时会明显放大流量。
7.5 Global Grouping
builder.setBolt("global-bolt", new GlobalBolt(), 1)
.globalGrouping("source-bolt");
特点:
- 所有 Tuple 都发送给下游 Task ID 最小的那个 Task;
- 可以实现全局单点聚合;
- 容易形成性能瓶颈。
7.6 Direct Grouping
Direct Grouping 由上游组件直接指定目标 Task。
上游需要声明 Direct Stream,并使用 emitDirect() 发送数据。
适合需要精确路由的特殊业务,但会增加组件之间的耦合。
7.7 Local or Shuffle Grouping
builder.setBolt("parse-bolt", new ParseBolt(), 4)
.localOrShuffleGrouping("source-spout");
如果当前 Worker 中存在下游 Task,则优先发送给同一 Worker 内的 Task,否则使用 Shuffle Grouping。
优点是可以减少跨进程网络传输。
7.8 None Grouping
None Grouping 表示不关心具体分组方式。在当前实现中,其行为通常接近 Shuffle Grouping。
7.9 分组策略总结
| Grouping | 分发方式 | 典型用途 |
|---|---|---|
| Shuffle | 尽量均匀分发 | 无状态解析、过滤 |
| Fields | 按字段分区 | 分组统计、状态计算 |
| Partial Key | 按 Key 并兼顾负载 | 热点 Key 场景 |
| All | 广播给所有 Task | 配置、规则更新 |
| Global | 全部发给单个 Task | 全局汇总 |
| Direct | 上游指定目标 Task | 精确路由 |
| Local or Shuffle | 优先本地传输 | 减少跨 Worker 通信 |
| None | 不关心分组 | 普通无状态处理 |
8. Storm 的可靠性机制
Storm 基础 API 主要提供 At-least-once,也就是“至少处理一次”语义。
为了理解它,需要先了解 Tuple Tree。
8.1 Tuple Tree
一个 Spout Tuple 可能产生多个下游 Tuple,每个下游 Tuple 又可能继续产生更多 Tuple。
Spout Tuple
├── Bolt A Tuple
│ ├── Bolt C Tuple
│ └── Bolt D Tuple
└── Bolt B Tuple
└── Bolt E Tuple
只有整棵 Tuple Tree 都处理完成,最初的 Spout Tuple 才算成功。
8.2 Message ID
可靠 Spout 在发送 Tuple 时需要指定 Message ID:
collector.emit(
new Values(orderId, amount),
messageId
);
没有 Message ID 的 Tuple 不会被 Storm 可靠性机制完整跟踪。
8.3 Anchoring
Bolt 发送新 Tuple 时,需要将它与输入 Tuple 关联:
collector.emit(
input,
new Values(orderId, amount)
);
这里的 input 就是 Anchor。
如果没有 Anchor,新产生的 Tuple 不会被加入原始 Tuple Tree,可能出现下游尚未处理完成,上游却已经收到成功通知的问题。
8.4 ACK
Bolt 成功处理 Tuple 后,需要调用:
collector.ack(input);
Storm 收到所有相关 ACK 后,会通知 Spout:
@Override
public void ack(Object messageId) {
// 消息完整处理成功
}
8.5 FAIL
处理失败时调用:
collector.fail(input);
Spout 随后会收到:
@Override
public void fail(Object messageId) {
// 根据 messageId 重新发送消息
}
8.6 超时机制
每个拓扑都有消息处理超时时间:
config.setMessageTimeoutSecs(30);
如果原始 Tuple 在规定时间内没有完整处理成功,Storm 会将它标记为失败。
超时时间不能随意设置:
- 太短会导致正常慢请求被重复处理;
- 太长会导致真实失败恢复过慢;
- 应结合下游接口、数据库和批量写入的 P99 耗时设置;
- 外部调用本身也必须配置连接超时和读取超时。
8.7 常见可靠性错误
发送新 Tuple 时没有 Anchor
错误:
collector.emit(new Values(result));
collector.ack(input);
正确:
collector.emit(input, new Values(result));
collector.ack(input);
忘记调用 ACK 或 FAIL
如果 Bolt 既不 ACK 也不 FAIL,消息只能等待超时,最终造成大量重试和延迟。
先 ACK 再执行副作用
错误:
collector.ack(input);
databaseService.save(data);
数据库写入失败时,Storm 已经认为消息处理成功。
应该在关键操作成功后再 ACK:
databaseService.save(data);
collector.ack(input);
认为 ACK 可以回滚外部数据库
Storm 的 ACK 只能表示计算拓扑中的处理状态,不能自动回滚已经写入的 MySQL、Redis 或外部接口。
9. Storm 的消息处理语义
9.1 At-most-once
At-most-once 表示消息最多处理一次,失败后不重试。
优点:
- 实现简单;
- 不会因为重试产生重复处理。
缺点:
- 发生异常时可能丢失数据。
适合允许少量丢失的非关键指标。
9.2 At-least-once
At-least-once 表示消息至少处理一次。
优点:
- 消息不容易丢失;
- 是 Storm 基础可靠性机制的核心语义。
缺点:
- 失败重试可能产生重复处理。
例如:
- Bolt 成功写入数据库;
- 还没来得及 ACK,Worker 就崩溃;
- Storm 认为消息失败;
- Spout 重新发送消息;
- 数据库再次收到相同写入。
因此,业务侧必须考虑幂等性。
9.3 Exactly-once
Storm 的 Trident API 可以提供 Exactly-once Processing Semantics,但“框架语义”不自动等于整个业务链路绝对只执行一次。
最终结果是否重复,还取决于外部存储是否支持:
- 事务;
- 幂等写入;
- 批次提交;
- 版本控制;
- 唯一约束;
- 原子更新。
9.4 常用幂等方案
业务唯一键
ALTER TABLE t_order_result
ADD UNIQUE KEY uk_order_id (order_id);
相同订单重复写入时,由唯一约束阻止重复数据。
Upsert
INSERT INTO t_word_count(word, total)
VALUES ('storm', 1)
ON DUPLICATE KEY UPDATE
total = total + 1;
注意:这个例子本身并不能避免重复消息导致重复累加。若需要严格幂等,还应同时记录消息 ID 或消费版本。
消息处理记录表
CREATE TABLE t_processed_message (
message_id VARCHAR(128) PRIMARY KEY,
processed_at DATETIME NOT NULL
);
在同一个数据库事务中:
- 插入消息处理记录;
- 更新业务数据;
- 提交事务。
如果消息 ID 已经存在,则跳过重复处理。
状态版本号
按照 Kafka Partition 和 Offset、业务事件版本或批次 ID 控制重复写入。
幂等接口
调用外部接口时传递稳定的幂等键:
Idempotency-Key: order-10001-refund
10. Storm 3.0.0 环境准备
本文示例使用:
- Apache Storm 3.0.0;
- Java 25;
- Maven 3.9 或兼容版本。
检查 Java:
java -version
javac -version
检查 Maven:
mvn -version
检查 Storm:
storm version
Storm 3.0.0 要求 Java 25 或更高版本。
如果现有环境只能使用 Java 11、17 或 21,应先评估运行环境和依赖兼容性,不要直接把示例版本机械替换后上线。Storm 2.8.9 是 2.x 最终版本,但 2.x 已结束维护。
11. WordCount 完整开发案例
下面实现一个实时单词统计拓扑:
SentenceSpout
↓ Shuffle Grouping
SplitSentenceBolt
↓ Fields Grouping(word)
WordCountBolt
11.1 项目结构
storm-word-count/
├── pom.xml
└── src/
└── main/
└── java/
└── com/
└── example/
└── storm/
└── WordCountTopology.java
11.2 Maven 配置
创建 pom.xml:
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="
http://maven.apache.org/POM/4.0.0
https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.example</groupId>
<artifactId>storm-word-count</artifactId>
<version>1.0.0</version>
<properties>
<maven.compiler.release>25</maven.compiler.release>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<storm.version>3.0.0</storm.version>
</properties>
<dependencies>
<dependency>
<groupId>org.apache.storm</groupId>
<artifactId>storm-client</artifactId>
<version>${storm.version}</version>
<scope>provided</scope>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.14.1</version>
<configuration>
<release>${maven.compiler.release}</release>
<encoding>${project.build.sourceEncoding}</encoding>
</configuration>
</plugin>
</plugins>
</build>
</project>
storm-client 使用 provided,表示编译时需要,但运行时由 Storm 集群提供。
如果项目还依赖数据库驱动、JSON 工具或其他第三方库,需要根据集群类路径策略将这些依赖打进应用包,或者通过 Storm CLI 的 --jars、--artifacts 参数提供。
11.3 完整 Java 代码
创建 WordCountTopology.java:
package com.example.storm;
import org.apache.storm.Config;
import org.apache.storm.StormSubmitter;
import org.apache.storm.spout.SpoutOutputCollector;
import org.apache.storm.task.OutputCollector;
import org.apache.storm.task.TopologyContext;
import org.apache.storm.topology.OutputFieldsDeclarer;
import org.apache.storm.topology.TopologyBuilder;
import org.apache.storm.topology.base.BaseRichBolt;
import org.apache.storm.topology.base.BaseRichSpout;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.Values;
import org.apache.storm.utils.Utils;
import java.util.HashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
public class WordCountTopology {
public static class SentenceSpout extends BaseRichSpout {
private final List<String> sentences = List.of(
"Apache Storm processes streams",
"Storm is distributed and reliable",
"Storm supports realtime computation",
"Apache Storm is scalable"
);
private transient SpoutOutputCollector collector;
private transient Map<Long, String> pending;
private long sequence;
@Override
public void open(
Map<String, Object> topologyConfig,
TopologyContext context,
SpoutOutputCollector collector
) {
this.collector = collector;
this.pending = new HashMap<>();
}
@Override
public void nextTuple() {
// 限制正在等待 ACK 的消息数量,避免无限发送。
if (pending.size() >= 1_000) {
Utils.sleep(10);
return;
}
long messageId = sequence++;
String sentence = sentences.get(
(int) (messageId % sentences.size())
);
pending.put(messageId, sentence);
collector.emit(
new Values(sentence),
messageId
);
// 示例中主动降低发送速度,避免本地运行时刷屏。
Utils.sleep(100);
}
@Override
public void ack(Object messageId) {
if (messageId instanceof Long id) {
pending.remove(id);
}
}
@Override
public void fail(Object messageId) {
if (!(messageId instanceof Long id)) {
return;
}
String sentence = pending.get(id);
if (sentence != null) {
// 使用相同 Message ID 重新发送原消息。
collector.emit(
new Values(sentence),
id
);
}
}
@Override
public void declareOutputFields(
OutputFieldsDeclarer declarer
) {
declarer.declare(new Fields("sentence"));
}
}
public static class SplitSentenceBolt extends BaseRichBolt {
private transient OutputCollector collector;
@Override
public void prepare(
Map<String, Object> topologyConfig,
TopologyContext context,
OutputCollector collector
) {
this.collector = collector;
}
@Override
public void execute(Tuple input) {
try {
String sentence = input.getStringByField("sentence");
for (String word : sentence.split("\\s+")) {
String normalizedWord = word
.trim()
.toLowerCase(Locale.ROOT);
if (!normalizedWord.isEmpty()) {
// 将新 Tuple 锚定到输入 Tuple。
collector.emit(
input,
new Values(normalizedWord)
);
}
}
collector.ack(input);
} catch (Exception exception) {
collector.fail(input);
}
}
@Override
public void declareOutputFields(
OutputFieldsDeclarer declarer
) {
declarer.declare(new Fields("word"));
}
}
public static class WordCountBolt extends BaseRichBolt {
private transient OutputCollector collector;
private transient Map<String, Long> counts;
@Override
public void prepare(
Map<String, Object> topologyConfig,
TopologyContext context,
OutputCollector collector
) {
this.collector = collector;
this.counts = new HashMap<>();
}
@Override
public void execute(Tuple input) {
try {
String word = input.getStringByField("word");
long total = counts.merge(word, 1L, Long::sum);
System.out.printf(
"task=%d, word=%s, total=%d%n",
input.getSourceTask(),
word,
total
);
collector.ack(input);
} catch (Exception exception) {
collector.fail(input);
}
}
@Override
public void declareOutputFields(
OutputFieldsDeclarer declarer
) {
// 当前 Bolt 不再向下游发送 Tuple。
}
}
public static void main(String[] args) throws Exception {
String topologyName = args.length > 0
? args[0]
: "word-count-topology";
TopologyBuilder builder = new TopologyBuilder();
builder.setSpout(
"sentence-spout",
new SentenceSpout(),
1
);
builder.setBolt(
"split-bolt",
new SplitSentenceBolt(),
2
).shuffleGrouping("sentence-spout");
builder.setBolt(
"count-bolt",
new WordCountBolt(),
2
).fieldsGrouping(
"split-bolt",
new Fields("word")
);
Config config = new Config();
config.setDebug(false);
config.setNumWorkers(2);
config.setMessageTimeoutSecs(30);
config.setMaxSpoutPending(1_000);
StormSubmitter.submitTopology(
topologyName,
config,
builder.createTopology()
);
}
}
11.4 为什么 CountBolt 使用 Fields Grouping
下面的代码按 word 分组:
.fieldsGrouping(
"split-bolt",
new Fields("word")
);
这样可以保证相同单词进入同一个 Task。
例如:
storm → Task 1
apache → Task 2
storm → Task 1
java → Task 3
apache → Task 2
如果使用 Shuffle Grouping,同一个单词可能进入不同 Task,而每个 Task 又维护自己的内存 Map,最终会得到多个局部统计结果。
11.5 示例状态的局限性
示例中的 counts 保存在 JVM 内存中:
private transient Map<String, Long> counts;
它只适合演示,不能直接作为生产级状态方案,因为:
- Worker 重启后状态会丢失;
- Rebalance 可能导致状态重新分布;
- 多个 Task 只维护各自的局部状态;
- 无法直接提供全局一致查询;
- 消息重试可能造成重复累加。
生产环境应根据业务要求使用:
- Redis;
- MySQL;
- HBase;
- Cassandra;
- 状态检查点;
- Trident State;
- 支持幂等或事务写入的外部存储。
12. Storm 拓扑的打包与运行
12.1 编译打包
在项目根目录执行:
mvn clean package
生成:
target/storm-word-count-1.0.0.jar
12.2 本地模式运行
Storm 3.0.0 提供 storm local 命令:
storm local \
target/storm-word-count-1.0.0.jar \
com.example.storm.WordCountTopology \
local-word-count
本地模式会在当前进程中启动嵌入式 Storm 环境,适合开发和验证。官方命令默认运行一段时间后自动退出,不应将本地模式当作生产集群。
12.3 提交到远程集群
storm jar \
target/storm-word-count-1.0.0.jar \
com.example.storm.WordCountTopology \
word-count-topology
参数含义:
storm jar
应用 Jar 路径
Main Class
传递给 main() 的拓扑名称
12.4 查看拓扑
storm list
12.5 停用拓扑的数据源
storm deactivate word-count-topology
Deactivate 会停止 Spout 继续发送新数据,但正在处理的数据仍可以继续完成。
12.6 恢复拓扑
storm activate word-count-topology
12.7 动态调整并行度
storm rebalance word-count-topology \
-w 30 \
-n 4 \
-e sentence-spout=2 \
-e split-bolt=6 \
-e count-bolt=8
其中:
-w 30:等待 30 秒;-n 4:调整为 4 个 Worker;-e component=count:调整组件 Executor 数量。
Executor 数量不能超过该组件预先设置的 Task 数量。
12.8 停止拓扑
storm kill word-count-topology -w 30
Storm 会先停用 Spout,等待正在处理的消息完成,再终止 Worker。
不要在生产环境中直接杀死 Worker JVM 来代替正常的拓扑停止流程。
13. Storm 集群部署
13.1 基本节点规划
一个简单集群可以规划为:
| 节点 | 服务 |
|---|---|
| node-1 | Nimbus、UI |
| node-2 | Nimbus、Supervisor |
| node-3 | Supervisor |
| node-4 | Supervisor |
| zk-1、zk-2、zk-3 | ZooKeeper |
生产环境应根据故障域、资源和访问压力决定是否将 Nimbus、UI、ZooKeeper 与 Worker 分开部署。
13.2 基础配置
编辑:
$STORM_HOME/conf/storm.yaml
示例:
storm.zookeeper.servers:
- "zk-1.example.com"
- "zk-2.example.com"
- "zk-3.example.com"
storm.zookeeper.port: 2181
nimbus.seeds:
- "node-1.example.com"
- "node-2.example.com"
storm.local.dir: "/data/storm"
supervisor.slots.ports:
- 6700
- 6701
- 6702
- 6703
ui.port: 8080
logviewer.port: 8000
13.3 启动 Nimbus
storm nimbus
13.4 启动 Supervisor
在每个工作节点执行:
storm supervisor
13.5 启动 UI
storm ui
访问:
http://storm-ui-host:8080
13.6 启动 Logviewer
storm logviewer
13.7 使用进程管理器
Storm 守护进程应由 systemd、Supervisor、Kubernetes 或其他进程管理系统托管。
不要只在终端中手动启动,因为终端退出、进程崩溃或服务器重启后,服务可能无法恢复。
13.8 Full 与 Lite 发行包
Storm 3.0.0 提供两种发行包:
- Full:包含更多可选集成依赖;
- Lite:体积更小,不包含部分 Kafka、Hadoop 等可选组件。
如果使用 Lite 版本,需要按实际情况补充相关依赖。选择发行包前,应先确认拓扑使用的 Kafka、HDFS、HBase、Kerberos 和监控组件。
14. Storm 与 Kafka 集成
Storm 常用 storm-kafka-client 从 Kafka 读取消息。
典型链路如下:
Kafka Topic
↓
KafkaSpout
↓
ParseBolt
↓
BusinessBolt
↓
DatabaseBolt
14.1 Kafka 与 Storm 的职责
Kafka 负责:
- 消息持久化;
- Topic 和 Partition;
- 消息顺序;
- 消费位点;
- 消息回放。
Storm 负责:
- 并行计算;
- 数据转换;
- 分组和聚合;
- 多级处理;
- 失败跟踪;
- 结果输出。
14.2 Partition 与并行度
如果 Kafka Topic 只有 4 个 Partition,即使 KafkaSpout 设置了大量并行 Task,也不一定能获得等比例吞吐提升。
设计时需要联合考虑:
- Kafka Partition 数量;
- Spout Executor 数量;
- Spout Task 数量;
- 下游 Bolt 并行度;
- 单条消息处理耗时;
- 下游数据库吞吐能力。
14.3 消费位点与 ACK
Storm 对 Tuple 的 ACK 与 Kafka Offset 提交需要由 KafkaSpout 正确协调。
业务 Bolt 成功 ACK 并不意味着业务写入天然具备 Exactly-once。数据库侧仍要处理:
- 重复消息;
- 事务失败;
- 超时重试;
- Worker 在写入成功后、ACK 前崩溃;
- Kafka 消息重新消费。
14.4 反压
当下游 Bolt 处理速度低于 KafkaSpout 的读取速度时,会产生积压和反压。
处理方向包括:
- 增加下游 Bolt 并行度;
- 优化数据库写入;
- 使用批量写入;
- 控制 Spout Pending 数量;
- 检查数据倾斜;
- 降低单条 Tuple 大小;
- 优化跨 Worker 网络传输;
- 扩容 Worker 和 Supervisor;
- 确认真正瓶颈是否在外部系统。
不要只增加 Spout 并行度。下游已经饱和时,提高读取速度只会加重积压。
15. Storm 监控与故障排查
15.1 常见指标
Storm UI 中应重点关注:
| 指标 | 说明 |
|---|---|
| Emitted | 组件发送的 Tuple 数量 |
| Transferred | 实际发送到其他 Task 的数量 |
| Acked | 成功处理的 Tuple 数量 |
| Failed | 失败的 Tuple 数量 |
| Complete Latency | Spout Tuple 完整处理耗时 |
| Execute Latency | Bolt 执行 Tuple 的平均耗时 |
| Process Latency | Bolt 收到 Tuple 到 ACK 的耗时 |
| Capacity | Bolt 执行线程的忙碌程度 |
| Error | 组件最近发生的异常 |
需要注意,Storm 部分内置指标可能采用采样后估算,不能在不了解采样率的情况下把 UI 数字当成绝对精确的业务计数。
15.2 Capacity
Capacity 可以帮助判断 Bolt 是否接近饱和。
经验上:
- Capacity 较低:Bolt 仍有余量;
- Capacity 持续升高:处理开始接近上限;
- Capacity 长期接近或超过 1:Bolt 很可能已经成为瓶颈。
但 Capacity 不是唯一依据,还要同时观察:
- 上游发送速度;
- 下游 ACK 速度;
- Pending 数量;
- 队列积压;
- CPU;
- GC;
- 网络;
- 外部数据库耗时。
15.3 Failed 持续增加
可能原因包括:
- Bolt 主动调用
fail(); - 消息处理超时;
- Worker 重启;
- 数据格式异常;
- 数据库连接失败;
- 外部接口超时;
- 业务代码抛出异常;
- Tuple Tree 没有正确 ACK;
- Spout 重放消息。
排查顺序建议:
- 查看 Storm UI 的组件错误;
- 定位出错的 Bolt 或 Spout;
- 查看对应 Worker 日志;
- 通过消息 ID、订单号等业务键定位原始数据;
- 检查下游依赖;
- 检查是否发生 Worker 重启;
- 检查超时时间;
- 检查 ACK、FAIL 和 Anchoring。
15.4 Complete Latency 变高
常见原因:
- Bolt 处理变慢;
- 数据库慢查询;
- 外部接口超时;
- 网络抖动;
- GC 停顿;
- 数据倾斜;
- Worker CPU 饱和;
- Tuple 过大;
- 下游并行度不足;
- 批量数据迟迟没有达到刷新条件。
15.5 Spout Pending 达到上限
如果未完成的 Tuple 数量达到:
topology.max.spout.pending
Storm 会限制 Spout 继续发送数据。
这通常说明下游处理速度不足,或者 ACK 链路存在问题。
15.6 常用排查命令
查看拓扑:
storm list
获取拓扑最近错误:
storm get-errors topology-name
监控吞吐:
storm monitor topology-name
临时调整日志级别:
storm set_log_level \
-l ROOT=DEBUG:60 \
topology-name
恢复日志配置:
storm set_log_level \
-r ROOT \
topology-name
16. Storm 性能调优
16.1 先定位瓶颈
不要一开始就盲目增加 Worker。
应该先回答:
- Spout 是否读不动;
- 哪个 Bolt Capacity 最高;
- 是否存在 Fields Grouping 数据倾斜;
- CPU 是否饱和;
- GC 是否频繁;
- 网络是否成为瓶颈;
- 数据库是否出现慢查询;
- 外部接口是否变慢;
- Tuple 是否过大;
- 是否存在大量序列化和反序列化。
16.2 调整 Worker 数量
config.setNumWorkers(4);
增加 Worker 可以获得更多 JVM 和进程级资源,但也会增加:
- JVM 内存;
- 进程管理成本;
- 网络传输;
- 序列化开销;
- Worker 之间的数据交换。
如果服务器 CPU 已经饱和,单纯增加 Worker 不一定有效。
16.3 调整 Bolt 并行度
builder.setBolt(
"business-bolt",
new BusinessBolt(),
8
);
适合 CPU 密集型或单条处理较慢的 Bolt。
但是,如果瓶颈是数据库连接池或下游限流,增加 Bolt 并行度可能导致数据库压力进一步升高。
16.4 处理数据倾斜
Fields Grouping 出现热点 Key 时,可以考虑:
- 增加 Key 粒度;
- 给热点 Key 添加前缀;
- 先局部聚合,再全局合并;
- 使用 Partial Key Grouping;
- 将热点业务单独拆分;
- 对热点 Key 单独限流。
例如,把单阶段统计:
userId → CountBolt
改成两阶段:
userId + randomPrefix
↓
LocalCountBolt
↓
userId
↓
MergeCountBolt
16.5 控制 Spout Pending
config.setMaxSpoutPending(5_000);
值过小:
- Spout 容易被限制;
- 吞吐量可能不足。
值过大:
- 内存占用增加;
- 在途消息过多;
- 延迟可能升高;
- 失败时重放压力更大。
应通过压测寻找合适值。
16.6 减少 Tuple 大小
不要在 Tuple 中传递不需要的大对象。
不推荐:
完整订单对象 + 用户对象 + 商品列表 + 原始 JSON
推荐:
orderId + userId + amount + eventType
大 Tuple 会增加:
- 序列化开销;
- 网络带宽;
- Worker 队列内存;
- GC 压力;
- 重试成本。
Storm 3.0.0 支持对跨 Worker Tuple 流量选择性启用 Zstandard 压缩,但压缩会消耗 CPU,适合大载荷和网络带宽敏感场景,不应对所有小 Tuple 无差别开启。
16.7 优化序列化
生产环境应:
- 优先使用 Storm 原生支持的数据类型;
- 为自定义类型注册 Kryo 序列化器;
- 避免频繁创建巨大临时对象;
- 避免依赖 Java 原生序列化;
- 保持上下游字段协议稳定。
16.8 批量写数据库
逐条写数据库:
1 Tuple → 1 次 INSERT
可能导致数据库连接和网络开销过高。
可以改为:
100 条 Tuple → 1 次批量 INSERT
但批量写入需要额外处理:
- 批次大小;
- 最大等待时间;
- 批量失败;
- 部分成功;
- ACK 时机;
- Worker 异常;
- 消息重复;
- 关闭或 Rebalance 时的数据刷新。
不能在数据还没有可靠写入时提前 ACK。
16.9 避免阻塞 Spout
nextTuple() 不应执行长时间阻塞操作。
错误示例:
@Override
public void nextTuple() {
remoteService.waitForMessageForever();
}
如果数据源读取会阻塞,可以:
- 使用带超时的 Poll;
- 将数据先放入本地缓冲队列;
- 使用专门的客户端线程;
- 控制缓冲区上限;
- 正确处理中断和关闭。
16.10 合理设置消息超时
config.setMessageTimeoutSecs(60);
应满足:
消息超时
>
正常链路 P99 处理耗时
+
合理的网络和调度余量
不要通过无限增大消息超时掩盖慢查询、线程阻塞或下游故障。
16.11 关闭逐条 Debug 日志
下面的配置会产生大量日志:
config.setDebug(true);
生产环境通常应关闭:
config.setDebug(false);
也不要为每一条高频消息记录完整 INFO 日志。可以使用:
- 指标统计;
- 采样日志;
- 聚合日志;
- 只记录异常;
- 按业务 ID 定向调试。
17. Storm 生产实践与常见问题
17.1 Bolt 是否线程安全
一个 Executor 是一个线程,一个 Executor 可能运行同一组件的一个或多个 Task。
不要随意使用全局静态可变对象:
private static final Map<String, Long> COUNTS = new HashMap<>();
它可能造成:
- 多 Task 共享状态;
- 并发安全问题;
- 数据相互污染;
- 测试和生产行为不一致。
组件实例状态应明确限定作用域,外部共享状态应放入可靠存储。
17.2 prepare 和 cleanup
Bolt 通常在 prepare() 中初始化资源:
@Override
public void prepare(
Map<String, Object> config,
TopologyContext context,
OutputCollector collector
) {
this.collector = collector;
this.dataSource = createDataSource();
}
但是不能把所有资源释放完全依赖于 cleanup(),因为进程崩溃、服务器宕机或被强制终止时,cleanup() 不一定有机会执行。
17.3 数据库连接
不要为每条 Tuple 创建一个新连接:
Connection connection =
DriverManager.getConnection(url, username, password);
应该使用:
- 连接池;
- 合理的最大连接数;
- 获取连接超时;
- SQL 执行超时;
- 失败重试上限;
- 熔断和降级;
- 指标监控。
17.4 无限重试
如果消息本身格式错误,无限重试没有意义,还会持续占用资源。
建议区分:
- 临时故障:有限次数重试;
- 永久性数据错误:发送到死信队列;
- 下游不可用:退避、限流、熔断;
- 未知异常:记录上下文并告警。
17.5 状态丢失
Bolt 内存状态可能在以下情况丢失:
- Worker 重启;
- 服务器宕机;
- Rebalance;
- 代码升级;
- Task 重新分配。
关键状态必须外置或使用专门的状态管理机制。
17.6 配置热更新
如果业务规则需要热更新,可以使用:
- 配置中心;
- 广播 Stream;
- 定时刷新;
- 版本化配置;
- 原子替换本地快照。
更新时要考虑:
- 多 Task 生效时间不一致;
- 配置加载失败;
- 旧配置回滚;
- 配置版本审计;
- 并发读取安全。
17.7 拓扑升级
常见升级方式是:
- 提交新版本拓扑;
- 验证新拓扑运行状态;
- 切换上游流量;
- 等待旧拓扑完成在途消息;
- 停止旧拓扑。
如果直接使用相同名称覆盖或先杀后启动,要评估:
- 消费位点;
- 状态迁移;
- 数据重复;
- 短暂中断;
- 外部存储兼容性;
- Tuple Schema 兼容性。
17.8 安全配置
Storm 默认环境不能直接视为安全生产配置。
生产环境至少应考虑:
- Nimbus 身份认证;
- 操作权限控制;
- ZooKeeper ACL;
- Kerberos 或其他认证机制;
- UI 和 Logviewer 访问控制;
- 防火墙;
- TLS;
- 反向代理;
- 最小权限系统账号;
- 密钥和密码外置;
- 审计日志。
Storm UI 可以执行 Kill、Activate、Deactivate、Rebalance 等状态修改操作,因此不应直接暴露到公网。
17.9 生产上线检查清单
- Spout 可以正确重放失败消息;
- 每个 Bolt 都正确调用 ACK 或 FAIL;
- 下游 Tuple 已正确 Anchoring;
- 外部副作用具备幂等性;
- 数据库操作具有超时;
- 消息重试次数有上限;
- 无法处理的消息进入死信队列;
- 关键状态没有只保存在内存;
- 已设置合理的并行度;
- 已验证 Fields Grouping 是否数据倾斜;
- 已配置 Worker 内存;
- 已监控 CPU、内存和 GC;
- 已监控 ACK、FAIL 和处理延迟;
- 已进行故障和重复消息测试;
- 已进行容量压测;
- UI 和 Logviewer 已设置访问控制;
- 已准备回滚方案。
18. Storm 与其他流处理框架对比
18.1 Storm 与 Flink
| 对比项 | Storm | Flink |
|---|---|---|
| 核心模型 | Tuple 级实时处理 | 有状态流处理 |
| API 风格 | Spout、Bolt、Topology | DataStream、ProcessFunction、SQL |
| 状态管理 | 需要开发者明确设计 | 状态和 Checkpoint 能力更完整 |
| 基础处理语义 | At-least-once | 支持 At-least-once 和 Exactly-once |
| 事件时间 | 可以实现窗口等能力 | 事件时间、水位线体系更完整 |
| 学习成本 | 核心模型直接 | 功能更多,概念也更多 |
| 适用场景 | 低延迟事件处理、已有 Storm 体系 | 状态计算、窗口、复杂流式分析 |
如果项目高度依赖窗口、事件时间、大规模状态和端到端 Checkpoint,通常应重点评估 Flink。
如果已有成熟 Storm 集群,业务以简单、低延迟的事件处理为主,则没有必要为了“框架更新”盲目迁移。
18.2 Storm 与 Spark Structured Streaming
| 对比项 | Storm | Spark Structured Streaming |
|---|---|---|
| 处理方式 | 原生流式 Tuple 处理 | 主要采用微批处理模型 |
| 延迟特征 | 适合低延迟逐事件处理 | 更适合统一批流和数据分析 |
| SQL 能力 | 相对有限 | 与 Spark SQL 深度结合 |
| 生态 | 实时计算 | 大数据、机器学习、SQL 生态 |
| 典型用途 | 实时事件处理 | 实时 ETL、分析和统一数据平台 |
18.3 Storm 与 Kafka Streams
| 对比项 | Storm | Kafka Streams |
|---|---|---|
| 部署方式 | 独立计算集群 | 嵌入普通应用 |
| 数据源 | 可以连接多种系统 | 主要围绕 Kafka |
| 调度能力 | Nimbus 和 Supervisor | 依赖应用实例和 Kafka 分区 |
| 运维复杂度 | 需要维护 Storm 集群 | 不需要独立计算集群 |
| 适用场景 | 多级分布式计算拓扑 | Kafka 内部流处理应用 |
如果数据全部来自 Kafka,处理流程相对简单,又希望把流处理嵌入微服务,可以优先评估 Kafka Streams。
18.4 如何选择
不要只比较框架理论性能,还应考虑:
- 团队经验;
- 已有基础设施;
- 状态规模;
- 延迟要求;
- Exactly-once 要求;
- SQL 和窗口需求;
- 运维成本;
- 社区和版本生命周期;
- 数据源与目标存储;
- 迁移成本。
19. Storm 3.0.0 的主要变化
Storm 3.0.0 是一个大版本,主要变化包括:
19.1 Java 25
Storm 3.0.0 要求 Java 25 或更高版本。
升级前应检查:
- 操作系统和容器镜像;
- JVM 参数;
- 监控 Agent;
- 日志组件;
- 数据库驱动;
- 自定义序列化器;
- 第三方 Storm Connector;
- 构建和发布环境。
19.2 移除 Clojure 支持
Storm 3.0.0 移除了 Clojure 依赖和 Clojure DSL。
使用 Java API 的 Storm 2.x 拓扑通常可以相对平滑地迁移,但使用 Clojure API 的项目需要改写。
19.3 Java API 兼容性
官方说明 Java API 与 Storm 2.x 保持兼容,现有 Java 拓扑通常不需要因为 API 本身进行大规模改造。
不过,升级仍然需要完成完整的依赖、配置、序列化、运行环境和回归验证。
19.4 Zstandard 压缩
Storm 3.0.0 可以对跨 Worker Tuple 流量和部分集群状态启用 Zstandard 压缩。
它适合:
- Tuple 较大;
- 跨 Worker 网络流量高;
- 网络带宽是主要瓶颈;
- CPU 仍有余量。
对于很小的 Tuple,压缩头和 CPU 成本可能高于节省的网络开销。
19.5 动态批量大小
Storm 3.0.0 引入基于 AIMD 的动态批量调整能力,可以在反压场景中根据运行状态动态调整生产者批量大小。
它属于需要结合实际负载验证的优化项,不能替代容量规划和瓶颈分析。
19.6 新的分组和调度能力
Storm 3.0.0 增加了 JitterAwareStreamGrouping,并改进了调度器行为。
这些能力适合对任务抖动、节点恢复和负载均衡有更高要求的场景,但生产启用前仍应进行故障演练和性能测试。
19.7 Full 与 Lite 发行包
3.0.0 同时提供完整包和精简包。
精简包减少了可选集成依赖,可以降低下载和部署体积,但需要团队明确补充实际使用的 Kafka、Hadoop、HBase 或安全组件。
20. 总结
Apache Storm 使用 Spout、Bolt、Stream 和 Topology 描述一个持续运行的实时计算流程。
它的核心执行链路可以概括为:
数据源
↓
Spout 读取并发送 Tuple
↓
Grouping 决定数据分区
↓
Bolt 执行计算
↓
外部系统保存结果
↓
ACK 或 FAIL
学习和使用 Storm 时,最需要掌握以下几点:
- Topology 是长期运行的实时计算图;
- Spout 负责读取数据,Bolt 负责处理数据;
- Grouping 决定数据如何分区;
- Worker 是 JVM,Executor 是线程,Task 是实际处理实例;
- Storm 基础 API 主要提供 At-least-once 语义;
- ACK 不能替代数据库事务和业务幂等;
- Fields Grouping 需要关注数据倾斜;
- 内存状态可能因 Worker 重启或 Rebalance 丢失;
- 增加并行度前应先定位真正的性能瓶颈;
- 生产系统必须完善超时、重试、死信、监控、安全和回滚机制。
Storm 的核心模型并不复杂,但真正的生产难点通常不在于编写一个 Bolt,而在于正确处理消息重复、状态恢复、数据倾斜、外部系统故障和容量变化。
参考资料
更多推荐



所有评论(0)