Apache Flink:最惊艳的开源流计算框架
0. 引言
在流计算理论日趋成熟的今天,Apache Flink(以下简称 Flink)是这些理论与模型的优秀工程实践。与 Spark 以批处理为核心不同,Flink 的核心是流计算引擎,并明确把状态管理纳入系统架构。本文从系统架构、流的描述、流的处理、流的状态、消息处理可靠性五个方面展开,最后给出框架选型建议。
1. 系统架构:主从 + 状态内建
Flink 是主从(Master/Worker)架构的分布式系统:主节点负责任务调度与监控,从客户端接收作业 JAR 与资源后分析优化、生成执行计划,再分配给 Worker 执行。
Flink 可部署在 YARN、Mesos、Kubernetes 等资源管理器上。与 Storm、Spark Streaming 最大的不同是:Flink 明确地把状态管理纳入系统架构——计算节点执行任务时可将状态保存到本地,通过 Checkpoint 机制配合 HDFS/S3 等分布式文件系统,在不降低性能的前提下实现状态的分布式管理。
版本注:Flink 1.9 起引入存算分离的流式存储(State Processor API 与 SQL 持续演进);1.12 起 DataSet API 冻结、流批统一到 DataStream/SQL;1.15+ 推荐使用 KafkaSource/KafkaSink 新连接器(取代旧版 FlinkKafkaConsumer/Producer)。
2. 流的描述:DataStream
Flink 使用 DataStream 描述数据流(角色类似 Spark 的 RDD),输入、处理、输出均围绕它展开:
| 组成 | 职责 | 示例 |
|---|---|---|
| Source | 流数据输入源:消息中间件、数据库、文件系统等 | addSource、socketTextStream、fromCollection |
| Transformation | 将一个或多个 DataStream 转化为新 DataStream | map、flatMap、filter、keyBy、window、union、join |
| Sink | 将 DataStream 输出到外部系统 | print、writeAsText、addSink |
通过对 DataStream 进行各种 Transformation,就形成了描述流计算过程的 DAG。值得一提的是:Flink 支持批处理(DataSet/批模式),但将批处理视为流处理的特殊情况——这与 Spark Streaming"将流处理视为连续批处理"的做法截然相反。
版本注:Flink 1.12+ 流批统一后,DataSet API 逐渐淡出,批任务建议直接使用 Table/SQL 或 DataStream 的批执行模式。
3. 流的处理
3.1 输入与输出
内置数据源分四类:基于文件(readTextFile)、基于套接字(socketTextStream)、基于集合(fromCollection)、自定义(addSource)。输出同样分文件、控制台、套接字、自定义四类(addSink)。
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<String> text = env.socketTextStream("localhost", 9999, "\n");
DataStream<WordWithCount> windowCounts = text
.flatMap((String value, Collector<WordWithCount> out) -> {
for (String word : value.split("\\s")) {
out.collect(new WordWithCount(word, 1L));
}
})
.keyBy(w -> w.word)
.window(TumblingProcessingTimeWindows.of(Time.seconds(5)))
.reduce((a, b) -> new WordWithCount(a.word, a.count + b.count));
windowCounts.print().setParallelism(1);以上代码实现单词计数:socket 读入文本流 → flatMap 分词 → keyBy 按单词分组 → 5 秒滚动窗口 reduce 聚合。
3.2 流数据状态与窗口
Flink 的设计思路更清晰:明确将流信息状态从流数据状态分离。DataStream 转化操作分两类——常规流式操作(map/filter/reduce)与流数据状态相关操作(window/union/join/coGroup)。窗口 API 中 Window 针对 KeyedStream,WindowAll 针对非 KeyedStream。
3.3 反向压力
Flink 采用容量有限的分布式阻塞队列传递数据:下游读取过慢时自然减慢上游写入。与 Storm、Spark Streaming 需显式开关不同,Flink 的反向压力天然内建于数据传送方案,无需额外配置。
4. 流的状态:Keyed State 与 Operator State
Flink 是第一个将流信息状态管理从流数据状态管理中剥离的框架:DataStream 是对数据在时间维度的管理,状态接口是对数据在空间维度的管理。
| 状态类型 | 划分维度 | 典型用途 |
|---|---|---|
| Keyed State | 按 key 值划分(关联 KeyedStream) | "统计不同 IP 上出现的不同设备数" |
| Operator State | 按算子并行度划分(绑定并行实例) | Kafka Consumer 记录各 partition 的消费 offset |
两者的实现思路一致:Keyed State 按 key 划分状态空间,Operator State 按并行实例划分——类似线程局部量,运行时互不影响。Flink 1.6 起引入状态 TTL(State Time-To-Live),为淘汰过期状态提供便利。
5. 消息处理可靠性:基于 Checkpoint 的 exactly-once
Flink 基于 snapshot + checkpoint 的故障恢复机制,在内部提供 exactly-once 级别保证,前提是:
- 状态必须使用 Flink 内部状态机制(Keyed State / Operator State),其保存与恢复都在 Flink 故障机制内;
- 若使用外部存储(如独立 Redis 记录 PV/UV),只能获得 at-least-once;
- 端到端 exactly-once 还需 connectors 配合:如 Kafka Source 提供 exactly-once,部分 Sink 仅 at-least-once。
6. 面对流计算框架的选型建议
从横向功能看,所有流计算框架的核心概念相同(DAG、队列、线程、状态、反向压力),掌握核心概念即可触类旁通;从纵向发展看,Flink 是当前流计算领域最推荐的入口,其次是 Spark Streaming(微批、生态成熟),Storm/Samza 了解概念即可。Apache Beam 则是构建在具体框架之上的统一编程模式,适合追踪领域进展。
7. 小结
- 架构:主从式 + 状态内建(Checkpoint + 分布式文件系统);
- 描述:DataStream 为核心,Source/Transformation/Sink 三段式;
- 处理:窗口与流数据状态集中在 DataStream,流批一体(批为流的特例);
- 状态:Keyed State(按 key)/ Operator State(按并行度),支持 TTL;
- 可靠性:Checkpoint 快照 + 内部状态机制 → exactly-once;端到端保证需 connector 配合。
下一章讲解场景案例:如何用 Flink 实现实时风控引擎。