如何实现Flink的Checkpoint/Savepoint跨多数据中心备份
Flink跨数据中心双写Checkpoint/Savepoint可行方案
完全可以通过Flink自身的插件化机制实现状态双写,不需要依赖存储层的跨域同步能力,有两种成熟方案可选:
方案1:自定义FileSystem实现双写逻辑(兼容所有Flink版本)
Flink的Checkpoint/Savepoint存储基于抽象的FileSystem接口实现,你可以自行开发双写包装类适配低版本需求:
- 实现Flink的
org.apache.flink.core.fs.FileSystem接口,内部同时持有DC1和DC2对应存储(HDFS/S3)的FileSystem实例 - 所有写入类方法(create、rename、mkdirs等)同步调用两个实例的对应方法,确保两边写入一致性
- 读取类方法默认优先读取DC1实例,DC1故障时可自动切到DC2读取
- 将自定义实现打包为jar放入Flink的lib目录,在flink-conf.yaml中配置状态存储URI为自定义FileSystem的专属scheme即可生效
注意:写入策略可根据业务容错要求调整,可选「至少一个DC写成功就判定Checkpoint成功」或「所有DC都写成功才算完成」,避免单DC临时波动导致Checkpoint频繁失败
方案2:使用Flink 1.15+ 原生多副本Checkpoint功能(无需二次开发)
如果你的Flink版本≥1.15,官方已经内置多存储路径支持,直接配置即可使用:
- 在flink-conf.yaml中配置多个状态存储路径,路径之间用分号分隔,示例配置如下:
state.checkpoint-storage: filesystem state.checkpoints.dir: s3://dc1-bucket/flink/checkpoints;s3://dc2-bucket/flink/checkpoints state.savepoints.dir: hdfs://dc1-namespace/flink/savepoints;hdfs://dc2-namespace/flink/savepoints # 可选配置:强制要求所有路径都写入成功才算Checkpoint完成,默认关闭 # state.checkpoints.multiple-storages.all-required: true
- Flink会自动将状态数据同步写入所有配置的路径,恢复时按配置顺序自动查找可用的最新Checkpoint,只要有一个DC的状态完整就能正常恢复作业
落地注意事项
- 跨DC网络延迟较高时会拉长Checkpoint同步阶段耗时,建议先压测跨DC写入带宽、延迟是否满足你的Checkpoint间隔要求
- 两个方案都同时支持Checkpoint和Savepoint的双写,无需额外适配
- 大状态作业优先选官方原生的多副本功能,稳定性比自定义实现更有保障
内容的提问来源于stack exchange,提问作者shanker861
相关产品推荐
相关产品推荐

