全部笔记All notes

Flink Checkpoint & Savepoint 数据一致性保障

阅读 8m 45s8m 45s read

概述

在 Flink 流计算中,故障恢复是一个关键问题。Flink 通过 Checkpoint 和 Savepoint 机制提供 Exactly-Once 语义,保证任务在出现故障后仍然能够恢复到最近的状态,避免数据丢失或重复计算。

💡 提示: Checkpoint 和 Savepoint 都是 Flink 容错机制的核心组成部分,但它们的使用场景和特点有所不同。

本文重点

  • Checkpoint: 自动周期性状态快照,用于故障恢复
  • Savepoint: 手动触发的状态快照,用于作业升级和迁移
  • 配置优化: 在生产环境中的最佳实践
  • 性能调优: 如何平衡容错性和性能

Checkpoint 机制

什么是 Checkpoint

Checkpoint 是 Flink 自动进行的周期性状态快照,用于故障恢复。当任务失败时,Flink 会使用最近的 Checkpoint 自动恢复状态,保证计算的正确性。

特点
  • 自动执行: 不需要人工触发
  • 定期保存: 默认 10s 触发一次
  • 故障恢复: 作业失败时自动从最新 Checkpoint 恢复
  • 分布式存储: 支持 HDFS、S3、RocksDB 等

⚠️ 注意: Checkpoint 的频率会影响作业性能,需要根据实际场景进行权衡。

启用和配置

基本配置

可以在 Flink Job 中显式开启 Checkpoint:

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.CheckpointingMode;
import org.apache.flink.streaming.api.environment.CheckpointConfig;

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// 启用 Checkpoint,每 5 秒做一次
env.enableCheckpointing(5000);

// 获取 Checkpoint 配置
CheckpointConfig checkpointConfig = env.getCheckpointConfig();
关键参数配置

Flink 提供了丰富的 Checkpoint 配置项来优化性能:

// 保障 Exactly-Once 语义(默认)
checkpointConfig.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);

// Checkpoint 超时时间 60 秒
checkpointConfig.setCheckpointTimeout(60000);

// 两次 Checkpoint 之间最小间隔 5 秒
checkpointConfig.setMinPauseBetweenCheckpoints(5000);

// 同时只允许 1 个 Checkpoint 进行
checkpointConfig.setMaxConcurrentCheckpoints(1);

// 任务取消后保留 Checkpoint
checkpointConfig.enableExternalizedCheckpoints(
    CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION
);

// 容忍 Checkpoint 失败次数
checkpointConfig.setTolerableCheckpointFailureNumber(3);

🔧 配置建议:

  • 小型任务: 5-10 秒
  • 大型任务 (TB 级状态): 30-60 秒

存储后端选择

Flink 提供多种状态存储后端,适用于不同场景:

1. MemoryStateBackend
  • 特点: 存储在内存中
  • 适用场景: 本地测试、小状态作业
  • 限制: 状态大小受限于内存
import org.apache.flink.runtime.state.memory.MemoryStateBackend;

env.setStateBackend(new MemoryStateBackend());
2. FsStateBackend
  • 特点: 存储在文件系统 (HDFS/S3)
  • 适用场景: 生产环境、中等规模状态
  • 优势: 支持大规模状态
import org.apache.flink.runtime.state.filesystem.FsStateBackend;

env.setStateBackend(new FsStateBackend("hdfs://namenode:8020/flink-checkpoints"));
3. RocksDBStateBackend
  • 特点: 基于 RocksDB 的持久化存储
  • 适用场景: 超大规模状态 (TB 级)
  • 优势: 支持增量 Checkpoint
import org.apache.flink.contrib.streaming.state.RocksDBStateBackend;

RocksDBStateBackend rocksDBBackend = new RocksDBStateBackend("hdfs://namenode:8020/flink-checkpoints");
// 启用增量 Checkpoint
rocksDBBackend.setEnableIncrementalCheckpointing(true);
env.setStateBackend(rocksDBBackend);

💡 选择建议:

  • 测试环境: MemoryStateBackend
  • 中小型作业: FsStateBackend
  • 大型作业: RocksDBStateBackend

高级配置

非对齐 Checkpoint

Flink 1.11+ 支持非对齐 Checkpoint,可以显著减少背压情况下的 Checkpoint 时间:

import java.time.Duration;

// 启用非对齐 Checkpoint
checkpointConfig.enableUnalignedCheckpoints();

// 设置对齐超时,超时后自动切换到非对齐
checkpointConfig.setAlignmentTimeout(Duration.ofSeconds(30));
两阶段提交

对于需要端到端 Exactly-Once 的场景:

// 启用两阶段提交
env.getCheckpointConfig().enableExternalizedCheckpoints(
    CheckpointConfig.ExternalizedCheckpointCleanup.DELETE_ON_CANCELLATION
);

// 设置事务超时时间
env.getCheckpointConfig().setCheckpointTimeout(900000); // 15 分钟

Savepoint 机制

什么是 Savepoint

Savepoint 是手动触发的 Flink 任务状态快照,主要用于:

  • 作业升级: 更新代码后无缝重启
  • 集群迁移: 从一个集群迁移到另一个集群
  • 数据备份: 手动保留状态,防止数据丢失
  • A/B 测试: 从同一状态启动不同版本
特点
  • 手动触发: 需要显式执行
  • 持久存储: 存储在用户指定的位置
  • 版本升级: 可用于任务升级和迁移
  • 长期保留: 不会自动删除,需要手动管理

⚠️ 注意: Savepoint 可能会比 Checkpoint 大,因为它需要保存完整的状态信息。

触发和管理

触发 Savepoint

Savepoint 只能手动执行,使用 Flink CLI 命令:

# 基本命令
flink savepoint <jobID> <savepointPath>

# 示例
flink savepoint 1a2b3c4d5e6f hdfs://namenode:8020/flink-savepoints/

# 使用 YARN 作业管理器
flink savepoint <jobID> hdfs://namenode:8020/flink-savepoints/ -yid <yarnAppId>
停止作业并创建 Savepoint
# 优雅停止并创建 Savepoint
flink stop --savepointPath hdfs://namenode:8020/flink-savepoints/ <jobID>

# 取消作业并创建 Savepoint
flink cancel -s hdfs://namenode:8020/flink-savepoints/ <jobID>
管理 Savepoint
# 列出所有 Savepoint
hdfs dfs -ls hdfs://namenode:8020/flink-savepoints/

# 删除旧的 Savepoint
flink savepoint dispose hdfs://namenode:8020/flink-savepoints/savepoint-xxx

💡 提示: 定期清理旧的 Savepoint 可以节省存储空间。

恢复作业

从 Savepoint 启动
# 基本命令
flink run -s <savepointPath> <jarFile>

# 完整示例
flink run -s hdfs://namenode:8020/flink-savepoints/savepoint-1234 \
    -c com.example.StreamingJob \
    my-flink-job.jar
允许非恢复状态

当作业更新后某些算子被删除或修改时:

# 允许跳过无法恢复的算子
flink run -s hdfs://namenode:8020/flink-savepoints/savepoint-1234 \
    --allowNonRestoredState \
    -c com.example.StreamingJob \
    my-flink-job.jar
编程式恢复
import org.apache.flink.api.java.utils.ParameterTool;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

public class SavepointRecovery {
    public static void main(String[] args) throws Exception {
        ParameterTool params = ParameterTool.fromArgs(args);
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        // 从参数中获取 Savepoint 路径
        String savepointPath = params.get("savepoint");
        if (savepointPath != null) {
            env.setStateBackend(new RocksDBStateBackend(savepointPath));
        }
        
        // 构建作业逻辑
        // ...
        
        env.execute("Restored Job");
    }
}

两者对比

Checkpoint vs Savepoint 详细对比

对比项CheckpointSavepoint
触发方式自动(定期触发)手动(用户主动触发)
主要用途故障恢复作业升级/迁移
存储方式增量存储(RocksDB)完整存储
生命周期默认删除(可配置保留)始终保留
性能影响轻量级可能较重
兼容性仅同版本跨版本兼容
适用场景保证计算任务的高可用性任务升级、版本迁移、数据保护

💡 总结:

  • Checkpoint 适用于故障恢复,不适用于任务升级
  • Savepoint 适用于任务升级和迁移

使用场景选择

使用 Checkpoint 的场景
  • 作业 7x24 小时运行,需要自动故障恢复
  • 对恢复时间有要求(分钟级)
  • 不需要频繁升级代码
使用 Savepoint 的场景
  • 计划内的代码升级
  • 集群迁移或维护
  • 需要长期保存某个时刻的状态
  • A/B 测试或灰度发布

生产环境最佳实践

结合使用 Checkpoint & Savepoint

在生产环境下,通常采用 Checkpoint + Savepoint 组合,既保证任务的自动恢复,又允许手动升级:

  1. 启用 Checkpoint,确保任务崩溃后可以自动恢复
  2. 定期手动 Savepoint,用于版本升级和数据保护
  3. 存储在 HDFS/RocksDB,保证状态持久化
完整配置示例
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.CheckpointingMode;
import org.apache.flink.streaming.api.environment.CheckpointConfig;
import org.apache.flink.contrib.streaming.state.RocksDBStateBackend;

public class ProductionJobConfig {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        // 1. 启用 Checkpoint
        env.enableCheckpointing(30000);  // 30 秒一次
        
        CheckpointConfig checkpointConfig = env.getCheckpointConfig();
        checkpointConfig.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
        checkpointConfig.setMinPauseBetweenCheckpoints(10000);
        checkpointConfig.setCheckpointTimeout(120000); // 2 分钟超时
        checkpointConfig.setMaxConcurrentCheckpoints(1);
        
        // 2. 允许外部 Checkpoint,便于 Savepoint
        checkpointConfig.enableExternalizedCheckpoints(
            CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION
        );
        
        // 3. 使用 RocksDB 后端
        RocksDBStateBackend backend = new RocksDBStateBackend(
            "hdfs://namenode:8020/flink-checkpoints"
        );
        backend.setEnableIncrementalCheckpointing(true);
        env.setStateBackend(backend);
        
        // 4. 作业逻辑
        // ...
        
        env.execute("Production Job");
    }
}

⚠️ 重要: 这样配置后:

  • 如果任务失败,Flink 会自动从最近的 Checkpoint 恢复
  • 如果任务升级,可以手动 Savepoint,并从 Savepoint 重新启动

大规模生产环境优化

在生产环境下,Flink 任务可能会处理 TB 级别状态数据,需要特别优化 Checkpoint 及 Savepoint:

✅ 1. 使用 RocksDBStateBackend

状态数据很大时,必须使用 RocksDB:

RocksDBStateBackend backend = new RocksDBStateBackend("hdfs://flink-checkpoints");
// 启用增量 Checkpoint
backend.setEnableIncrementalCheckpointing(true);
// 设置本地 RocksDB 目录
backend.setDbStoragePath("/mnt/ssd/rocksdb");
env.setStateBackend(backend);
✅ 2. 增加 Checkpoint 存储路径

可以配置多个 HDFS 存储路径,提高稳定性:

# flink-conf.yaml
state.backend.fs.checkpointdir: hdfs://namenode:8020/flink-checkpoints
state.backend.rocksdb.localdir: /mnt/flink-rocksdb
state.checkpoints.num-retained: 3
✅ 3. 合理调整 Checkpoint 频率

Flink Checkpoint 需要占用 CPU 资源,建议:

  • 小型任务: 5~10 秒
  • 中型任务: 10~30 秒
  • 大型任务 (TB 级状态): 30~60 秒
// 大型任务配置
env.enableCheckpointing(45000);  // 45 秒
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(15000); // 间隔至少 15 秒
✅ 4. 定期 Savepoint 备份
  • 每周或任务升级前手动 Savepoint
  • 存储在 HDFS/S3,防止任务取消后数据丢失
# 定期执行脚本
#!/bin/bash
JOB_ID="your-job-id"
SAVEPOINT_DIR="hdfs://namenode:8020/flink-savepoints/weekly"
DATE=$(date +%Y%m%d_%H%M%S)

flink savepoint $JOB_ID $SAVEPOINT_DIR/savepoint_$DATE
✅ 5. 监控 Checkpoint 指标

监控以下关键指标:

  • Checkpoint 持续时间: 不应超过间隔时间的 50%
  • Checkpoint 大小: 监控增长趋势
  • Checkpoint 失败率: 应低于 1%
// 通过 Metrics 监控
env.getCheckpointConfig().setFailOnCheckpointingErrors(false);
env.getCheckpointConfig().setTolerableCheckpointFailureNumber(3);

性能优化

Checkpoint 性能调优

1. 使用增量 Checkpoint

对于大状态作业,增量 Checkpoint 可以显著减少 I/O:

RocksDBStateBackend backend = new RocksDBStateBackend("hdfs://checkpoints");
backend.setEnableIncrementalCheckpointing(true);
backend.setNumberOfTransferThreads(4); // 增加传输线程
2. 优化 RocksDB 配置
import org.apache.flink.contrib.streaming.state.RocksDBOptionsFactory;
import org.rocksdb.*;

public class OptimizedRocksDBConfig implements RocksDBOptionsFactory {
    @Override
    public DBOptions createDBOptions(DBOptions currentOptions, 
                                     Collection<AutoCloseable> handlesToClose) {
        return currentOptions
            .setMaxBackgroundJobs(4)
            .setMaxOpenFiles(-1)
            .setIncreaseParallelism(4)
            .setUseFsync(false);
    }
    
    @Override
    public ColumnFamilyOptions createColumnOptions(
            ColumnFamilyOptions currentOptions, 
            Collection<AutoCloseable> handlesToClose) {
        return currentOptions
            .setCompactionStyle(CompactionStyle.LEVEL)
            .setLevel0FileNumCompactionTrigger(10)
            .setMaxBytesForLevelBase(256 * 1024 * 1024)
            .setTargetFileSizeBase(64 * 1024 * 1024);
    }
}

// 应用配置
backend.setRocksDBOptions(new OptimizedRocksDBConfig());
3. 使用本地 SSD

将 RocksDB 本地目录配置在 SSD 上:

# flink-conf.yaml
state.backend.rocksdb.localdir: /mnt/ssd1/rocksdb,/mnt/ssd2/rocksdb
4. 调整网络缓冲
# 增加网络缓冲区
taskmanager.network.memory.fraction: 0.15
taskmanager.network.memory.min: 128mb
taskmanager.network.memory.max: 1gb

Savepoint 性能优化

1. 并行化 Savepoint
# 增加并行度
state.backend.fs.write-buffer-size: 4mb
savepoint.max-concurrent-checkpoints: 2
2. 压缩 Savepoint
// 启用压缩
env.getConfig().setUseSnapshotCompression(true);
3. 选择合适的存储
  • HDFS: 适合大多数场景
  • S3: 适合云环境,注意网络延迟
  • 本地 SSD: 仅用于测试

常见问题

Q1: Checkpoint 失败怎么办?

可能原因及解决方案:

  1. 超时: 增加 checkpointTimeout
  2. 背压: 启用非对齐 Checkpoint
  3. 存储问题: 检查 HDFS 空间和权限
  4. 内存不足: 切换到 RocksDB 后端
// 容忍失败
checkpointConfig.setTolerableCheckpointFailureNumber(5);
checkpointConfig.setFailOnCheckpointingErrors(false);

Q2: Savepoint 太大怎么优化?

  1. 启用压缩:
env.getConfig().setUseSnapshotCompression(true);
  1. 清理无用状态:
  • 设置状态 TTL
  • 定期清理过期数据
  1. 使用增量 Savepoint (Flink 1.15+)

Q3: 如何从旧版本 Savepoint 恢复?

  1. 检查兼容性: 查看 Flink 官方升级指南
  2. 使用 --allowNonRestoredState: 允许部分状态不恢复
  3. 渐进式升级: 逐个小版本升级

Q4: Checkpoint 和 Savepoint 可以相互转换吗?

  • Checkpoint 转 Savepoint: 可以,使用外部化 Checkpoint
  • Savepoint 转 Checkpoint: 不可以,但可以从 Savepoint 启动后生成新的 Checkpoint

Q5: 多个作业共享 Checkpoint 目录吗?

不建议。每个作业应使用独立目录:

hdfs://namenode:8020/flink-checkpoints/job1/
hdfs://namenode:8020/flink-checkpoints/job2/

总结

核心要点

机制作用触发方式适用场景
Checkpoint自动状态快照定期触发故障恢复
Savepoint手动任务快照手动触发任务升级、版本迁移

✅ 最佳实践总结

  1. 生产环境开启 Checkpoint,存储在 HDFS/RocksDB
  2. 任务升级时手动 Savepoint 以保留状态
  3. 优化 Checkpoint 频率,避免性能损耗
  4. 使用 RocksDB 处理大规模状态
  5. 监控 Checkpoint 指标,及时发现问题

🚫 错误:

  • 不要在生产环境使用 MemoryStateBackend
  • 不要设置过于频繁的 Checkpoint(< 5秒)
  • 不要忽略 Checkpoint 失败告警

相关文章

容错机制系列

  • Flink 状态管理详解 - 深入理解状态管理
  • Flink 两阶段提交 - 端到端 Exactly-Once
  • Flink 故障恢复 - 故障恢复机制

DataStream API 系列

进阶学习

上一章 / 下一章