流处理全家桶:Redis Stream、Kafka Streams、Flink、Spark 原理与选型全解
本文详细解析当前主流的四种实时数据处理方案:Redis Stream、Kafka Streams、Apache Flink 和 Apache Spark。
文章从核心本质、执行流程、关键能力、状态管理机制、适用场景等角度逐层分析,并回答以下常见问题:
- Stream 到底是不是消息队列?
- 已经有 MQ 为什么还需要流处理?
- Kafka Streams 和 Flink 的区别是什么?
- Flink 和 Spark 应该如何选择?
最后给出生产环境选型指南,帮助开发者根据业务规模和实时性要求选择合适方案,而不是盲目追逐技术热点。
## 一、流处理技术的核心本质
传统批处理(Batch Processing)的特点是:
数据先存起来,积累到一定规模后统一计算。
例如:
- 每天凌晨统计订单数据
- 每小时生成一次报表
这种方式吞吐量高,但实时性较差。
而流处理(Stream Processing)的特点是:
数据产生后立即处理。
例如:
- 用户下单立即更新 GMV
- 支付后立即触发风控
- 设备异常立即报警
因此,流处理系统关注的是:
- 数据持续产生
- 数据持续计算
- 数据持续输出
核心目标是在数据流动过程中完成:
- 清洗(Filter)
- 转换(Map)
- 聚合(Aggregate)
- Join
- 窗口统计(Window)
- 状态计算(Stateful Processing)
需要注意:
Redis Stream、Kafka Streams、Flink、Spark 并不处于同一个层级。
它们分别属于:
| 技术 | 本质 |
|------|------|
| Redis Stream | 消息流数据结构(轻量 MQ) |
| Kafka Streams | 嵌入式流计算库 |
| Flink | 专业实时计算引擎 |
| Spark | 通用大数据计算引擎 |
## 二、Redis Stream:轻量级消息流
### 2.1 核心定义
Redis Stream 是 Redis 5.0 引入的数据结构。
它提供:
- 消息追加
- 消息消费
- 消费者组
- ACK确认
- 消息重放
因此它具备 MQ 的基本能力。
但要明确:
Redis Stream 是消息传输组件,不是流计算引擎。
它负责"搬运数据",不负责"计算数据"。
### 2.2 执行流程
**生产消息**
```
XADD order_stream * userId 1001 amount 99
```
Redis 自动生成消息 ID:
```
1681234567890-0
```
**消费消息**
```
XREADGROUP
```
消费者组拉取消息。
**ACK确认**
```
XACK
```
确认消费成功。
**消息保留**
未 ACK 消息会进入 PEL(Pending Entries List)。
支持:
- 重试
- 转移消费者
- 死信处理
### 2.3 核心特性
**优点**
- 极轻量
- 学习成本低
- 与 Redis 生态天然集成
- 延迟极低
**缺点**
- 无计算能力
- 内存成本高
- 不适合超大规模消息堆积
- 横向扩展能力有限
### 2.4 适用场景
适合:
- 异步通知
- 秒杀削峰
- 订单解耦
- 实时告警
- 任务队列
不适合:
- 实时 OLAP
- 风控计算
- 大规模状态计算
## 三、Kafka Streams:嵌入式流计算
### 3.1 核心定义
Kafka Streams 是 Kafka 官方提供的流处理库。
注意:
Kafka Streams 不是独立集群。
它只是一个 Java Library。
直接嵌入 Spring Boot 即可运行。
### 3.2 执行流程
```
Kafka Topic
↓
Kafka Streams
↓
Map Filter Join Window Aggregate
↓
Kafka Topic
```
应用启动后自动消费 Topic。
处理结果写回 Topic。
### 3.3 状态管理
Kafka Streams 最大特点:
内置状态管理。
默认使用:
RocksDB
存储本地状态。
同时通过:
Changelog Topic
同步到 Kafka。
节点重启后自动恢复状态。
### 3.4 核心特性
支持:
- Filter
- Map
- FlatMap
- GroupBy
- Count
- Aggregate
- Join
- Window
具备轻量级实时计算能力。
### 3.5 优缺点
**优点**
- 无需额外集群
- 开发简单
- 与 Kafka 深度整合
- 自动状态恢复
**缺点**
- 仅支持 JVM 生态
- CEP 能力较弱
- 大规模复杂计算不如 Flink
### 3.6 适用场景
适合:
- 实时计数
- 实时报表
- 数据脱敏
- 数据格式转换
- 用户行为统计
## 四、Apache Flink:专业实时计算引擎
### 4.1 核心定义
Flink 是目前最成熟的流计算框架之一。
设计理念:
Streaming First
即:
一切都是流。
即使批处理,本质上也是有界流(Bounded Stream)。
### 4.2 执行流程
```
Kafka
↓
Source
↓
Map
↓
KeyBy
↓
Window
↓
Aggregate
↓
Sink
```
作业提交后:
- JobManager 负责调度。
- TaskManager 负责计算。
### 4.3 状态管理
Flink 的核心竞争力之一就是状态管理。
支持:
- Keyed State
- Operator State
状态后端可以是:
- HashMapStateBackend
- EmbeddedRocksDBStateBackend
检查点(Checkpoint)定期持久化。
故障后自动恢复。
### 4.4 Exactly Once
通过:
- Checkpoint
- Barrier
- Two Phase Commit
实现 Exactly Once。
保证:
数据既不会丢,也不会重复计算。
### 4.5 核心能力
支持:
- Window
- Join
- CEP
- Watermark
- Event Time
- SQL
- CDC
几乎覆盖所有实时计算场景。
### 4.6 适用场景
适合:
- 实时风控
- 实时推荐
- 实时大屏
- IoT监控
- 用户行为分析
- 金融交易监控
## 五、Apache Spark:批流一体计算引擎
### 5.1 核心定义
Spark 最初是批处理框架。
后来发展出:
- Spark Streaming
- Structured Streaming
如今主流方案是:
Structured Streaming
实现批流统一。
### 5.2 核心特性
**批流一体**
同一套 API:
```
DataFrame Dataset
```
既能处理离线数据:
```
Parquet Hive
```
也能处理实时数据:
```
Kafka Pulsar
```
**丰富生态**
Spark 生态包含:
- Spark SQL
- MLlib
- GraphX
- Delta Lake
这是 Flink 难以匹敌的优势。
### 5.3 实时能力
Spark 本质仍然偏向微批处理(Micro Batch)。
虽然 Structured Streaming 已大幅降低延迟:
100ms ~ 数秒
但整体实时性仍弱于 Flink。
### 5.4 适用场景
适合:
- 数据仓库
- ETL
- 用户画像
- 离线分析
- 机器学习训练
- 准实时计算
## 六、全方案对比
| 维度 | Redis Stream | Kafka Streams | Flink | Spark |
|------|--------------|---------------|-------|-------|
| 本质 | 消息流数据结构 | 嵌入式流处理库 | 实时计算引擎 | 通用计算引擎 |
| 是否独立部署 | 否 | 否 | 是 | 是 |
| 状态管理 | 无 | RocksDB + Changelog | State Backend | Checkpoint |
| 实时性 | 极高 | 高 | 极高 | 中等 |
| 复杂计算 | 无 | 一般 | 极强 | 强 |
| CEP | 不支持 | 不支持 | 支持 | 有限 |
| SQL支持 | 无 | 有限 | 完整支持 | 完整支持 |
| 运维成本 | 极低 | 低 | 高 | 高 |
| 学习成本 | 低 | 中 | 高 | 中 |
## 七、生产环境选型指南
**小团队**
```
Redis Stream + 业务代码
```
优先简单。
**Kafka 已落地**
```
Kafka + Kafka Streams
```
最具性价比。
**核心实时业务**
```
Kafka + Flink
```
当前主流黄金组合。
适用于:
- 电商
- 金融
- 风控
- 推荐系统
**数据平台建设**
```
Kafka + Spark
```
适用于:
- 数据仓库
- ETL
- AI训练
## 八、总结
一句话概括:
- Redis Stream:轻量级 MQ
- Kafka Streams:带状态管理的轻量流计算
- Flink:专业实时计算平台
- Spark:批流一体的大数据平台
如果业务只是异步解耦,Redis Stream 就够了。
如果已经全面使用 Kafka,需要一些实时统计能力,Kafka Streams 往往是成本最低的选择。
如果涉及海量数据、复杂窗口计算、CEP、实时风控等场景,Flink 是当前最主流的解决方案。
如果企业已经构建了完整的数据湖和离线数仓体系,那么 Spark 仍然是批处理和数据分析领域的重要选择。
技术选型的核心不是"最先进",而是"最匹配业务需求"。真正优秀的架构师关注的是成本、收益和复杂度之间的平衡,而不是技术栈的堆砌。
更多推荐




所有评论(0)