概述
Apache Flink 是一个开源的分布式流处理框架,被设计用于在有界和无界数据流上进行有状态的计算。与传统的微批处理(Micro-Batch)引擎不同,Flink 从底层就将流处理作为一等公民——批处理仅仅是有界流的特例。这种"流批一体"的设计哲学使 Flink 能够在低延迟、高吞吐和精确一次(Exactly-Once)语义之间取得卓越的平衡。
Flink 的核心能力包括:有状态计算、事件时间处理、轻量级容错、高吞吐低延迟、以及灵活的窗口和状态管理。目前已被阿里巴巴、Uber、Netflix、LinkedIn 等公司在生产环境中大规模使用,是实时计算领域的事实标准。
一、整体架构分层
Flink 的软件架构是一个分层、分布式、主从式的设计,从上到下可以分为四层:
1.1 部署层(Deployment Layer)
Flink 可以运行在多种环境中:
- Standalone 模式:Flink 自带资源管理,不依赖外部系统
- Flink on YARN:利用 Hadoop YARN 作为资源调度器
- Flink on Kubernetes:原生云原生部署,支持弹性伸缩
- 本地模式:单进程运行,用于开发和调试
1.2 编程接口层(API Layer)
Flink 提供三层 API,抽象程度从高到低:
| 最高层 | SQL & Table API | 数据分析师 | 声明式,自动优化,流批统一语义 |
| 中间层 | DataStream API | 应用开发者 | 精细控制流处理逻辑,核心编程模型 |
| 最底层 | Stateful ProcessFunction | 高级开发者 | 直接操作状态和定时器,最灵活 |
SQL & Table API 是最常用的接口,支持标准 SQL 语法,Flink 内部将其编译优化为 DataStream 程序执行。DataStream API 提供了 map、filter、keyBy、window、join、process 等丰富的算子,是构建复杂流处理应用的基础。
1.3 运行时层(Runtime Layer)
运行时层是 Flink 的核心引擎,负责作业调度、资源管理和分布式执行。
1.4 存储层(Storage Layer)
包括状态后端(RocksDB/Heap)、Checkpoint 存储(HDFS/S3/OSS)和高可用存储(ZooKeeper/K8s)。
二、核心组件详解
2.1 JobManager(Master 节点)
JobManager 是 Flink 集群的"大脑",负责整个作业的生命周期管理。其内部包含四个核心子组件:
| Dispatcher | 接收用户提交的作业,持久化 JobGraph,为每个作业创建独立的 JobMaster |
| JobMaster | 负责单个作业的全生命周期管理,将 JobGraph 展开为 ExecutionGraph 并调度执行 |
| CheckpointCoordinator | 周期性触发分布式快照,协调所有 TaskManager 完成 Checkpoint |
| ResourceManager | 管理 TaskManager 及其 Task Slot,与外部资源系统(YARN/K8s)交互申请资源 |
在高可用(HA)模式下,集群可以部署多个 JobManager,通过 ZooKeeper 或 Kubernetes 进行 Leader 选举,确保只有一个 Active JobManager,其余为 Standby。
2.2 TaskManager(Worker 节点)
TaskManager 是 Flink 的"工人",所有实际计算在此发生。每个 TaskManager 管理若干 Task Slot(任务槽):
- Task Slot 是资源调度的基本单位,代表 TaskManager 内存的一个固定子集
- 每个 Slot 运行一个 Task(可包含多个通过算子链融合的 Subtask)
- 一个 Subtask 是一个算子的一个并行实例
TaskManager 的关键模块:
- Operator Chain:将多个算子融合在同一线程中执行,减少序列化/网络开销
- Network Stack:基于 Netty + Credit-based 流控,管理 Task 间数据传输
- Memory Manager:管理 Managed Memory(用于排序、哈希、状态存储等)
- State Backend:存储和管理算子状态
2.3 Client(客户端)
Client 不是运行时的组成部分,它负责:
- 解析用户代码,构建 StreamGraph
- 优化为 JobGraph(算子链合并)
- 通过 HTTP REST 接口提交给 Dispatcher
- 支持 CLI / REST API / Web UI 三种操作方式
三、作业执行流程
Flink 作业从提交到运行,经历四层图转化:
用户代码
↓
StreamGraph(逻辑拓扑,包含所有算子节点和边)
↓ 算子链合并优化
JobGraph(减少不必要的序列化/反序列化)
↓ 提交到 Dispatcher
ExecutionGraph(物理执行计划,包含并行实例和网络拓扑)
↓ 调度执行
Task 在 TaskManager 的 Slot 中运行
详细流程:
四、核心机制
4.1 算子链(Operator Chain)
Flink 为了优化性能,会将多个满足条件的算子融合为一个 Task,在同一个线程中执行。这称为算子链(Operator Chaining)。
融合条件(必须同时满足):
- 算子之间是 One-to-One 连接(如 map → filter → keyBy 中的 map → filter)
- 算子具有相同的并行度
- 算子属于相同的 Slot 共享组
- 没有被用户通过 disableChaining() 显式禁用
好处: 避免不必要的序列化/反序列化和网络传输开销,显著降低延迟。
4.2 并行度(Parallelism)
Flink 的每个算子都可以独立设置并行度。一个算子的多个并行实例称为 Subtask(子任务),它们可以运行在不同的 TaskManager 上。
Source(并行度=3)→ Map(并行度=3)→ KeyBy/Window(并行度=2)→ Sink(并行度=1)
Source-0 ──┐
Source-1 ──┼── Map-0 ──┐
Source-2 ──┘ Map-1 ──┼── Window-0 ── Sink-0
Map-2 ──┘ Window-1 ──┘
- 并行度可以在代码中设置,也可以在全局配置中指定
- Flink 2.0 进一步支持了动态并行度调整(Adaptive Scheduler),根据运行时负载自动优化
4.3 数据交换模式
Flink 在 Task 之间采用 Pipeline 数据交换模式:
- One-to-One(Forwarding):数据在上下游 Subtask 之间一一对应传输,保持分区不变
- Redistributing(Shuffle):数据经过重分区后传输,包括:
- Hash 分区(keyBy 后的默认方式)
- Broadcast(广播到所有下游)
- Global(全部发往下游第一个 Subtask)
- Rescale(在局部范围内轮询分发)
Pipeline 交换的核心优势:数据产生后立即推送给下游,实现真正的流式处理,延迟极低。
五、状态管理
状态管理是 Flink 最核心的能力之一,使流式应用能够"记住过去的事件"并影响未来处理。
5.1 状态类型
| Keyed State | 按 Key 分区 | 每个 Key 维护独立的状态实例 | 用户累计消费金额 |
| Operator State | 算子级别 | 整个算子共享一份状态 | Kafka Consumer 的 offset |
Keyed State 的子类型:
- ValueState:存储单个值
- ListState:存储一个列表
- MapState:存储键值对映射
- ReducingState:通过 ReduceFunction 聚合
- AggregatingState:通过 AggregateFunction 聚合
5.2 状态后端
| HashMapStateBackend | JVM 堆内存 | 状态较小、对延迟敏感 | 访问快,但受内存限制 |
| EmbeddedRocksDBStateBackend(Flink 1.x) | 本地磁盘(嵌入式 RocksDB) | 状态较大(GB~TB 级) | 不受内存限制,但访问有磁盘 I/O 开销 |
5.3 状态 TTL
可以为状态设置过期时间(Time-To-Live),自动清理不再使用的状态数据,防止状态无限膨胀:
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Time.hours(24))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.cleanupInRocksdbCompactFilter(1000)
.build();
六、时间语义与窗口
6.1 三种时间语义
| Event Time | 事件发生的时间(数据自带) | 最精确,能处理乱序数据,依赖 Watermark |
| Processing Time | 数据到达算子时的系统时间 | 最简单,延迟最低,但结果不确定 |
6.2 Watermark(水位线)机制
Watermark 是 Flink 处理乱序数据的核心机制,用于衡量事件时间的进展:
数据流: [t=10] [t=12] [t=11] [t=15] [W(t=13)] [t=14] [t=16] [W(t=15)]
↑ ↑
Watermark=13 Watermark=15
- Watermark = 当前最大事件时间 – 允许的最大乱序时间
- 当 Watermark 达到或超过窗口结束时间时,触发窗口计算
- 晚于 Watermark 的数据被视为迟到数据,可通过侧输出流(Side Output)处理
Watermark 生成策略:
- BoundedOutOfOrderness:允许固定时间的乱序(最常用)
- Monotonous Timestamps:假设时间戳单调递增(无乱序)
- 自定义策略:根据业务特点自定义 Watermark 生成逻辑
6.3 窗口类型
| 滚动窗口(Tumbling) | 固定大小、无重叠 | 每5分钟统计一次 |
| 滑动窗口(Sliding) | 固定大小、有重叠 | 每1分钟统计最近5分钟的数据 |
| 会话窗口(Session) | 基于活动间隙动态划分 | 用户无操作超过30分钟则关闭窗口 |
| 全局窗口(Global) | 默认不触发,需自定义 Trigger | 自定义触发逻辑 |
窗口生命周期:
数据进入 → 分配到窗口 → 触发器(Trigger)决定是否触发计算
→ 窗口函数执行 → 驱逐器(Evictor)可选清理数据 → 输出结果
七、容错机制(Checkpoint)
7.1 Checkpoint 原理
Flink 的容错基于 Chandy-Lamport 分布式快照算法的变体——异步屏障快照(Asynchronous Barrier Snapshotting)。
核心流程:
7.2 Checkpoint vs Savepoint
| 触发方式 | 自动周期性触发 | 用户手动触发 |
| 用途 | 故障恢复 | 作业升级/迁移/版本变更 |
| 存储位置 | 由 Flink 管理 | 由用户管理(外部存储) |
| 格式 | 轻量级,优化性能 | 标准化,支持跨版本兼容 |
| 生命周期 | 旧 Checkpoint 自动清理 | 永久保留直到手动删除 |
7.3 非对齐 Checkpoint(Unaligned Checkpoint)
传统对齐 Checkpoint 的问题:当存在**反压(Backpressure)**时,Barrier 在反压节点堆积,导致 Checkpoint 超时。
非对齐 Checkpoint 的解决方案:
- Barrier 不再等待所有输入对齐
- 未对齐的数据也被包含在状态快照中
- 代价:状态快照更大,但 Checkpoint 时间与反压无关
八、窗口函数与聚合
8.1 增量聚合函数
| ReduceFunction | 每次输入与当前累积值合并,输出新的累积值 |
| AggregateFunction | 更灵活的聚合,支持累加器(Accumulator)类型 |
| FoldFunction | 类似 ReduceFunction,但输入和累加器类型可以不同(已废弃) |
8.2 全量窗口函数
| ProcessWindowFunction | 接收窗口内所有数据,提供全局上下文(窗口起止时间、状态、定时器等) |
| ApplyWindowFunction | 简化版的 ProcessWindowFunction |
8.3 组合使用示例
dataStream
.keyBy(event -> event.getUserId())
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.aggregate(
new CountAgg(), // 增量聚合:快速计算中间结果
new WindowResultFunction() // 全量函数:附加窗口元信息
);
九、Connector 与数据源
Flink 通过 Connector 与外部系统集成:
| Kafka | Source + Sink | 最常用的消息队列 Connector,支持 Exactly-Once |
| JDBC | Source + Sink | 关系型数据库读写 |
| Elasticsearch | Sink | 写入 ES 索引 |
| HDFS | Source + Sink | 文件系统读写,支持 Parquet/ORC 等格式 |
| Redis | Source + Sink | 键值存储交互 |
| CDC(Change Data Capture) | Source | 实时捕获数据库变更(MySQL/PostgreSQL/Oracle) |
| Paimon | Source + Sink | 流式湖存储,与 Flink 深度集成 |
Flink CDC 是一个重要的 Connector,基于 Debezium 实现,支持:
- 全量 + 增量一体化读取数据库变更
- Schema Evolution(DDL 变更自动同步)
- 多表整库同步
- 支持写入 Paimon、Iceberg 等湖格式
十、典型应用场景
10.1 实时数仓(Real-time Data Warehouse)
Flink 是实时数仓的核心计算引擎:
- 实时 ETL:从 Kafka 消费原始数据,清洗转换后写入 DWD 层
- 实时聚合:在 DWD 层基础上计算 DWS 汇总指标
- 实时报表:将计算结果写入 Doris/ClickHouse,供 BI 看板查询
10.2 实时风控
- 实时交易监控:毫秒级检测异常交易模式
- 反欺诈:基于规则引擎和 CEP(复杂事件处理)实时识别欺诈行为
- 信用评估:结合用户历史行为和实时交易进行动态评分
10.3 实时监控与告警
- 系统监控:实时采集和分析系统指标
- 业务监控:实时监控 GMV、DAU 等核心业务指标
- 异常检测:基于滑动窗口和 CEP 识别异常模式
10.4 实时推荐
- 用户行为实时分析:基于用户最近的行为序列更新推荐模型
- 特征工程:实时计算用户画像特征
- A/B 测试:实时分流和效果评估
10.5 数据集成与 CDC
- 数据库实时同步:通过 Flink CDC 将 MySQL 变更实时同步到数仓
- 异构系统数据迁移:在不同存储系统之间实时同步数据
- Schema Evolution:源表 DDL 变更时自动同步到目标系统
十一、安装部署
11.1 环境要求
| Java | JDK 8 或 JDK 11(Flink 2.0 要求 JDK 11+) |
| Hadoop | 2.x 或 3.x(使用 HDFS 时需要) |
11.2 下载与解压
下载地址:https://archive.apache.org/dist/flink/flink-1.17.1/ 
下载完后上传到 hadoop1 的 /opt/software/ 目录下,然后使用 tar 命令解压。 
修改配置文件 flink-conf.yaml
cd /opt/module/flink-1.17.1/conf
vim flink-conf.yaml
修改下面的条目
# JobManager节点地址.
jobmanager.rpc.address: hadoop1
jobmanager.bind-host: 0.0.0.0
rest.address: hadoop1
rest.bind-address: 0.0.0.0
# TaskManager节点地址.需要配置为当前机器名
taskmanager.bind-host: 0.0.0.0
taskmanager.host: hadoop1
修改配置文件 workers 为如下内容:
hadoop1
hadoop2
hadoop3
修改配置文件 masters
hadoop1:8081
11.3 分发到其它机器
cd /opt/module/
xsync flink-1.17.1
在其它机器上修改 flink-conf.yaml 文件中的 taskmanager.host 为自己的主机名:hadoop2、hadoop3
11.4 启动集群
在 hadoop1 上执行命令 Flink 集群
cd /opt/module/flink-1.17.1
bin/start-cluster.sh

11.5 访问 Web UI
启动成功后,可以访问 http://hadoop1:8081 对flink集群和任务进行监控管理

十二、与 Spark Streaming 对比
| 处理模型 | 真正的流处理(逐条处理) | 微批处理(Micro-Batch) |
| 延迟 | 毫秒级 | 秒~分钟级 |
| 容错机制 | Checkpoint + Barrier(Chandy-Lamport) | 基于 RDD 血缘重算 |
| 状态管理 | 原生支持,多种后端 | 依赖外部存储(如 HBase) |
| 时间语义 | Event Time + Watermark | 支持但实现较复杂 |
| 窗口模型 | 丰富的窗口类型,原生支持 | 有限的窗口支持 |
| 批处理 | 有界流(统一引擎) | 原生批处理引擎 |
| SQL 能力 | 完善,流批统一 | 成熟,生态丰富 |
| 生态成熟度 | 快速成长中 | 非常成熟,社区庞大 |
| 适用场景 | 实时计算、低延迟、有状态流处理 | 批处理、ETL、对延迟要求不高的场景 |
网硕互联帮助中心






评论前必须登录!
注册