一、Flink 核心基础认知
- 核心定位:分布式、实时、有状态、流批一体计算引擎
- 核心思想:流是本质,批是有界流,统一一套流式 Runtime 执行
- 两大 API 体系
- DataStream API:底层、灵活、支持复杂状态业务(血缘解析项目使用)
- Table/SQL API:高层、简洁、流批一体,适合常规统计分析
- 核心优势:低延迟、支持乱序数据、状态持久化、故障自动容错恢复
二、Flink 运行架构(底层基石)
1. 集群三大角色
- JobManager(JM):调度作业、管理 checkpoint、协调故障恢复、生成执行图
- TaskManager(TM):真正执行业务Task、管理 Slot、存储运行时状态
- Client:作业提交、参数封装、生成 JobGraph
2. 资源与执行单元
- Slot:TM 最小执行工位,隔离资源,一个 Slot 同一时刻跑一个 Task
- 并行度 Parallelism:算子并发执行数量,决定 SubTask 数量
- Task/SubTask
- Task:算子链合并后的逻辑执行单元
- SubTask:Task 的并发实例,真正运行的线程
- 算子链机制:无 shuffle 的连续算子合并为一个 Task,减少 IO 开销;keyBy 等 shuffle 操作会切断算子链
3. 四级执行图流转
代码逻辑 → StreamGraph(逻辑算子图) → JobGraph(算子链优化) → ExecutionGraph(并发物理执行图)
三、数据流全链路(Source → Transform → Sink)
- Source 数据源
- 无界流:Kafka(核心)、日志流
- 有界流:文件、Hive、数据库全量数据
- 核心知识点:offset 托管、分区与并行度约束、反序列化、脏数据处理
- 规则:Kafka Source 并行度 ≤ Topic 分区数
- Transform 转换算子
- 无状态算子:map、filter、flatMap(一条输多条,血缘解析核心算子)
- 重分区算子(产生 Shuffle、切断算子链):keyBy、rebalance、broadcast
- 有状态顶级算子:ProcessFunction(Flink 万能算子,支持状态+定时器)
- 拓展机制:侧输出流(分流脏数据、迟到数据)
- Sink 数据输出
- 常见目标:Kafka、MySQL、ES、NebulaGraph、Hive
- 生产核心要求:连接池复用、批量刷写、失败重试、死信队列
- 一致性核心:幂等写入、2PC 两阶段提交、端到端语义保障
四、时间与乱序处理机制(流处理灵魂)
- 三种时间语义
- ProcessingTime:机器处理时间,最快、不准
- EventTime:事件发生时间,生产唯一使用,保证业务准确
- IngestionTime:数据进入 Flink 的时间
- Watermark 水位线
- 本质:时间戳,标记「该时间前数据全部到达」
- 核心作用:处理乱序数据、驱动窗口触发
- 特性:严格单调递增,不可回退
- 乱序处理:设置水位线延迟 + allowedLateness 允许迟到数据
五、Window 窗口机制
- 四大窗口类型
- 滚动窗口:无重叠、固定周期
- 滑动窗口:可重叠、固定步长
- 会话窗口:间隔超时触发,无固定周期
- 全局窗口:手动触发,无自动周期
- 窗口四大组成
窗口分配器 + 触发器 + 剔除器 + 窗口聚合函数 - 核心规则
- EventTime 窗口:由 Watermark 驱动触发,非系统时间
- 支持迟到数据处理、窗口侧输出丢弃数据
六、状态管理机制(Flink 核心中的核心)
- 状态分类
-
KeyedState(主流):keyBy 后使用,按 Key 分片
- ValueState、ListState、MapState、聚合状态
- 血缘项目用途:ListState 缓存迟到血缘边数据
OperatorState:绑定 SubTask,无 key 分片,用于 Source Offset、自定义 Sink
重点:普通 Java 集合变量不属于 Flink 托管状态,重启丢失、不参与快照
- StateBackend 状态后端
- HashMapStateBackend:状态存 JVM 堆内存,小状态使用,大状态易 OOM
- RocksDBStateBackend(生产首选)
- 运行时状态:存储在 TM 本地磁盘
- 快照备份:持久化到 HDFS
- 支持增量 Checkpoint、超大状态、磁盘落地防 OOM
- Flink 内嵌管控,无独立进程,业务仅通过 State API 操作
- MemoryStateBackend:纯内存,仅测试使用
- 状态 TTL
- 作用:自动过期清理闲置状态,防止状态无限膨胀、磁盘溢出
- 原理:标记过期,后台 Compaction 阶段真正删除
七、容错机制(Checkpoint + Savepoint)
- Checkpoint 自动快照
- 核心原理:Barrier 屏障机制,切割数据流、触发算子状态快照
- 快照方式:HashMap 同步快照、RocksDB 异步增量快照
- 核心配置:快照间隔、超时时间、最大并发快照、对齐/非对齐 CP
- 作用:故障自动恢复状态、Kafka Offset,保障数据不丢不重
- Savepoint 手动快照
- 手动触发、语义兼容、版本稳定
- 使用场景:作业升级、代码变更、集群迁移、扩缩容
- 区别:Checkpoint 服务于故障恢复,Savepoint 服务于运维变更
- 数据一致性语义
- At-Least-Once:至少一次(默认)
- Exactly-Once:精准一次(Checkpoint + 2PC 事务 Sink / 幂等写入)
八、生产核心问题与调优
- 反压机制:上下游流速不匹配,逐级阻塞,定位瓶颈节点
- 数据倾斜:Key 分布不均,单 SubTask 压力过大,解决方案:预聚合、加盐打散
- RocksDB 调优:堆外内存、压缩策略、增量快照、Compaction 优化
- Checkpoint 优化:避免频繁快照、解决超时、堆积、小文件过多问题
- 状态优化:TTL 配置、状态拆分、超大状态并行度扩容
1.5版本之后不需要再指定taskmanager个数,由flink自己计算
计算公式
taskmanger数量 = 平行度/每个 TM 承载 Slot 数 避免任务因为slot配置不合理无法启动问题
二、Flink 系统参数
| 参数 | 值 | 说明 |
|---|---|---|
-m yarn-cluster |
yarn-cluster |
运行模式:提交到 YARN 集群 |
-d |
(无) | detach 模式,异步提交后返回 |
-yqu |
root.stream_compute |
YARN 队列名称 |
-ynm |
BDP_RTC_FLINK_20001773_itao_bdp_dag_kinship_impt |
YARN 应用名称,用于 UI 识别 |
-s |
hdfs://RTC-FT-KJ/tmp/drcs_flink/prod/20001773/.../chk-110192 |
savepoint的保存路径 |
-n |
(无) | 非严格恢复模式,允许跳过无法恢复的状态 |
-yjm |
2048 |
JobManager 内存 (MB) |
-ytm |
2048 |
TaskManager 内存 (MB) |
-ys |
1 |
每个 TaskManager 的 Slot 数量 |
-yD high-availability=zookeeper |
zookeeper |
启用 ZooKeeper 高可用 |
-yD high-availability.storageDir=... |
HDFS 路径 | HA 元数据存储目录 |
-yD high-availability.zookeeper.path.root=... |
/drcs_flink/ha/prod/20001773 |
ZooKeeper 根节点路径 |
-yD high-availability.zookeeper.quorum=... |
10.117.88.166:12181,... |
ZooKeeper 集群地址 |
-yD state.checkpoints.dir=... |
HDFS 路径 | Checkpoint 数据持久化目录 |
-c |
com.sf.bdp.itao.dag.kinship.impt.hive.HiveLineageImport |
指定 Java 主类 |
| (JAR 路径) | /app/.../itao-bdp-dag-kinship-impt-hive-1.70.jar |
应用程序 JAR 包 |
三、自定义参数
JAR 包路径之后的所有参数,通过
--key value格式传递,由ParameterTool.fromArgs(args)解析
| 参数 | 值 | 说明 |
|---|---|---|
--bootstrap.servers |
bigdataiot.kafka.sfcloud.local:9095 |
Kafka Broker 地址 |
--zookeeper.connect |
bigdataiot.kafka.sfcloud.local:2181/kafka/bigdataiot |
ZooKeeper 连接串 |
--kafka.topic |
BDP_TABLE_USAGE |
消费的主题名称 |
--group.id |
itao_bdp_dag_kinship_import |
Kafka 消费者组 ID |
--nebula.space |
itao_bdp_dag_kinship_v2 |
表级血缘图空间 |
--nebula.space.column |
itao_bdp_dag_kinship_col |
字段级血缘图空间 |
--parallelism |
4 |
Flink 并行度(代码限制最大为 4) |
--write.mode |
remote |
写入模式:remote=上报后端;local=直接写入 Nebula |
--exception.print |
true |
异常时是否打印原始数据 |
flink ha
flink ha 采用yarn 模式每次只会注册一个jobmanager,采用zk方式,配置好ha参数,resourcemanager检测到失败
会在重新拉起一个application,读取zk中最新的checkpoint,只有开启ha才能精确确定checkpoint位置
不然都需要人工传入处理
jobmanger主要:
- 生成Checkpoint 快照,管控状态一致性
- 管理task的状态,失败重启等
# ha的配置说明
-yD high-availability=zookeeper` `zookeeper` | 启用 ZooKeeper 高可用 |
-yD high-availability.storageDir=...` | HDFS 路径 | HA 元数据存储目录 |
-yD high-availability.zookeeper.path.root=...` | `/drcs_flink/ha/prod/20001773` | ZooKeeper 根节点路径 |
-yD high-availability.zookeeper.quorum=...` | `10.117.88.166:12181,...` | ZooKeeper 集群地址 |
statebackend
运行是状态
Flink 托管算子状态(State)的存储实现:
负责:状态本地读写、快照持久化(Checkpoint)、故障后从快照恢复状态。
主要在运行时保存中间结果、窗口数据、Kafka 偏移量、自定义缓存数据,这些数据都由 StateBackend 管理落地
运行时常用的状态管理有
- EmbeddedRocksDBStateBackend,本地构建rocksdb数据库,只需要选用,yarn集群注意选用大磁盘,适用于大状态维护,比如大维表(gb级别使用)
- HashMapStateBackend, 内存维护状态,(mb级别使用)
checkpoint
checkpoint 是运行是状态阶段性状态的快照
time
时间,主要值事件到达的实际
- ProcessingTime 到达机器的时间
- EventTime 事件的业务时间(一般选用这个)
- IngestionTime 时间进入集群的实际,基本与到达集群时间一致
flink cdc
1.x 版本有锁
作用:先把表存量历史数据一次性全部拉取,保证初始数据完整,再无缝切增量,避免数据断层。
执行流程:
-
加全局读锁(可选)
默认模式下短暂执行FLUSH TABLES WITH READ LOCK,锁住整张表禁止写入,防止快照期间数据变动;
InnoDB 可以使用快照事务(repeatable read)无锁快照,避免锁表影响业务(线上主流配置 -
记录当前 Binlog 点位(文件名 + 偏移量)
锁定瞬间,立刻获取当前 binlog 文件名称、position 偏移量,标记快照结束的位点。后续增量就从这个位点开始消费。 -
分批查询全量数据
按照主键 / 唯一索引分片,分页 SELECT 读取整张表数据,逐条发送至 Flink 流;
大数据量表会分片拉取,防止单次查询超时、内存溢出。 -
释放读锁
全量数据读取完毕后,释放表锁(无锁快照无需此步骤)。
全量结束后,Flink CDC 根据之前记录的 binlog 位点,模拟MySQL 从库(Slave) 协议,向 MySQL 主库发起拉取 binlog 请求:
- 遵循 MySQL 主从复制协议,持续拉取 ROW 格式 binlog 日志;
- Debezium 解析 binlog 二进制日志,把原始日志封装成标准化 JSON / 结构化变更数据:
{
"before": 变更前数据(update/delete才有),
"after": 变更后数据(insert/update才有),
"op": "c/u/d/r" // c新增、u更新、d删除、r快照数据
"source": binlog位点、库表、时间戳等元数据
}
- 转换为 Flink 内部数据流,交由 Flink 算子实时处理;
- 持续消费直到程序停止,属于无限流。
2.x 版本采用chunk分拆无锁模式
无锁拆分模式算法过程,前提支持全量快照模式才行
- 根据表主键查询数量,进行拆分,默认值8096行
- 根据分配最小主键,最大主键查询数据
select * from t where id >= start and id < end; - 将相关数据发送发送kafka
- 增量阶段从全局binlog缓存中过滤出自己的部分发送kafka
- 最终消费与binlog消费一致后,由binlog读取线程发送kafka不会再有进行拆分chunk发送
mysql的mvvc
mysql 对每一行数据有两个隐藏的字段 分别是事物id,undo.log的上一个历史版本
mysql更新数据的过程
- mysql对数据修改会加上行级锁,保障mvvc版本链的有序性
- 将旧数据写入undo.log日志,生产版本号
- 更新数据,将当前事物id写入,将undolog版本写入两个隐藏的字段
- 解锁行数据
在同一个事物内对同一行数据修改多次,形成多个版本的mvvc版本联调,如果回滚从而回滚到最初的那个
mysql的事物可见性
在我开始事物时,查找所有未交事物的id,找出最小值,最大值
- 事物id == 自身 可见
- 事物id < 最小id 可见,因为比最小未提交事物id小说么这个事物一定已经提交了
- 事物id >= 最大事物id 不可见
- 事物id > 最小事物id 不可见