Flink 架构总览
概述
Flink 是流式计算框架,与 Spark 的微批不同,Flink 以"数据事件"为单位持续处理,天然支持事件时间与精确一次语义。本文讲清 JobManager/TaskManager 分工、Task Slot 与并行度、RPC 通信,以及与 Spark Streaming 的架构差异。
一、Flink 是什么
| 特性 | 说明 |
|---|---|
| 流批一体 | 流处理 + 批处理统一引擎 |
| 毫秒级延迟 | 逐条事件处理 |
| 事件时间 | 原生支持乱序处理 |
| 精确一次 | Checkpoint + 状态恢复 |
| 状态管理 | 强类型 Keyed State |
定位:实时计算的标杆引擎,被广泛用于实时数仓、风控、推荐等场景。
二、运行时角色
2.1 架构图
┌──────────────────────────────────────────────┐
│ JobManager(作业管理器) │
│ 作业调度 │ Checkpoint 协调 │ 资源申请 │ 恢复 │
└───────┬──────────────┬───────────────────────┘
│ 调度任务 │ 状态查询
┌───────┴──────┐ ┌────┴───────┐
│ TaskManager │ │ TaskManager│ ...
│ ┌────────┐ │ │ ┌────────┐ │
│ │ Task │ │ │ │ Task │ │
│ │ Task │ │ │ │ Task │ │
│ └────────┘ │ │ └────────┘ │
└──────────────┘ └────────────┘2.2 角色职责
| 角色 | 职责 |
|---|---|
| JobManager | 调度作业、协调 Checkpoint、故障恢复、资源管理 |
| TaskManager | 执行任务、管理状态、数据缓冲(Slot 容器) |
| Task | 算子并行实例(TaskManager 中的执行单元) |
| Slot | TaskManager 内固定资源的执行槽 |
2.3 客户端与提交
提交模式:
1. Standalone:客户端提交 JobGraph 到 JobManager
2. YARN/K8s:客户端申请集群,启动 JobManager| 组件 | 说明 |
|---|---|
| Client | 构建作业图、提交 |
| JobGraph | 优化后的作业图 |
| ExecutionGraph | JobManager 展开的并行执行图 |
三、Task Slot 与并行度
3.1 Slot 概念
TaskManager 内存按 Slot 切分(如 4 个 Slot)
每个 Slot 可运行一个 Task 线程
槽位数量决定单节点并行度上限| 参数 | 默认 | 说明 |
|---|---|---|
taskmanager.numberOfTaskSlots | 1 | 每 TM 槽位数 |
parallelism.default | 1 | 默认并行度 |
| 算子并行度 | 算子单独设置 | setParallelism |
3.2 并行度设置方式
| 方式 | 生效 |
|---|---|
算子 .setParallelism(n) | 该算子 |
| 执行环境 | 全局 |
提交 -p n | 作业级 |
| 配置默认 | 全局默认 |
3.3 Slot 与并行度的关系
并行度 = 总 Slot 数(理想)
TaskManager 数 × Slot 数 = 集群可承载并行度
并行度 > Slot 数 → 排队执行3.4 Slot 共享
| 特点 | 说明 |
|---|---|
| 默认共享 | 不同算子的 Task 可共享一个 Slot |
| 好处 | 提高 Slot 利用率 |
| 关闭 | slotSharingGroup 隔离 |
四、RPC 通信
4.1 通信架构
JobManager ↔ TaskManager:Akka RPC(控制面)
TaskManager ↔ TaskManager:Netty(数据面)| 通道 | 协议 | 用途 |
|---|---|---|
| 控制面 | Akka/RPC | 调度、心跳、Checkpoint 协调 |
| 数据面 | Netty | 算子间数据传输 |
4.2 数据交换
Task 间数据传输:
分区分配 → 网络缓冲 → 目标 Task
背压(Backpressure):下游慢时上游限速| 背压机制 | 说明 |
|---|---|
| 缓冲池 | 固定大小,耗尽即阻塞 |
| 反压传播 | 逐级向上限速到数据源 |
| 监控 | UI 查看 Backpressure 指标 |
五、Flink vs Spark Streaming
5.1 处理模型
| 维度 | Flink | Spark Streaming |
|---|---|---|
| 模型 | 连续流处理 | 微批(DStream) |
| 延迟 | 毫秒级 | 秒级 |
| 事件时间 | 原生支持 | 有限支持 |
| 状态 | 强 Keyed State | 有限状态 |
5.2 架构差异
| 维度 | Flink | Spark Streaming |
|---|---|---|
| 核心 | 流式引擎 | 批引擎 + 微批 |
| 调度 | 作业图细粒度 | 批次 Job |
| 容错 | 分布式快照 | RDD 血统 + WAL |
| 背压 | 天然反压 | 显式限速 |
5.3 选型
| 场景 | 推荐 |
|---|---|
| 毫秒级延迟要求 | Flink |
| 乱序/事件时间复杂窗口 | Flink |
| 已有 Spark 批栈,秒级延迟可接受 | Spark Structured Streaming |
| 批流一体 | Flink |
六、作业生命周期
提交 → 构建作业图 → 调度 Task → 执行(持续运行)
→ Checkpoint 周期快照
→ 故障时从快照恢复
→ 取消/停止(正常或 Savepoint 存档)| 状态 | 说明 |
|---|---|
| RUNNING | 正常运行 |
| FAILING/FAILED | 失败处理/失败 |
| CANCELLING/CANCELED | 取消 |
| FINISHED | 正常结束 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 并行度不足 | 调整 slot 数、并行度、TM 数量 |
| 背压严重 | 下游处理慢,优化或扩容 |
| JobManager 单点 | 配置 HA(ZooKeeper/K8s) |
| 数据倾斜 | 检查 key 分布,做预聚合 |
| 作业重启频繁 | 查看日志定位失败算子,检查恢复策略 |