大数据实时计算Apache Storm
Apache Storm本质上是一个分布式流式计算引擎:持续接收无界数据流,实时完成过滤、转换、聚合、Join、状态计算,并把结果写入下游系统。官方目前文档对应 Storm 3.0.0。
1. 核心架构流程
1 | Kafka / MQ / 日志 / CDC |
三个最重要的概念:
| 概念 | 含义 |
|---|---|
| Spout | 数据源,负责持续产生Tuple |
| Bolt | 计算节点,负责处理数据 |
| Topology | Spout + Bolt构成的实时计算DAG |
Storm的 opology与传统MapReduce最大区别是:MapReduce Job 最终结束,而 Topology通常持续运行。
2. Storm 真正解决什么问题?
例如电商实时交易:
1 | 订单事件 |
每秒:
1 | 100 万订单事件 |
Storm官方定位就是对无界数据流进行实时、可靠、可扩展处理,典型用途包括实时分析、在线机器学习、持续计算、ETL 等。
3. Storm并行计算本质
Storm不是简单的“一个线程处理一条数据”。
例如:
1 | Kafka |
数据可以通过不同的 Stream Grouping 分发到不同 Task,从而实现水平扩展。这样可以保证:同一个 Key 的数据进入同一个计算分区。
Storm实时计算能力
现代Storm不只是简单的 Spout → Bolt。
Stream API 已经支持:
filtermapflatMap- window
- aggregation
- join
- state
- repartition
- output
例如:
1 | Kafka |
官方 Stream API 支持时间窗口、计数窗口、聚合、Join、状态计算以及重新分区。
Storm 的可靠性
实时计算最麻烦的问题之一不是算得快,而是:
机器挂了,数据怎么办?
Storm有Tuple 的 ACK / replay 机制。
逻辑可以理解成:
1 | Spout |
如果中间失败:
1 | Spout |
所以 Storm 的核心卖点之一就是容错和可靠的数据处理保证。
4. Storm vs Kafka
1 | Kafka |
- Kafka 是消息/事件流平台
- Storm 是流计算引擎
典型架构:
1 | 业务系统 |
Storm 官方也提供 Kafka Consumer 集成。
5. Storm vs Flink
| 维度 | Storm | Flink |
|---|---|---|
| 定位 | 实时流计算 | 流批一体计算 |
| 编程模型 | Spout/Bolt/Topology | DataStream / SQL |
| 低延迟 | 很强 | 很强 |
| Window | 支持 | 很强 |
| 状态计算 | 支持 | 很强 |
| Event Time | 支持 | 很强 |
| Exactly Once | Trident 等方案 | 原生能力更完整 |
| SQL | 有,但不是核心优势 | 核心能力 |
| 学习成本 | 较低 | 中等 |
| 经典实时计算 | 很适合 | 很适合 |
| 新项目主流选择 | 相对少 | 通常更常见 |