flink 核心知识

一、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)

  1. Source 数据源
  • 无界流:Kafka(核心)、日志流
  • 有界流:文件、Hive、数据库全量数据
  • 核心知识点:offset 托管、分区与并行度约束、反序列化、脏数据处理
  • 规则:Kafka Source 并行度 ≤ Topic 分区数
  1. Transform 转换算子
  • 无状态算子:map、filter、flatMap(一条输多条,血缘解析核心算子)
  • 重分区算子(产生 Shuffle、切断算子链):keyBy、rebalance、broadcast
  • 有状态顶级算子:ProcessFunction(Flink 万能算子,支持状态+定时器)
  • 拓展机制:侧输出流(分流脏数据、迟到数据)
  1. Sink 数据输出
  • 常见目标:Kafka、MySQL、ES、NebulaGraph、Hive
  • 生产核心要求:连接池复用、批量刷写、失败重试、死信队列
  • 一致性核心:幂等写入、2PC 两阶段提交、端到端语义保障

四、时间与乱序处理机制(流处理灵魂)

  1. 三种时间语义
  • ProcessingTime:机器处理时间,最快、不准
  • EventTime:事件发生时间,生产唯一使用,保证业务准确
  • IngestionTime:数据进入 Flink 的时间
  1. Watermark 水位线
  • 本质:时间戳,标记「该时间前数据全部到达」
  • 核心作用:处理乱序数据、驱动窗口触发
  • 特性:严格单调递增,不可回退
  • 乱序处理:设置水位线延迟 + allowedLateness 允许迟到数据
    五、Window 窗口机制
  1. 四大窗口类型
  • 滚动窗口:无重叠、固定周期
  • 滑动窗口:可重叠、固定步长
  • 会话窗口:间隔超时触发,无固定周期
  • 全局窗口:手动触发,无自动周期
  1. 窗口四大组成
    窗口分配器 + 触发器 + 剔除器 + 窗口聚合函数
  2. 核心规则
  • EventTime 窗口:由 Watermark 驱动触发,非系统时间
  • 支持迟到数据处理、窗口侧输出丢弃数据
    六、状态管理机制(Flink 核心中的核心)
  1. 状态分类
  • KeyedState(主流):keyBy 后使用,按 Key 分片

    • ValueState、ListState、MapState、聚合状态
    • 血缘项目用途:ListState 缓存迟到血缘边数据
  • OperatorState:绑定 SubTask,无 key 分片,用于 Source Offset、自定义 Sink

  • 重点:普通 Java 集合变量不属于 Flink 托管状态,重启丢失、不参与快照

  1. StateBackend 状态后端
  • HashMapStateBackend:状态存 JVM 堆内存,小状态使用,大状态易 OOM
  • RocksDBStateBackend(生产首选)
    • 运行时状态:存储在 TM 本地磁盘
    • 快照备份:持久化到 HDFS
    • 支持增量 Checkpoint、超大状态、磁盘落地防 OOM
    • Flink 内嵌管控,无独立进程,业务仅通过 State API 操作
  • MemoryStateBackend:纯内存,仅测试使用
  1. 状态 TTL
  • 作用:自动过期清理闲置状态,防止状态无限膨胀、磁盘溢出
  • 原理:标记过期,后台 Compaction 阶段真正删除

七、容错机制(Checkpoint + Savepoint)

  1. Checkpoint 自动快照
  • 核心原理:Barrier 屏障机制,切割数据流、触发算子状态快照
  • 快照方式:HashMap 同步快照、RocksDB 异步增量快照
  • 核心配置:快照间隔、超时时间、最大并发快照、对齐/非对齐 CP
  • 作用:故障自动恢复状态、Kafka Offset,保障数据不丢不重
  1. Savepoint 手动快照
  • 手动触发、语义兼容、版本稳定
  • 使用场景:作业升级、代码变更、集群迁移、扩缩容
  • 区别:Checkpoint 服务于故障恢复,Savepoint 服务于运维变更
  1. 数据一致性语义
  • 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 版本有锁

作用:先把表存量历史数据一次性全部拉取,保证初始数据完整,再无缝切增量,避免数据断层。
执行流程:

  1. 加全局读锁(可选)
    默认模式下短暂执行 FLUSH TABLES WITH READ LOCK,锁住整张表禁止写入,防止快照期间数据变动;
    InnoDB 可以使用快照事务(repeatable read)无锁快照,避免锁表影响业务(线上主流配置
  2. 记录当前 Binlog 点位(文件名 + 偏移量)
    锁定瞬间,立刻获取当前 binlog 文件名称、position 偏移量,标记快照结束的位点。后续增量就从这个位点开始消费。
  3. 分批查询全量数据
    按照主键 / 唯一索引分片,分页 SELECT 读取整张表数据,逐条发送至 Flink 流;
    大数据量表会分片拉取,防止单次查询超时、内存溢出。
  4. 释放读锁
    全量数据读取完毕后,释放表锁(无锁快照无需此步骤)。

全量结束后,Flink CDC 根据之前记录的 binlog 位点,模拟MySQL 从库(Slave) 协议,向 MySQL 主库发起拉取 binlog 请求:

  1. 遵循 MySQL 主从复制协议,持续拉取 ROW 格式 binlog 日志;
  2. Debezium 解析 binlog 二进制日志,把原始日志封装成标准化 JSON / 结构化变更数据:
{
  "before": 变更前数据(update/delete才有),
  "after": 变更后数据(insert/update才有),
  "op": "c/u/d/r"  // c新增、u更新、d删除、r快照数据
  "source": binlog位点、库表、时间戳等元数据
}
  1. 转换为 Flink 内部数据流,交由 Flink 算子实时处理;
  2. 持续消费直到程序停止,属于无限流。

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 不可见
©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

友情链接更多精彩内容