Skip to content

Flink Job 自动恢复功能 #498

Description

@baisui1981

需求概述

当前 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),核心经验:

  1. 状态持久化到数据库 —— t_flink_app 表存储 JobID、状态、tracking 标记
  2. 主动轮询检测 —— FlinkAppHttpWatcher 每 5 秒轮询,发现 FAILED 自动重投
  3. 自动 Savepoint 恢复 —— 自动选择 t_flink_savepointlatest=true 的路径重新提交
  4. 手动 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 应借鉴此思路,重点实现:

  1. 状态从本地文件升级到数据库持久化
  2. Watcher 主动轮询检测 + 自动恢复
  3. Savepoint/Checkpoint 的自动收集和选择
  4. 手动 Mapping 作为兜底机制

如果你需要,我还可以把这个内容写入 design/flink-job-restore/07-GitHub-Issue.md 文件中。

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions