Flink Savepoint 状态恢复

FreeGuideOnline 最新 2026-07-13

bash

触发 Savepoint,保存到指定目录

bin/flink savepoint [targetDirectory]


**示例**:
```bash
# JobId 为 c8f3b7a0e4d2c1b5f3a0e4d2c1b5f3a0
bin/flink savepoint c8f3b7a0e4d2c1b5f3a0e4d2c1b5f3a0 hdfs:///flink/savepoints

命令会返回 Savepoint 的完整路径,例如 hdfs:///flink/savepoints/savepoint-c8f3b7-1a2b3c4d5e

2.2 通过 REST API 触发

向 JobManager 发送 POST 请求:

curl -X POST http://<jobmanager-host>:8081/jobs/<jobId>/savepoints \
  -d '{"target-directory":"hdfs:///flink/savepoints","cancel-job":"false"}'

cancel-job 控制是否在 Savepoint 完成后停止作业。

2.3 带取消的 Savepoint

如果希望在作业停止时生成 Savepoint,可使用:

bin/flink cancel -s [targetDirectory] <jobId>

此命令会优雅停止作业并持久化状态,常用于计划停机。


三、从 Savepoint 恢复作业

恢复作业的核心是指定 Savepoint 路径,Flink 会初始化所有算子的状态,然后继续消费数据。

3.1 使用命令行提交并恢复

在提交作业时通过 --fromSavepoint 参数指定路径。

bin/flink run \
  --fromSavepoint hdfs:///flink/savepoints/savepoint-c8f3b7-1a2b3c4d5e \
  -c com.example.MyJob \
  my-app.jar

参数说明

  • --fromSavepoint:Savepoint 的绝对路径。
  • 其他参数与常规提交一致(主类、jar 包、运行参数等)。

3.2 在代码中恢复(不常用)

你也可以在 StreamExecutionEnvironment 中设置恢复路径:

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 启用 Checkpoint 以便后续 Savepoint 可用
env.enableCheckpointing(5000);
env.getCheckpointConfig().setCheckpointStorage("hdfs:///flink/checkpoints");

// 不需要在代码中指定 Savepoint,而是提交时命令行指定

推荐做法:尽量使用命令行参数,保持代码无侵入。

3.3 使用 SQL Client 恢复

若使用 Flink SQL 提交作业,可通过如下方式:

-- 设置状态后端(若需要)
SET 'execution.savepoint.path' = 'hdfs:///flink/savepoints/savepoint-xxxx';
-- 提交 SQL 作业
INSERT INTO sink_table SELECT ... ;

或者直接在启动 SQL Client 时通过 -s 指定:

sql-client.sh -s hdfs:///flink/savepoints/savepoint-xxxx

3.4 恢复时的并行度调整

Savepoint 保存了每个算子的 最大并行度(Max Parallelism)设置。恢复时允许修改作业的并行度,但不能超过该算子的最大并行度。

  • 缩减并行度env.setParallelism(4) 若原为 8,则状态会自动重新分发。
  • 扩大并行度:只能扩大到不超过 maxParallelism

若作业未显式设置 maxParallelism,默认值为 128(或 2 的幂),一般足够。


四、状态后端与恢复的兼容性

4.1 选择合适的状态后端

Savepoint 的状态存储依赖于你配置的 状态后端

  • HashMapStateBackend(默认):内存存储,Savepoint 文件全量写入文件系统。
  • EmbeddedRocksDBStateBackend:磁盘 + 内存,适合超大状态,Savepoint 速度较慢但稳定。

恢复时使用与创建 Savepoint 时相同的状态后端(或兼容的)即可,Flink 自动识别。

4.2 恢复至不同状态后端

从 HashMap 后端创建的 Savepoint 可以恢复到 RocksDB 后端,反之亦然。但注意性能差异,建议保持一致。

4.3 处理状态序列化升级

若你修改了算子的状态数据类型(如 Avro schema 升级),可能需要自定义 TypeSerializerSnapshot 或使用 State Processor API 迁移状态。对于普通业务改动(如增加字段),Flink 的状态恢复会尝试使用原序列化器反序列化,若失败则抛出无法恢复的异常。


五、算子 UID——恢复成功的关键

Flink 使用 算子 UID 将 Savepoint 中的状态映射到新作业图上的对应算子。若未设置 UID,Flink 会根据作业拓扑自动生成 ID,这些 ID 在作业修改后极易变化,导致恢复失败。

5.1 正确做法:为每个有状态的算子设置 UID

DataStream<String> stream = env
    .addSource(new MySource())
    .uid("my-source")          // 设置 UID
    .keyBy(e -> e)
    .process(new MyProcessFunction())
    .uid("my-processor")
    .print()
    .uid("my-sink");

5.2 恢复时算子匹配规则

  • UID 完全相同 → 直接恢复状态。
  • UID 不同但拓扑相同 → 可尝试启动但状态清空(丢失所有状态)。
  • UID 匹配但状态 schema 不兼容 → 抛出异常。

最佳实践:从项目第一天起就为每个有状态的算子分配稳定的 UID,即使后续重构也要保持不变。


六、实战演示:从 Savepoint 恢复一个修改后的作业

假设我们有一个实时统计订单量的作业,现在需要修改聚合逻辑并增加一个并行度。

6.1 步骤一:为原始作业设置 UID 并部署

env.setMaxParallelism(256);
DataStream<Order> orders = env
    .addSource(new OrderSource())
    .uid("order-source")
    .keyBy(Order::getCategory)
    .window(TumblingProcessingTimeWindows.of(Time.minutes(1)))
    .aggregate(new OrderAggregateFunction())
    .uid("order-window-agg")
    .print()
    .uid("sink-print");

提交后获得 JobId a1b2c3d4e5f6

6.2 步骤二:触发 Savepoint 并保存路径

bin/flink savepoint a1b2c3d4e5f6 hdfs:///flink/savepoints
# 返回: hdfs:///flink/savepoints/savepoint-a1b2c3-7a8b9c0d

6.3 步骤三:修改作业

  • 更改聚合函数 OrderAggregateFunction -> AdvancedOrderAggregateFunction
  • 扩大并行度,例如 env.setParallelism(4)(原为 2,且不超过 maxParallelism)。
  • 保持所有 UID 不变。

6.4 步骤四:从 Savepoint 恢复新作业

bin/flink run \
  --fromSavepoint hdfs:///flink/savepoints/savepoint-a1b2c3-7a8b9c0d \
  --parallelism 4 \
  -c com.example.AdvancedOrderJob \
  advanced-order-app.jar