storm简介
Storm 核心概念:Spout 产数据,Bolt 处理数据,Topology 串起来,Stream 流着走
Storm 是实时流处理框架——数据源源不断地来,你源源不断地处理,毫秒级延迟。
它的核心模型很简单:
数据源 → Spout(采集) → Bolt(处理1) → Bolt(处理2) → 输出
所有的 Spout 和 Bolt 串成一个有向无环图,就是 Topology。数据在这个图里按 Stream 流动,Stream Grouping 决定数据流向哪个并行任务。
整体架构:Nimbus 管调度,Supervisor 管执行
Storm 集群两个角色:
| 角色 | 进程 | 做什么 |
|---|---|---|
| 管理节点 | Nimbus | 分配任务、监控集群、分发代码 |
| 工作节点 | Supervisor | 启动/停止工作进程,执行具体任务 |
关键设计:Nimbus 和 Supervisor 不直接通信,通过 ZooKeeper 协调。
- Nimbus 把任务信息写进 ZooKeeper
- Supervisor 从 ZooKeeper 读任务信息,启动执行
- Nimbus 挂了,Supervisor 继续跑;重启后从 ZooKeeper 恢复状态
这个设计的价值:Nimbus 挂了集群照跑,容错性极强。
Topology:Storm 的任务定义
Topology 是 Storm 里”一个实时计算任务”的定义——相当于 MapReduce 里的 Job,但 MapReduce Job 跑完就结束,Topology 一直跑。
结构: Spout + Bolt 组成的有向无环图(DAG)
Spout → Bolt1 → Bolt2 → Bolt3
↑ ↑
└─ 可以分支、合并
提交 Topology:
TopologyBuilder builder = new TopologyBuilder();
builder.setSpout("kafka-spout", new KafkaSpout(), 2);
builder.setBolt("filter-bolt", new FilterBolt(), 4).shuffleGrouping("kafka-spout");
builder.setBolt("count-bolt", new CountBolt(), 2).fieldsGrouping("filter-bolt", new Fields("user_id"));
并行度: setSpout(name, spout, parallelism_hint) 第三个参数是并行任务数,Storm 会启动多个线程并行执行。
Stream:数据流的抽象
Stream 是 Storm 里”数据流”的概念——无限持续的元组(Tuple)序列。
Tuple 类似数据库里的一行:
Tuple: {user_id: 123, action: "click", timestamp: 1600000000}
Stream 是无界的——数据源源不断进来,没有”结束”的概念。
Spout:数据的源头
Spout 从外部系统读数据,以 Tuple 形式发射到 Topology 里。
常见 Spout 来源:
- Kafka
- 数据库
- 日志文件
- Socket
Spout 有两种:
| 类型 | 特点 |
|---|---|
| 可靠 Spout | 记录 Tuple 状态,失败时重发(”至少一次”语义) |
| 不可靠 Spout | 发完不跟踪,可能丢数据(性能更高) |
生产上用可靠 Spout,数据不能丢。
Bolt:数据的处理器
Bolt 是 Storm 里真正干活的地方——接收 Tuple,处理,发射新 Tuple,或者写存储。
Bolt 做什么:
- 过滤、清洗
- 聚合、计数
- Join 多路数据
- 写数据库、发告警
一个 Bolt 可以订阅多个输入流:
builder.setBolt("join-bolt", new JoinBolt(), 4)
.fieldsGrouping("spout1", new Fields("id"))
.fieldsGrouping("spout2", new Fields("id"));
Stream Grouping:数据流向哪个并行任务
Bolt 有多个并行 Task,数据来了往哪个 Task 发?由 Grouping 决定。
| 分组类型 | 规则 | 场景 |
|---|---|---|
| Shuffle Grouping | 随机分发,均匀 | 负载均衡,每个 Task 处理差不多 |
| Fields Grouping | 按字段哈希,同值同 Task | 按用户 ID 聚合,同一用户的数据去同一个 Task |
| All Grouping | 广播给所有 Task | 分发字典表、配置数据 |
| Global Grouping | 全发到同一个 Task(最小 ID) | 全局聚合,比如算总数 |
| Local or Shuffle | 优先本地 Task,否则随机 | 减少网络传输 |
选择原则:
- 需要按 key 聚合 →
Fields Grouping - 纯并行计算(无依赖)→
Shuffle Grouping - 需要复制数据到所有节点 →
All Grouping - 全局汇总 →
Global Grouping