流处理系统:Flink与Spark Streaming的实时计算原理
一、从批处理到流处理:计算范式的演进
1.1 数据处理的两种模式
批处理(Batch Processing):
对有限、静态的数据集进行计算,强调高吞吐。
批处理特征:
├── 输入:有限数据集(如一天日志)
├── 处理:全量读取 → 计算 → 输出
├── 延迟:分钟级到小时级
├── 代表:Hadoop MapReduce、Spark SQL
└── 适用:离线分析、报表生成
流处理(Stream Processing):
对无限、动态的数据流进行实时计算,强调低延迟。
流处理特征:
├── 输入:无限数据流(如实时点击流)
├── 处理:来一条处理一条
├── 延迟:毫秒级到秒级
├── 代表:Flink、Spark Streaming、Kafka Streams
└── 适用:实时监控、风控、推荐
1.2 Lambda架构与Kappa架构
Lambda架构(Nathan Marz, 2011):
Lambda 架构:
┌─────────────┐
│ 数据源 │
└──────┬──────┘
│
┌────────┴────────┐
▼ ▼
┌─────────────┐ ┌─────────────┐
│ 批处理层 │ │ 速度层 │
│ (Batch) │ │ (Speed) │
│ 全量计算 │ │ 增量计算 │
│ 高准确 │ │ 低延迟 │
└──────┬──────┘ └──────┬──────┘
│ │
▼ ▼
┌─────────────┐ ┌─────────────┐
│ 批处理视图 │ │ 实时视图 │
│ (准确) │ │ (近似) │
└──────┬──────┘ └──────┬──────┘
│ │
└────────┬────────┘
▼
┌─────────────┐
│ 服务层 │
│ (合并视图) │
└─────────────┘
问题:
- 维护两套代码(批处理 + 流处理)
- 逻辑不一致风险
- 系统复杂度高
Kappa架构(Jay Kreps, 2014):
Kappa 架构:
┌─────────────┐
│ 数据源 │
└──────┬──────┘
│
▼
┌─────────────┐
│ 消息队列 │
│ (Kafka) │
└──────┬──────┘
│
▼
┌─────────────┐
│ 流处理层 │
│ (唯一处理) │
│ 统一代码 │
└──────┬──────┘
│
▼
┌─────────────┐
│ 服务层 │
└─────────────┘
改进:
- 只有一套代码(流处理)
- 重放机制实现批处理
- 系统简化
1.3 流处理的核心挑战
1. 无限数据:
- 内存无法容纳全部数据
- 需要窗口(Window)机制
2. 时间语义:
- 事件时间(Event Time)vs 处理时间(Processing Time)
- 乱序数据与水位线(Watermark)
3. 状态管理:
- 算子状态的持久化
- 故障恢复与一致性
4. exactly-once 语义:
- 不丢失数据
- 不重复处理
二、时间语义与窗口机制
2.1 三种时间概念
时间类型:
事件时间(Event Time):
├── 数据产生的时间
├── 嵌入在数据记录中
├── 最准确,但可能有乱序
└── 需要处理延迟数据
摄取时间(Ingestion Time):
├── 数据进入系统的时间
├── 由消息队列设置
├── 介于事件时间和处理时间之间
处理时间(Processing Time):
├── 算子处理数据的时间
├── 最简单,延迟最低
└── 不确定性最高
乱序数据示例:
实际发生顺序:
T=1s 事件 A (用户点击)
T=2s 事件 B (用户点击)
T=3s 事件 C (用户点击)
到达系统顺序(网络延迟):
T=1s 事件 A
T=5s 事件 C (延迟 2s)
T=6s 事件 B (延迟 4s)
问题:
- 如果按处理时间统计,T=1-3s 窗口只包含 A
- 但 B 和 C 实际也属于这个时间段
2.2 水位线(Watermark)
水位线机制:
一种表示"事件时间进展"的特殊标记,用于处理乱序数据。
水位线定义:
Watermark(t) = 当前最大事件时间 - 允许延迟
含义:
- 所有事件时间 < t 的数据都已到达
- 可以触发窗口计算
- 延迟数据(> t)可以丢弃或特殊处理
水位线传播:
class WatermarkStrategy:
"""
水位线生成策略
"""
@staticmethod
def for_bounded_out_of_orderness(delay: timedelta):
"""
有界乱序水位线
适用于最大延迟已知的场景
"""
return BoundedOutOfOrdernessGenerator(delay)
@staticmethod
def for_monotonous_timestamps():
"""
单调递增水位线
适用于无乱序场景
"""
return AscendingTimestampGenerator()
class BoundedOutOfOrdernessGenerator:
"""
有界乱序水位线生成器
"""
def __init__(self, max_out_of_orderness: timedelta):
self.max_delay = max_out_of_orderness
self.max_timestamp = datetime.min
def extract_timestamp(self, event) -> datetime:
"""提取事件时间戳"""
timestamp = event['timestamp']
self.max_timestamp = max(self.max_timestamp, timestamp)
return timestamp
def get_current_watermark(self) -> datetime:
"""生成当前水位线"""
return self.max_timestamp - self.max_delay
水位线可视化:
时间轴:
0s 1s 2s 3s 4s 5s 6s
│ │ │ │ │ │ │
├─────┼─────┼─────┼─────┼─────┼─────┤
│ A │ │ │ C │ B │ │ ← 事件到达
│ │ │ │ │ │ │
└─────┴─────┴─────┴─────┴─────┴─────┘
水位线进展(max_delay=2s):
T=1s: Watermark = 1s - 2s = -1s (不触发)
T=5s: Watermark = 3s - 2s = 1s (触发 [0-1s] 窗口)
T=6s: Watermark = 3s - 2s = 1s (等待更多数据)
2.3 窗口类型
滚动窗口(Tumbling Window):
固定大小,不重叠
时间: 0-5s 5-10s 10-15s 15-20s
┌─────┐┌─────┐┌─────┐┌─────┐
事件: │A B C││D E ││F G H││I J │
└─────┘└─────┘└─────┘└─────┘
滑动窗口(Sliding Window):
固定大小,可重叠
窗口大小=10s,滑动步长=5s
时间: 0-10s 5-15s 10-20s
┌───────┐
│A B C D│
└───┬───┘
│┌───────┐
││C D E F│
│└───┬───┘
│ │┌───────┐
│ ││E F G H│
│ │└───────┘
会话窗口(Session Window):
动态大小,由活动间隙决定
间隙(Gap)= 5s
事件: A B C D E F
时间: 1s 2s 10s 11s 12s 25s
├────┤ ├────┴────┤ │
│窗口1│ │ 窗口2 │ │窗口3
└────┘ └─────────┘ │
│
A,B 间隔 < 5s → 同一窗口 │
B,C 间隔 > 5s → 新窗口 │
C,D,E 连续 → 同一窗口 │
E,F 间隔 > 5s → 新窗口 │
全局窗口(Global Window):
只有一个窗口,所有数据都在其中
需要自定义触发器(Trigger)决定何时计算
适用场景:
- 自定义聚合逻辑
- 复杂事件处理(CEP)
三、Flink:真正的流处理引擎
3.1 Flink架构概览
Flink 架构:
┌─────────────────────────────────────────┐
│ Flink Client │
│ (提交作业,生成 JobGraph) │
└──────────────────┬──────────────────────┘
│
▼
┌─────────────────────────────────────────┐
│ JobManager (主节点) │
│ ┌─────────────┐ ┌─────────────────┐ │
│ │ Dispatcher │ │ ResourceManager │ │
│ │ (接收作业) │ │ (资源管理) │ │
│ └─────────────┘ └─────────────────┘ │
│ ┌─────────────────────────────────────┐│
│ │ JobMaster (作业管理) ││
│ │ - 调度执行图 (ExecutionGraph) ││
│ │ - 协调 Checkpoint ││
│ │ - 故障恢复 ││
│ └─────────────────────────────────────┘│
└──────────────────┬──────────────────────┘
│
▼
┌─────────────────────────────────────────┐
│ TaskManager (工作节点) × N │
│ ┌─────────────────────────────────────┐│
│ │ Slot 1 │ Slot 2 │ Slot 3 ││
│ │ (任务槽) │ (任务槽) │ (任务槽) ││
│ │ ┌──────┐ │ ┌──────┐ │ ┌──────┐ ││
│ │ │ Task │ │ │ Task │ │ │ Task │ ││
│ │ └──────┘ │ └──────┘ │ └──────┘ ││
│ └────────────┴────────────┴───────────┘│
└─────────────────────────────────────────┘
3.2 DataStream API
// Flink DataStream 示例
StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
// 设置事件时间和水位线
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
DataStream<ClickEvent> clicks = env
.addSource(new KafkaSource<>("click-stream"))
.assignTimestampsAndWatermarks(
WatermarkStrategy
.<ClickEvent>forBoundedOutOfOrderness(
Duration.ofSeconds(5)
)
.withTimestampAssigner(
(event, timestamp) -> event.getEventTime()
)
);
// 窗口聚合
DataStream<PageViewCount> pageViews = clicks
.keyBy(ClickEvent::getPageId)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.aggregate(new CountAggregate());
// 执行
env.execute("Page View Analytics");
3.3 Flink的状态管理
状态类型:
// Keyed State(键控状态)
// 每个 key 有独立的状态
ValueState<Integer> countState;
ListState<Event> eventListState;
MapState<String, Integer> mapState;
ReducingState<Event> reducingState;
AggregatingState<Event, Integer> aggState;
// Operator State(算子状态)
// 每个算子实例共享
ListState<Integer> operatorState;
BroadcastState<String, Rule> broadcastState;
状态后端:
MemoryStateBackend:
├── 状态存储在 JVM 堆内存
├── Checkpoint 存储在 JobManager 内存
├── 适用:本地测试、小状态
└── 限制:状态大小受内存限制
FsStateBackend:
├── 状态存储在 JVM 堆内存
├── Checkpoint 存储在文件系统
├── 适用:大状态、生产环境
└── 限制:状态大小仍受内存限制
RocksDBStateBackend:
├── 状态存储在 RocksDB(本地磁盘)
├── Checkpoint 存储在文件系统
├── 适用:超大状态、增量 Checkpoint
└── 优势:状态大小不受内存限制
3.4 Checkpoint机制
分布式快照(Chandy-Lamport算法):
Checkpoint 流程:
1. JobManager 向所有 Source 发送 Checkpoint Barrier
Source ──> Barrier ──> 下游算子
2. Source 快照自己的状态,向下游传播 Barrier
3. 算子收到 Barrier 后:
a. 将当前处理的数据缓冲
b. 快照自己的状态
c. 向下游传播 Barrier
4. 所有算子完成快照后,向 JobManager 确认
5. JobManager 收到所有确认后,Checkpoint 完成
Barrier对齐(Alignment):
多输入流的 Barrier 对齐:
算子有两个输入流:
Stream A: A1, A2, [Barrier], A3, A4
Stream B: B1, B2, B3, [Barrier], B4
处理过程:
1. 正常处理 A1, A2, B1, B2, B3
2. 收到 Stream A 的 Barrier
3. 暂停处理 Stream A,缓存 A3, A4
4. 继续处理 Stream B
5. 收到 Stream B 的 Barrier
6. 快照状态
7. 继续处理缓存的 A3, A4 和 B4
非对齐 Checkpoint(Flink 1.11+):
- 不暂停处理
- 将"在途数据"也作为状态的一部分
- 降低延迟,但状态更大
3.5 Exactly-Once语义
端到端 Exactly-Once:
三个层面的保证:
1. Flink 内部(Checkpoint)
└── 分布式快照保证状态一致性
2. Source(可重放)
└── Kafka:offset 作为状态
└── 文件:读取位置作为状态
3. Sink(幂等或事务)
├── 幂等写入:Elasticsearch、Redis
└── 事务写入:Kafka、JDBC
Kafka 事务写入示例:
- 开启事务
- 消费数据 → 处理 → 写入 Kafka
- Checkpoint 时提交事务
- 失败回滚到上一次 Checkpoint
四、Spark Streaming:微批处理架构
4.1 DStream与微批
Spark Streaming 核心思想:
将流处理转化为一系列小批处理(Micro-batch)。
微批处理:
实时流: ──────────────────────────────────────►
│ │ │ │ │ │ │
▼ ▼ ▼ ▼ ▼ ▼ ▼
Batch: [B1] [B2] [B3] [B4] [B5] [B6] [B7]
│ │ │ │ │ │ │
▼ ▼ ▼ ▼ ▼ ▼ ▼
处理: P1 P2 P3 P4 P5 P6 P7
│ │ │ │ │ │ │
└────┴────┴────┴────┴────┴────┘
输出
特点:
- 批大小:通常 100ms - 几秒
- 延迟:秒级
- 吞吐:高(利用 Spark SQL 优化)
4.2 Structured Streaming
Spark 2.0 引入 Structured Streaming,统一批处理和流处理 API:
// Structured Streaming 示例
val spark = SparkSession.builder()
.appName("StructuredStreaming")
.getOrCreate()
import spark.implicits._
// 读取流数据
val lines = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "localhost:9092")
.option("subscribe", "events")
.load()
.selectExpr("CAST(value AS STRING)")
// 处理(与批处理相同 API)
val wordCounts = lines
.flatMap(_.split(" "))
.groupBy("value")
.count()
// 输出
val query = wordCounts.writeStream
.outputMode("complete") // complete/update/append
.format("console")
.start()
query.awaitTermination()
输出模式:
| 模式 | 描述 | 适用场景 |
|---|---|---|
| Complete | 输出完整结果表 | 聚合操作 |
| Update | 只输出更新的行 | 有键的聚合 |
| Append | 只输出新增的行 | 无聚合、窗口 |
4.3 Spark Streaming的状态管理
UpdateStateByKey:
// 维护全局状态
val stateDStream = events
.map(event => (event.userId, event.value))
.updateStateByKey[Int](
// 更新函数
(newValues: Seq[Int], state: Option[Int]) => {
val currentSum = state.getOrElse(0)
val newSum = newValues.sum + currentSum
Some(newSum)
},
// 分区器
new HashPartitioner(sc.defaultParallelism)
)
MapGroupsWithState:
// 更灵活的状态管理
val sessionized = events
.groupByKey(event => event.sessionId)
.mapGroupsWithState[SessionState, SessionResult](
GroupStateTimeout.NoTimeout()
) { (sessionId, events, state) =>
val currentState = state.getOption.getOrElse(SessionState.empty)
val newState = events.foldLeft(currentState)(_ + _)
if (newState.isExpired) {
state.remove()
SessionResult(sessionId, newState, expired = true)
} else {
state.update(newState)
SessionResult(sessionId, newState, expired = false)
}
}
4.4 Checkpoint与WAL
Checkpoint机制:
Spark Streaming Checkpoint:
1. 元数据 Checkpoint:
├── 配置信息
├── DStream 操作链
└── 未完成的批次
2. 数据 Checkpoint:
└── 有状态操作的状态数据
└── 存储在 HDFS/S3
恢复流程:
1. 从 Checkpoint 读取元数据
2. 重建 DStream 图
3. 恢复状态数据
4. 继续处理
WAL(Write Ahead Log):
接收数据时的 WAL:
1. 从 Kafka 接收数据
2. 先写入 WAL(HDFS)
3. 再处理数据
4. 处理完成后删除 WAL
作用:
- Driver 故障时,从 WAL 恢复数据
- 保证数据不丢失
- 代价:写入延迟
五、Flink vs Spark Streaming 深度对比
5.1 架构对比
| 特性 | Flink | Spark Streaming |
|---|---|---|
| 处理模型 | 原生流处理 | 微批处理 |
| 延迟 | 毫秒级 | 秒级 |
| 吞吐 | 高 | 很高(批处理优化) |
| 状态管理 | 内置,多种后端 | 需借助外部系统 |
| Exactly-Once | 原生支持 | 需额外配置 |
| 事件时间支持 | 完善 | 较完善 |
| SQL支持 | Flink SQL | Spark SQL(更成熟) |
| 生态集成 | Kafka、ES、HBase | 更丰富的连接器 |
5.2 适用场景
选择 Flink:
- 延迟要求 < 1秒
- 复杂事件处理(CEP)
- 需要精确的事件时间处理
- 大状态场景(> 100GB)
选择 Spark Streaming:
- 延迟要求 > 1秒可接受
- 需要与 Spark 生态深度集成
- 大量 SQL/ML 处理
- 已有 Spark 团队
5.3 性能对比
基准测试(Yahoo Streaming Benchmark):
| 指标 | Flink | Spark Streaming | Storm |
|---|---|---|---|
| 吞吐 (events/s) | 10M+ | 20M+ | 5M |
| 延迟 (p99) | 50ms | 500ms | 100ms |
| CPU 使用率 | 中等 | 较高 | 高 |
| 内存使用 | 中等 | 较高 | 低 |
六、流处理系统实践
6.1 背压处理
背压(Backpressure):
当消费速度 < 生产速度时,系统需要减速上游生产,避免 OOM。
Flink 背压机制:
下游 Task 繁忙 → 缓冲区满 → 停止读取上游
↑ │
└──────────────────────────────┘
自动调节:
- 降低 Source 消费速度
- 动态调整 Checkpoint 间隔
- 自动扩缩容(K8s 集成)
6.2 监控与调优
关键指标:
STREAMING_METRICS = {
# 延迟指标
'records_lag': '消费延迟(Kafka)',
'watermark_delay': '水位线延迟',
'processing_delay': '处理延迟',
# 吞吐指标
'records_in_per_second': '输入吞吐',
'records_out_per_second': '输出吞吐',
# 状态指标
'state_size': '状态大小',
'checkpoint_duration': 'Checkpoint 耗时',
'checkpoint_size': 'Checkpoint 大小',
# 资源指标
'cpu_usage': 'CPU 使用率',
'memory_usage': '内存使用率',
'gc_time': 'GC 时间',
}
调优建议:
1. 并行度设置
└── 通常与 Kafka 分区数相同
2. 缓冲区大小
└── 增加缓冲区可减少网络往返
3. Checkpoint 间隔
└── 权衡:恢复速度 vs 性能影响
4. 状态后端选择
└── 大状态 → RocksDB
└── 小状态 → Heap
6.3 实际案例:实时风控系统
架构:
用户行为 ──> Kafka ──> Flink CEP ──> 规则引擎 ──> 决策
│
▼
历史特征(Redis)
▲
│
特征计算(Flink)
Flink 作业:
1. 实时特征计算
- 滑动窗口统计
- 会话特征提取
2. CEP 复杂事件检测
- 异常登录模式
- 高频交易检测
3. 规则评分
- 多维度风险评分
- 实时决策输出
性能:
- 延迟:P99 < 200ms
- 吞吐:100K+ events/s
- 可用性:99.99%
结语
流处理系统代表了大数据计算从"离线分析"到"实时决策"的范式转变。Flink 和 Spark Streaming 代表了两种架构哲学:
Flink:
"流处理优先,批处理是流的特例"
- 原生流处理架构
- 低延迟、高准确
- 适合实时性要求高的场景
Spark Streaming:
"批处理统一流处理"
- 微批处理架构
- 高吞吐、生态丰富
- 适合已有 Spark 生态的场景
理解流处理系统,不仅是掌握技术工具,更是理解实时计算的本质:在无限的数据流中,如何及时、准确地提取价值。
参考资源
经典论文:
- Akidau, T., et al. (2013). "MillWheel: Fault-Tolerant Stream Processing at Internet Scale". VLDB.
- Carbone, P., et al. (2015). "Apache Flink: Stream and Batch Processing in a Single Engine". IEEE Data Engineering Bulletin.
- Zaharia, M., et al. (2013). "Discretized Streams: Fault-Tolerant Streaming Computation at Scale". SOSP.
官方文档: 4. Flink 官方文档:https://nightlies.apache.org/flink/flink-docs-stable/ 5. Spark Streaming 指南:https://spark.apache.org/docs/latest/streaming-programming-guide.html 6. Kafka Streams 文档:https://kafka.apache.org/documentation/streams/
实践资源: 7. 《Streaming Systems》(Tyler Akidau) 8. Flink 中文社区:https://flink-learning.org.cn/
创建时间:2026年04月11日
更新时间:2026年04月11日