Flink 面试专题
概述
Flink 面试高频覆盖:容错与精确一次、状态管理、窗口与时间语义、反压排查、CDC 应用。本文按原理、对比、场景设计三类整理高频问答,附易错点清单。
一、容错与精确一次
Q1:Flink 如何实现 Exactly-Once?
三端配合:
源:offset 快照(checkpoint)
计算:分布式快照(Barrier 对齐)
输出:两阶段提交(事务 Sink)| 端 | 机制 |
|---|---|
| Source | Kafka offset 快照 |
| 引擎 | Barrier 对齐快照 |
| Sink | 两阶段提交 |
Q2:Barrier 对齐是什么?
Checkpoint 标记随数据流动:
多输入流算子收到所有输入的 Barrier 才快照
→ 状态与数据位置一致| 对齐 | 保证 |
|---|---|
| 有 | 精确一次 |
| 无(至少一次模式) | 状态可能不一致 |
Q3:Checkpoint 和 Savepoint 区别?
| 对比 | Checkpoint | Savepoint |
|---|---|---|
| 触发 | 自动周期 | 手动 |
| 生命周期 | 自动清理 | 持久 |
| 用途 | 故障恢复 | 升级/迁移 |
| 兼容 | 版本内 | 跨版本 |
Q4:重启策略有哪些?
| 策略 | 行为 |
|---|---|
| 固定延迟 | 固定次数 + 间隔 |
| 失败率 | 时间窗口失败上限 |
| 无 | 失败即停 |
| 指数退避 | 间隔递增 |
二、状态管理
Q5:Flink 状态分哪两类?
| 类型 | 作用域 |
|---|---|
| 键控状态 | 每 key 一份 |
| 算子状态 | 每并行实例一份 |
Q6:状态类型有哪些?
| 状态 | 结构 |
|---|---|
| ValueState | 单值 |
| ListState | 列表 |
| MapState | Map |
| ReducingState | 归约 |
| AggregatingState | 聚合 |
Q7:状态后端如何选型?
| 后端 | 适用 |
|---|---|
| Memory | 测试 |
| Fs | 中等状态 |
| RocksDB | 大状态 + 增量快照 |
Q8:状态无限增长怎么办?
1. 设置状态 TTL
2. 及时清理(clear)
3. 预聚合,少存明细
4. 窗口状态及时清除三、窗口与时间
Q9:Flink 三种时间?
| 时间 | 含义 |
|---|---|
| Event Time | 事件产生时间 |
| Processing Time | 处理时刻 |
| Ingestion Time | 进入系统时刻 |
Q10:Watermark 是什么?
水位线 = 最大事件时间 - 容忍延迟
语义:早于水位线的事件不会再到达
作用:触发窗口、声明完整、清理状态Q11:乱序数据如何处理?
| 手段 | 说明 |
|---|---|
| 水位线 | 容忍乱序 |
| Allowed Lateness | 窗口延迟允许 |
| 侧输出 | 超晚数据单独处理 |
Q12:窗口类型?
| 窗口 | 特征 |
|---|---|
| 滚动 | 无重叠 |
| 滑动 | 可重叠 |
| 会话 | 按间隙 |
| 全局 | 全数据 |
Q13:窗口为什么延迟输出结果?
水位线太保守 → 窗口关闭晚
allowedLateness 长 → 结果更新晚
需要权衡延迟与准确四、反压
Q14:什么是反压?
下游处理慢 → 缓冲占满 → 阻塞上游 → 逐级反压到 SourceQ15:如何定位反压?
| 方法 | 说明 |
|---|---|
| UI 反压面板 | 各任务反压等级 |
| 指标 | 缓冲占用 |
| 日志 | 积压线索 |
Q16:反压怎么解决?
1. 定位反压源头(瓶颈算子)
2. 优化瓶颈:计算/IO/外部系统
3. 扩容并行度
4. 检查数据倾斜五、CDC
Q17:CDC 是什么?
捕获数据库变更(binlog/WAL)成流
插入/更新/删除 → 变更事件Q18:Flink CDC 全量+增量怎么实现?
initial 模式:
全量快照(并行分块)
记录位点
位点后增量,无缝衔接Q19:CDC 与 Debezium 区别?
| 对比 | Flink CDC | Debezium |
|---|---|---|
| 形态 | 连接器 | 独立工具 |
| 链路 | 直连 Flink | 通常经 Kafka |
| 部署 | 内嵌 | 独立服务 |
六、场景设计类
Q20:实时数仓架构怎么设计?
采集:Flink CDC(MySQL → Kafka)
加工:Flink SQL(Kafka → 实时分层)
存储:Doris/ClickHouse/Iceberg
指标:实时聚合输出Q21:实时指标去重怎么做?
| 方案 | 说明 |
|---|---|
| 状态去重 | MapState 记录已见 key |
| 精确去重 | 明细状态 |
| 近似去重 | BloomFilter/HLL |
Q22:作业重启如何快速恢复?
1. 从 Checkpoint/Savepoint 恢复
2. RocksDB 增量 + 本地恢复
3. 检查重启策略
4. 验证状态一致性Q23:大促洪峰如何应对?
1. 提前扩容(并行度)
2. 削峰:限速/背压
3. 降级:部分指标简化
4. Kafka 缓冲峰值七、高频易错点
| 易错点 | 正确理解 |
|---|---|
| 精确一次 ≠ 仅引擎 | 需要源+Sink 配合 |
| 水位线不推进 | 数据无事件时间/空闲分区 |
| 状态默认不清理 | 需 TTL/主动清理 |
| 反压 ≠ 故障 | 是下游慢的信号 |
| checkpoint 频繁 | 增加负担,合理间隔 |
| 并行度随意改 | 状态恢复需兼容 |
八、速答清单
| 高频题 | 一句话答案 |
|---|---|
| 精确一次怎么实现 | 快照 offset + Barrier 对齐 + 两阶段提交 |
| 水位线作用 | 定义乱序容忍、触发窗口 |
| 状态分类 | 键控状态 + 算子状态 |
| 后端选型 | 小状态 Fs,大状态 RocksDB |
| 反压是啥 | 下游慢导致逐级阻塞 |
| 乱序处理 | 水位线 + allowedLateness + 侧输出 |
| CDC 链路 | binlog → Flink → 数仓/消息 |