Flink Savepoint 状态恢复
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