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 可以从消息队列读取原始数据,然后执行:

  1. 字段校验;
  2. 数据清洗;
  3. 格式转换;
  4. 无效数据过滤;
  5. 数据补充;
  6. 多流关联;
  7. 写入目标数据库。

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 等组件组成。

Storm 客户端

Nimbus

ZooKeeper

Supervisor 1

Supervisor 2

Worker JVM

Worker JVM

Worker JVM

Worker JVM

Executor 线程

Spout/Bolt Task

Kafka 等数据源

数据库、缓存或搜索引擎

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

一次完整处理过程如下:

  1. OrderSpout 从 Kafka 读取订单消息;
  2. Spout 将消息封装成 Tuple;
  3. ParseOrderBolt 将 JSON 转换成订单字段;
  4. ValidOrderBolt 校验订单是否合法;
  5. merchantId 对消息重新分组;
  6. MerchantAmountBolt 统计每个商户的金额;
  7. 结果写入外部存储;
  8. Bolt 调用 ack() 告诉 Storm 当前 Tuple 处理完成;
  9. 整棵 Tuple Tree 完成后,Spout 收到成功通知;
  10. 如果处理超时或失败,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 基础可靠性机制的核心语义。

缺点:

  • 失败重试可能产生重复处理。

例如:

  1. Bolt 成功写入数据库;
  2. 还没来得及 ACK,Worker 就崩溃;
  3. Storm 认为消息失败;
  4. Spout 重新发送消息;
  5. 数据库再次收到相同写入。

因此,业务侧必须考虑幂等性。

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
);

在同一个数据库事务中:

  1. 插入消息处理记录;
  2. 更新业务数据;
  3. 提交事务。

如果消息 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 的读取速度时,会产生积压和反压。

处理方向包括:

  1. 增加下游 Bolt 并行度;
  2. 优化数据库写入;
  3. 使用批量写入;
  4. 控制 Spout Pending 数量;
  5. 检查数据倾斜;
  6. 降低单条 Tuple 大小;
  7. 优化跨 Worker 网络传输;
  8. 扩容 Worker 和 Supervisor;
  9. 确认真正瓶颈是否在外部系统。

不要只增加 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 重放消息。

排查顺序建议:

  1. 查看 Storm UI 的组件错误;
  2. 定位出错的 Bolt 或 Spout;
  3. 查看对应 Worker 日志;
  4. 通过消息 ID、订单号等业务键定位原始数据;
  5. 检查下游依赖;
  6. 检查是否发生 Worker 重启;
  7. 检查超时时间;
  8. 检查 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。

应该先回答:

  1. Spout 是否读不动;
  2. 哪个 Bolt Capacity 最高;
  3. 是否存在 Fields Grouping 数据倾斜;
  4. CPU 是否饱和;
  5. GC 是否频繁;
  6. 网络是否成为瓶颈;
  7. 数据库是否出现慢查询;
  8. 外部接口是否变慢;
  9. Tuple 是否过大;
  10. 是否存在大量序列化和反序列化。

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 拓扑升级

常见升级方式是:

  1. 提交新版本拓扑;
  2. 验证新拓扑运行状态;
  3. 切换上游流量;
  4. 等待旧拓扑完成在途消息;
  5. 停止旧拓扑。

如果直接使用相同名称覆盖或先杀后启动,要评估:

  • 消费位点;
  • 状态迁移;
  • 数据重复;
  • 短暂中断;
  • 外部存储兼容性;
  • 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 如何选择

不要只比较框架理论性能,还应考虑:

  1. 团队经验;
  2. 已有基础设施;
  3. 状态规模;
  4. 延迟要求;
  5. Exactly-once 要求;
  6. SQL 和窗口需求;
  7. 运维成本;
  8. 社区和版本生命周期;
  9. 数据源与目标存储;
  10. 迁移成本。

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 时,最需要掌握以下几点:

  1. Topology 是长期运行的实时计算图;
  2. Spout 负责读取数据,Bolt 负责处理数据;
  3. Grouping 决定数据如何分区;
  4. Worker 是 JVM,Executor 是线程,Task 是实际处理实例;
  5. Storm 基础 API 主要提供 At-least-once 语义;
  6. ACK 不能替代数据库事务和业务幂等;
  7. Fields Grouping 需要关注数据倾斜;
  8. 内存状态可能因 Worker 重启或 Rebalance 丢失;
  9. 增加并行度前应先定位真正的性能瓶颈;
  10. 生产系统必须完善超时、重试、死信、监控、安全和回滚机制。

Storm 的核心模型并不复杂,但真正的生产难点通常不在于编写一个 Bolt,而在于正确处理消息重复、状态恢复、数据倾斜、外部系统故障和容量变化。


参考资料

  1. Apache Storm 官方网站
  2. Apache Storm 3.0.0 官方文档
  3. Apache Storm 核心概念
  4. Apache Storm 消息可靠性机制
  5. Apache Storm 并行度模型
  6. Apache Storm 命令行工具
  7. Apache Storm 集群部署
  8. Apache Storm 安全配置
  9. Apache Storm 3.0.0 发布说明
  10. Apache Storm 下载页面

更多推荐