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