需求概述
当前 TIS 启动 Flink 增量任务后,如果 Flink 集群端因故障宕机重启,TIS 端记录的 JobID 会失效,用户必须手动删除旧任务并重新创建。期望引入自动恢复机制,使 Flink 集群重启后 TIS 能够自动检测并恢复任务连接。
当前痛点
- 手动操作繁琐:集群重启后需人工删除 → 重新创建 → 手动选择恢复点
- 状态不可见:TIS 端只能看到任务"消失",无法区分是集群重启还是任务失败
- 恢复点选择困难:用户需手动决定从哪个 Savepoint/Checkpoint 恢复,容易出错
- 缺乏自动化:无自动重试/恢复机制,完全依赖人工介入
目标场景
| 场景 |
说明 |
期望行为 |
| Flink HA 恢复 |
JM 崩溃后 HA 机制自动拉起 |
TIS 自动重新连接,更新 JobID |
| 集群完全重启 |
整个集群停止后重新启动 |
TIS 自动用最新恢复点重新提交任务 |
| TIS Console 重启 |
Console 服务重启 |
启动后自动恢复对 Flink Job 的监控 |
参考实现:StreamPark
已调研 StreamPark 的恢复机制(源码位于 /opt/misc/streampark),核心经验:
- 状态持久化到数据库 ——
t_flink_app 表存储 JobID、状态、tracking 标记
- 主动轮询检测 ——
FlinkAppHttpWatcher 每 5 秒轮询,发现 FAILED 自动重投
- 自动 Savepoint 恢复 —— 自动选择
t_flink_savepoint 中 latest=true 的路径重新提交
- 手动 Mapping 兜底 —— 提供
mapping 接口手动认领运行中的 Job
设计方案
详见 design/flink-job-restore/ 目录下的完整设计文档(6 个 markdown 文件):
- Phase 1(1周):基础设施 —— 数据库表 +
FlinkJobStatusStore 存储层 + 与现有 incrJob.log 双写兼容
- Phase 2(1.5周):核心恢复逻辑 ——
FlinkJobRestoreManager + Watcher 轮询 + 自动重投策略
- Phase 3(1周):恢复点管理 ——
CheckpointCollector 定时收集 + 最新恢复点自动标记
- Phase 4(1.5周):集成测试 —— REST 接口 + 前端适配 + 端到端测试
核心组件
FlinkJobRestoreManager(新增)
├── RestoreStrategy —— 自动恢复策略(重启次数限制、间隔控制)
├── JobIdMapper —— JobID 变化后通过 JobName 匹配新 Job
└── CheckpointCollector —— 定时收集 Checkpoint/Savepoint 并落库
关键状态流转
RUNNING → LOST/FAILED → RESTARTING → RUNNING(新 JobID)
↓
达到 maxRestarts → TERMINATED
数据库设计(新增三张表)
tis_flink_job_status —— Job 状态主表(JobID、状态、tracking、restartCount)
tis_flink_restore_point —— 恢复点表(Savepoint/Checkpoint 路径、latest 标记)
tis_flink_restore_history —— 恢复历史审计表
完整 DDL 见 design/flink-job-restore/05-数据库设计.md
兼容性
- 向后兼容:保留现有
incrJob.log 机制,数据库层为可选增强
- 可开关:通过配置
flink.job.restore.enabled 全局控制
- 可回滚:数据库方案异常时可降级到纯文件存储
实施路线图
| 阶段 |
工期 |
关键产出 |
| Phase 1:基础设施 |
1 周 |
存储层 + 双写兼容 + 数据迁移 |
| Phase 2:核心恢复逻辑 |
1.5 周 |
RestoreManager + Watcher + 策略 |
| Phase 3:恢复点管理 |
1 周 |
CheckpointCollector + 自动落库 |
| Phase 4:集成测试 |
1.5 周 |
接口 + 前端 + 端到端测试 |
| 总计 |
~4-6 周 |
|
核心结论
StreamPark 也无法在"集群完全重启"后保留原 JobID 的魔法连接。其优雅之处在于检测到不可达后,自动用最新的 Savepoint 重新提交一个等价的新 Job。
TIS 应借鉴此思路,重点实现:
- 状态从本地文件升级到数据库持久化
- Watcher 主动轮询检测 + 自动恢复
- Savepoint/Checkpoint 的自动收集和选择
- 手动 Mapping 作为兜底机制
如果你需要,我还可以把这个内容写入 design/flink-job-restore/07-GitHub-Issue.md 文件中。
需求概述
当前 TIS 启动 Flink 增量任务后,如果 Flink 集群端因故障宕机重启,TIS 端记录的
JobID会失效,用户必须手动删除旧任务并重新创建。期望引入自动恢复机制,使 Flink 集群重启后 TIS 能够自动检测并恢复任务连接。当前痛点
目标场景
参考实现:StreamPark
已调研 StreamPark 的恢复机制(源码位于
/opt/misc/streampark),核心经验:t_flink_app表存储 JobID、状态、tracking标记FlinkAppHttpWatcher每 5 秒轮询,发现FAILED自动重投t_flink_savepoint中latest=true的路径重新提交mapping接口手动认领运行中的 Job设计方案
详见
design/flink-job-restore/目录下的完整设计文档(6 个 markdown 文件):FlinkJobStatusStore存储层 + 与现有incrJob.log双写兼容FlinkJobRestoreManager+ Watcher 轮询 + 自动重投策略CheckpointCollector定时收集 + 最新恢复点自动标记核心组件
FlinkJobRestoreManager(新增)
├── RestoreStrategy —— 自动恢复策略(重启次数限制、间隔控制)
├── JobIdMapper —— JobID 变化后通过 JobName 匹配新 Job
└── CheckpointCollector —— 定时收集 Checkpoint/Savepoint 并落库
关键状态流转
RUNNING → LOST/FAILED → RESTARTING → RUNNING(新 JobID)
↓
达到 maxRestarts → TERMINATED
数据库设计(新增三张表)
tis_flink_job_status—— Job 状态主表(JobID、状态、tracking、restartCount)tis_flink_restore_point—— 恢复点表(Savepoint/Checkpoint 路径、latest 标记)tis_flink_restore_history—— 恢复历史审计表完整 DDL 见
design/flink-job-restore/05-数据库设计.md兼容性
incrJob.log机制,数据库层为可选增强flink.job.restore.enabled全局控制实施路线图
核心结论
TIS 应借鉴此思路,重点实现:
如果你需要,我还可以把这个内容写入 design/flink-job-restore/07-GitHub-Issue.md 文件中。