Flink Checkpoint和大状态优化

Flink的checkpoint是一种实现容错和状态恢复的机制。Checkpoint会定期将Flink程序的状态保存到持久化存储系统中,通常是分布式文件系统,当Flink程序发生故障时,可以重新启动应用程序并从之前的状态中恢复,使得Flink程序能够回到故障发生前的一致状态。因此Flink中的checkpoint非常重要,本小节主要讲解Flink Checkpoint和大状态相关优化内容点。

image-20260706213500191

Checkpoint的监控

在Flink WebUI界面中提供了Flink Job的checkpoint监控信息。当Job结束后,这些信息仍然可用。如下图所示,有四个不同的选项卡(概览、历史记录、摘要信息、配置信息)可以显示checkpoint相关的信息,下面分别介绍这些内容

概览(Overview)选项卡

概览选项卡列出了以下统计信息。

  • Checkpoint Counts
    • Triggered:自作业开始以来触发的checkpoint总数。
    • In Progress:当前正在进行的checkpoint数量。
    • Completed:自作业开始以来成功完成的checkpoint总数。
    • Failed:自作业开始以来失败的checkpoint总数。
    • Restored:自作业开始以来进行的恢复操作的次数。这还表示自提交任务以来已重新启动多少次。请注意,带有 savepoint 的初始任务提交也算作一次恢复,如果 JobManager 在此操作过程中丢失,则该统计将重新计数。
  • Latest Completed Checkpoint:最新(最近)成功完成的 checkpoint。
  • Latest Failed Checkpoint:最新失败的 checkpoint。
  • Latest Savepoint:最新触发的 savepoint 及其外部路径。
  • Latest Restore:有两种类型的恢复操作。
    • Restore from Checkpoint:从checkpoint恢复。
    • Restore from Savepoint:从savepoint恢复。

image-20260706213750441

请注意,概览选项卡中的信息在JobManager丢失时无法保存,如果JobManager发生故障转移,这些信息将会被重置。

历史记录(History)选项卡

Checkpoint 历史记录保存有关最近触发的 checkpoint 的统计信息,包括当前正在进行的 checkpoint。对于失败的checkpoint,指标会尽最大努力进行更新,但是可能不准确。checkpoint统计详细信息内容如下:

  • ID:已触发 checkpoint 的 ID。每个 checkpoint 的 ID 都会递增,从 1 开始。
  • Status:Checkpoint 的当前状态,可以是正在进行(In Progress)、已完成(Completed) 或失败(Failed))。
  • Acknowledged:已确认完成的子任务数量与总任务数量。
  • Trigger Time:在 JobManager 上发起 checkpoint 的时间。
  • Latest Acknowledgement:JobManager 接收到任何 subtask 的最新确认的时间。
  • End to End Duration:从触发时间戳到最后一次确认的持续时间,也就是checkpoint开始到完成所需总时间。完整 checkpoint 的端到端持续时间由确认 checkpoint 的最后一个 subtask 确定。
  • Checkpointed Data Size: 在此次checkpoint的sync以及async阶段中持久化的数据量。如果启用了增量 checkpoint或者changelog,则此值可能会与全量checkpoint数据量产生区别。
  • Full Checkpoint Data Size: 所有已确认的 subtask 的 checkpoint 的全量数据大小。
  • Processed (persisted) in-flight data:在 checkpoint 对齐期间(从接收第一个和最后一个 checkpoint barrier 之间的时间)所有已确认的 subtask 处理/持久化 的大约字节数。如果启用了 unaligned checkpoint,持久化的字节数可能会大于0。

image-20260706215416972

点击每个checkpoint记录前的“+”符号,可以看到每个checkpoint中subtask详情,对于subtask的一些详细信息如下图所示,解释如下:

  • Sync Duration:Checkpoint 同步部分的持续时间。这包括 operator 的快照状态,并阻塞 subtask 上的所有其他活动(处理记录、触发计时器等)。
  • Async Duration:Checkpoint 的异步部分的持续时间。这包括将 checkpoint 写入设置的文件系统所需的时间。对于 unaligned checkpoint,这还包括 subtask 必须等待最后一个 checkpoint barrier 到达的时间(checkpoint alignment 持续时间)以及持久化数据所需的时间。
  • Alignment Duration:处理第一个和最后一个 checkpoint barrier 之间的时间。对于 checkpoint alignment 机制的 checkpoint,在 checkpoint alignment 过程中,已经接收到 checkpoint barrier 的 channel 将阻塞并停止处理后续的数据。
  • Start Delay:从 checkpoint barrier 创建开始到 subtask 收到第一个 checkpoint barrier 所用的时间。
  • Unaligned Checkpoint:Checkpoint 完成的时候是否是一个 unaligned checkpoint。在 alignment 超时的时候 aligned checkpoint 可以自动切换成 unaligned checkpoint。

image-20260706215740631

请注意:这些信息不会再JobManager中保存,如果JobManager故障转移,这些统计信息将重新计数。

配置信息(Configuration)选项卡

该选项卡展示了用户指定的配置。

  • Checkpointing Mode:恰好一次(Exactly Once)或者至少一次(At least Once)。
  • Interval:配置的 checkpoint 触发间隔。在此间隔内触发 checkpoint。
  • Timeout:超时之后,JobManager 取消 checkpoint 并触发新的 checkpoint。
  • Minimum Pause Between Checkpoints:Checkpoint 之间所需的最小暂停时间。Checkpoint 成功完成后,我们至少要等这段时间再触发下一个,这可能会延迟正常的间隔。
  • Maximum Concurrent Checkpoints:可以同时进行的最大 checkpoint 个数。
  • Persist Checkpoints Externally:启用或禁用持久化 checkpoint 到外部系统。如果启用,还会列出外部化 checkpoint 的清理配置(取消时删除或保留)。

Checkpoint优化

设置CheckpointStorage

我们可以设置Checkpoint检查点快照(SnapShot)存储的位置,可以选择JobManagerCheckpointStorage或者FileSystemCheckpointStorage,两者分别代表JobManager堆内存和文件系统。默认情况下,snapshot 存储在JobManager的堆内存中,建议改为持久化文件系统。

Flink设置Checkpoint storage检查点存储位置代码如下:

//设置checkpoint storage存储为JobManagerStorage,默认堆内存存储状态大小为5M
env.getCheckpointConfig().setCheckpointStorage(new JobManagerCheckpointStorage(5*1024*1024));

//设置checkpoint storage存储为hdfs路径
env.getCheckpointConfig().setCheckpointStorage("hdfs://mycluster/flink/checkpoints");

对于JobManagerCheckpointStorage来说,默认每个单独的状态大小限制为5M ,可以手动指定该值,如果Flink是本地开发调试或者Flink状态非常少的场景可以使用JobManagerCheckpointStorage,实际生产中推荐使用FileSystemCheckpointStorage。

设置Checkpoint模式

选择exactly-once语义保证整个应用内端到端的数据一致性,这种情况比较适合于数据要求比较高,不允许出现丢数据或者数据重复,与此同时,Flink的性能也相对较弱,而at-least-once语义更适合于时延和吞吐量要求非常高但对数据的一致性要求不高的场景。

Flink中通过setCheckpointingMode()方法来设置检查点模式,默认情况下使用的是exactly-once模式。

//设置检查点模式为exactly-once
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);

//设置检查点模式为at-least-once
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.AT_LEAST_ONCE);

设置checkpoint超时时间

超时时间指定了每次Checkpoint执行过程中的上限时间范围,一旦Checkpoint执行时间超过该阈值,Flink将会中断Checkpoint过程,并按照超时处理。该指标可以通过setCheckpointTimeout方法设定,默认为10分钟。

//设置Checkpoint 超时时间
env.getCheckpointConfig().setCheckpointTimeout(10*60*1000);

设置checkpoint之间最小等待时间

通过Webui观察Flink任务运行情况,如果checkpoint的完成时间经常超过checkpoints基本间隔时(例如:因为状态比计划的更大或者访问checkpoints所在的存储系统暂时变慢),Flink系统会不断地进行checkpoints,因为一旦checkpoint完成,新的checkpoints就会立即启动。

这种情况下意味着多过的资源被不断地束缚在checkpointing中,并且checkpoint算子进行得缓慢,checkpoint占用了大量计算资源而影响到整个应用的性能。为了防止这种情况,应用程序可以定义checkpoint之间的最小等待时间。

//设置 checkpoint 最小间隔时间为500ms,默认值为0
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(500);

除了以上设置方式外,还可以在提交Flink任务时进行“execution.checkpointing.min-pause”参数设置。

注意,当指定了该参数大于0时,Flink最大并行执行checkpoint的数量为1。

设置checkpoint并行度

如果Flink集群的资源充足,checkpoint周期时间较短,也可以配置Flink 应用程序同时进行多个checkpoints同时进行。这种情况下会分配更多的资源到checkpointing。Flink在默认情况下只有一个检查点可以运行,根据用户指定的数量可以同时触发多个Checkpoint,进而提升Checkpoint整体的效率。

//设置checkpoint最大并行度,默认为1
env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);

以上设置方式还可以在提交Flink任务的时候通过“execution.checkpointing.max-concurrent-checkpoints”参数进行设置。

如下是测试checkpoint并行执行的代码案例。该案例中读取自定义source中的数据并通过map处理,每条数据暂停处理0.5秒。案例代码如下:

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

DataStreamSource<StationLog> source = env.addSource(new RichParallelSourceFunction<StationLog>() {
    Boolean flag = true;

    /**
     * 主要方法:启动一个Source,大部分情况下都需要在run方法中实现一个循环产生数据
     * 这里计划1s 产生1条基站数据,由于是并行,当前节点有几个core就会有几条数据
     */
    @Override
    public void run(SourceContext<StationLog> ctx) throws Exception {
        Random random = new Random();
        String[] callTypes = {"fail", "success", "busy", "barring"};
        while (flag) {
            String sid = "sid_" + random.nextInt(10);
            String callOut = "1811234" + (random.nextInt(9000) + 1000);
            String callIn = "1915678" + (random.nextInt(9000) + 1000);
            String callType = callTypes[random.nextInt(4)];
            Long callTime = System.currentTimeMillis();
            Long durations = Long.valueOf(random.nextInt(50) + "");
            ctx.collect(new StationLog(sid, callOut, callIn, callType, callTime, durations));
            Thread.sleep(1000);//1s 产生一个事件
        }

    }

    //当取消对应的Flink任务时被调用
    @Override
    public void cancel() {
        flag = false;
    }
});

//处理数据
SingleOutputStreamOperator<String> result =
        source.keyBy(stationLog -> stationLog.sid)
        .map(new RichMapFunction<StationLog, String>() {

    private ListState<String> listState;

    @Override
    public void open(Configuration parameters) throws Exception {
        ListStateDescriptor<String> stateDescriptor = new ListStateDescriptor<>("liststate", String.class);
        listState = getRuntimeContext().getListState(stateDescriptor);
    }

    @Override
    public String map(StationLog value) throws Exception {
        //100倍状态存储
        for(int i = 0 ;i<100 ;i++){
            listState.add(value.toString());
        }
        //每条数据都暂停处理 500ms
        Thread.sleep(500);

        return value.toString();
    }

});

result.print();

env.execute();

将以上代码打包并提交到Standalone集群执行,默认checkpoint并行度为1,此时FlinkCheckpoint会在上一次checkpoint执行完成后立即执行后续checkpoint;如果设置checkpoint多并行度,可以看到同时会有多个checkpoint同时进行。

  1. 使用默认checkpoint并行度为1

    启动Standalone集群,并向集群中提交打包好的Flink任务,提交命令中设置并行度为8,checkpoint 周期为1s,checkpoint语义为exactly_once,checkpoint并行度使用默认的1。

    [root@node4 ~]# cd /software/flink-1.17.1/bin/
    [root@node4 bin]# ./flink run -m node1:8081 \
    -Dparallelism.default=8 \
    -Dexecution.checkpointing.interval=1000 \
    -Dexecution.checkpointing.mode=exactly_once \
    -c com.mashibing.flinkjava.code.chapter12.CheckpointParalleCodeTest /root/flink-jar-test/FlinkJavaCode-1.0-SNAPSHOT-jar-with-dependencies.jar
    

    可以看到Flink WebUI中checkpoint执行情况如下,可以看到虽然设置checkpoint的周期为1s,但是由于每个checkpoint中状态保存超过1s,所以每个checkpoint触发时间都是接着上个checkpoint执行完成后立即开始。

  2. 设置checkpoint并行度为100

    向Standalone集群中提交打包好的Flink任务,提交命令中设置并行度为8,checkpoint 周期为1s,checkpoint语义为exactly_once,checkpoint并行度为100。

    [root@node4 ~]# cd /software/flink-1.17.1/bin/
    [root@node4 bin]# ./flink run -m node1:8081 \
    -Dparallelism.default=8 \
    -Dexecution.checkpointing.interval=1000 \
    -Dexecution.checkpointing.mode=exactly_once \
    -Dexecution.checkpointing.max-concurrent-checkpoints= 100 \
    -c com.mashibing.flinkjava.code.chapter12.CheckpointParalleCodeTest /root/flink-jar-test/FlinkJavaCode-1.0-SNAPSHOT-jar-with-dependencies.jar
    

    可以看到Flink WebUI中checkopint执行情况如下,checkpiont执行间隔为1s。

    通过以上对比可以看到设置并行执行checkpoint可以加快checkpoint执行效率。

设置checkpiont失败次数

checkpoint在执行过程中如果出现失败设置可以容忍的检查的失败数,超过这个数量则系统自动关闭和停止任务,没有默认值。

//设置可容忍checkpoint失败次数,没有默认值值,设置为0,表示不容忍任何checkpoint失败
env.getCheckpointConfig().setTolerableCheckpointFailureNumber(0);

设置checkpoint清理策略

当开启了checkpoint外部持久化存储时,可以通过如下两种方式决定在取消Flink作业时是否清空外部存储系统中的状态数据,如果不设置改参数,Flink取消任务时默认不清空checkpoint状态数据。

  • RETAIN_ON_CANCELLATION:在Flink作业取消时保留检查点。在这种情况下,必须手动清除检查点状态。
  • DELETE_ON_CANCELLATION:在Flink作业取消时删除检查点。只有作业失败时才会保存检查点状态。
/**
 * 设置checkpoint的清理策略,当作业取消时,checkpoint数据的保留策略,默认值为RETAIN_ON_CANCELLATION
 * RETAIN_ON_CANCELLATION:当作业取消时,保留checkpoint数据
 * DELETE_ON_CANCELLATION:当作业取消时,删除checkpoint数据
 */
env.getCheckpointConfig().setExternalizedCheckpointCleanup(CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
env.getCheckpointConfig().setExternalizedCheckpointCleanup(CheckpointConfig.ExternalizedCheckpointCleanup.DELETE_ON_CANCELLATION);

设置增量Checkpoint

在Flink状态后端章节中,我们提到Flink提供了两类状态后端,一类是HashMapStateBackend,另一类是EmbeddedRocksDBSatateBackend。

  • HashMapStateBackend

    HashMapStateBackend将状态数据以HashMap数据结构进行存储,默认将状态存储在JobManager内存中。通过用户指定checkpint持久化目录也可以将状态数据存储在外部持久化系统中。这种状态后端每次进行checkpoint检查点时都是全量方式进行,适用于较小的状态数据集。

  • EmbeddedRocksDBStateBackend

    EmbeddedRocksDBStateBackend 是 Flink 的一种基于 RocksDB 的状态后端。RocksDB 是一个高性能、持久化的键值存储引擎,它将状态数据存储在本地磁盘上,默认在TaskManager本地数据目录中。与 HashMapStateBackend全量存储状态不同,RockDBStateBackend是目前唯一支持增量检查点的状态后端,可以保存非常大的状态。

    如果Flink任务checkpoint状态非常大,开启增量checkpoint应该是首要考虑因素,与完整的checkpoint相比,增量checkpoint可以显著减少checkpoint时间,因为增量checkpoint仅存储与先前完成的checkpoint不同的增量文件,而非全量数据备份。在生产环境中建议使用RockDBStateBackend方式存储状态。

    代码中设置RockDBStateBackend的方式并开启增量checkpoint方式如下:

    //设置状态后端为RocksDBStateBackend,并指定增量checkpoint
    env.setStateBackend(new EmbeddedRocksDBStateBackend(true));
    env.getCheckpointConfig().setCheckpointStorage("your-checkpoint-dir");
    

    在Flink代码中设置状态后端为new EmbeddedRocksDBStateBackend()默认是不开启增量checkpoint,加上true后支持增量checkpoint。也可以在Flink集群的flink-conf.yaml文件中配置state.backend.incremental为ture来开启增量checkpoint,或者在向集群中提交Flink任务时通过命令来设置。

    下面通过一个案例来演示增量checkpoint保存。该案例中读取自定义source中的数据并通过map处理,每条数据暂停处理0.5秒。案例代码如下:

    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
    
    DataStreamSource<StationLog> source = env.addSource(new RichParallelSourceFunction<StationLog>() {
        Boolean flag = true;
    
        /**
         * 主要方法:启动一个Source,大部分情况下都需要在run方法中实现一个循环产生数据
         * 这里计划1s 产生1条基站数据,由于是并行,当前节点有几个core就会有几条数据
         */
        @Override
        public void run(SourceContext<StationLog> ctx) throws Exception {
            Random random = new Random();
            String[] callTypes = {"fail", "success", "busy", "barring"};
            while (flag) {
                String sid = "sid_" + random.nextInt(10);
                String callOut = "1811234" + (random.nextInt(9000) + 1000);
                String callIn = "1915678" + (random.nextInt(9000) + 1000);
                String callType = callTypes[random.nextInt(4)];
                Long callTime = System.currentTimeMillis();
                Long durations = Long.valueOf(random.nextInt(50) + "");
                ctx.collect(new StationLog(sid, callOut, callIn, callType, callTime, durations));
                Thread.sleep(1000);//1s 产生一个事件
            }
    
        }
    
        //当取消对应的Flink任务时被调用
        @Override
        public void cancel() {
            flag = false;
        }
    });
    
    //处理数据
    SingleOutputStreamOperator<String> result =
            source.keyBy(stationLog -> stationLog.sid)
                    .map(new RichMapFunction<StationLog, String>() {
    
                        private ListState<String> listState;
    
                        @Override
                        public void open(Configuration parameters) throws Exception {
                            ListStateDescriptor<String> stateDescriptor = new ListStateDescriptor<>("liststate", String.class);
                            listState = getRuntimeContext().getListState(stateDescriptor);
                        }
    
                        @Override
                        public String map(StationLog value) throws Exception {
                            //100倍状态存储
                            for(int i = 0 ;i<100 ;i++){
                                listState.add(value.toString());
                            }
                            return value.toString();
                        }
    
                    });
    
    result.print();
    
    env.execute();
    

    如果不使用RocksDB 增量状态存储,可以看到checkpoint每次保存的都是全量状态。如果设置使用了RocksDB增量状态存储,可以看到checkpoint每次保存的都是增量状态。

    1. 使用全量checkpoint

      启动Standalone集群,并向集群中提交打包好的Flink任务,提交命令中设置并行度为8,checkpoint 周期为1s,状态后端使用rocksdb并设置状态保存在HDFS路径中,checkpoint语义为exactly_once。

      [root@node4 ~]# cd /software/flink-1.17.1/bin/
      [root@node4 bin]# ./flink run -m node1:8081 \
      -Dparallelism.default=8 \
      -Dexecution.checkpointing.interval=1000 \
      -Dstate.backend.type=rocksdb \
      -Dstate.checkpoints.dir=hdfs://mycluster/rockDBState-dir \
      -Dexecution.checkpointing.mode=exactly_once \
      -c com.wubaibao.flinkjava.code.chapter12.CKWithRockDBStateBackend /root/flink-jar-test/FlinkJavaCode-1.0-SNAPSHOT-jar-with-dependencies.jar
      

      可以看到Flink WebUI中checkpoint执行情况如下,checkpoint进行全量保存。

    2. 设置RocksDB状态后端并开启增量checkpoint

      启动Standalone集群,并向集群中提交打包好的Flink任务,提交命令中设置并行度为8,checkpoint 周期为1s,状态后端使用rocksdb并设置状态保存在HDFS路径中,checkpoint语义为exactly_once,checkpoint进行增量保存。

      [root@node4 ~]# cd /software/flink-1.17.1/bin/
      [root@node4 bin]# ./flink run -m node1:8081 \
      -Dparallelism.default=8 \
      -Dexecution.checkpointing.interval=1000 \
      -Dstate.backend.type=rocksdb \
      -Dstate.checkpoints.dir=hdfs://mycluster/rockDBState-dir \
      -Dexecution.checkpointing.mode=exactly_once \
      -Dstate.backend.incremental=true \
      -c com.mashibing.flinkjava.code.chapter12.RocksDBCKTest /root/flink-jar-test/FlinkJavaCode-1.0-SNAPSHOT-jar-with-dependencies.jar
      

      可以看到Flink WebUI中checkpoint执行情况如下,checkpoint进行增量保存。

      通过以上对比可以看出,使用RocksDB状态后端并设置checkpoint增量保存后,每次可以大大节省状态保存的大小。一旦启用了增量快照,网页上展示的 Full Checkpointed Data Size 只代表增量上传的数据量,而不是一次快照的完整数据量。

开启不对齐checkpoint

当Flink作业正运行在严重的背压下时,由于缓存中需要存储大量的数据可能会导致checkpoint周期非常长,这种情况下我们可以设置非对齐checkpoint(Unaligned checkpoint)。Flink从1.11版本开始支持 Unaligned checkpoints,非对齐checkpoint中 将In-flight 数据(例如,存储在缓冲区中的数据)作为 Checkpoint State的一部分,允许 Checkpoint Barrier 跨越这些缓冲区, Checkpoint 时长变得与当前吞吐量无关,从而减少checkpoint的时长。

可以通过两种方式来设置非对齐Checkpoint:代码中设置及flink-conf.yaml配置文件中设置。

  1. 代码中设置

    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();// 启用非对齐 Checkpoint
    env.getCheckpointConfig().enableUnalignedCheckpoints();
    
  2. flink-conf.yaml配置文件中设置

    execution.checkpointing.unaligned: true
    

注意:非对齐checkpoints 会增加状态存储的IO,因此当状态存储的IO是整个checkpoint过程中真正的瓶颈时,不能使用非对齐checkpoint。

开启Changelog

Flink 中使用RocksDB作为状态后端时,支持增量快照,可以减少checkpoint持续时间。在增量快照中,RocksDB出于空间放大和读性能的考虑会定期对状态做Compaction,Compaction会产生新的且相对较大的状态文件,这会导致除了新的变更之外还要重新上传旧状态。这个过程可能会带来如下2个问题:

  • 由于要重新上传旧状态,会增加checkpoint快照上传所需时间。在 Flink 的大型作业中,每次checkpoint中至少有一个task上传大量数据的可能性非常高,所以在大型Flink任务重这种延迟问题较为突出。
  • 目前这种增量checkpoint机制中,一个task只有在收到至少一个checkpoint barrier之后,才会做状态快照,无形中也增加了checkpoint时间。

为了解决目前增量checkpoint出现的以上问题,Flink1.15版本引入了Changelog State Backend ,其核心思想是引入State Changelog (状态变化日志)更加细粒度的持久化状态。Changelog State Backend机制中TaskManager在做State操作时会进行两份状态写入:

  1. State数据会写入到本地State Table中,State Table中记录了算子计算的状态结果,会周期性将数据存储到指定的外部存储系统中。这个过程是独立于Checkpoint过程。
  2. State数据以Append Only的形式写到本地Changelog中,Changelog存储在TaskManager内存中,状态数据的增量更改(插入/更新/删除)会被写入到Changelog中,每次checkpoint时,只需要将这部分增量数据同步到持久化存储中就完成了checkpoint过程。

当State Table同步到Checkpoint Storage后,意味着State ChangeLog中相关部分也被持久化到Checkpoint Storage中,在Checkpoint Storage中的ChangeLog就可以被截断到相应的点。以上两份状态写入及ChangeLog截断示意图如下:

Changelog State Backend 由于Changelog信息相对固定,在做Checkpoint时只需要持久化Changelog部分,持久化的数据变少,所以可以减少checkpointing时间,可以减少exactly-once 模式下端到端的延迟。但是也会带来一些额外的开销:

  1. 在状态存储系统中创建更多额外的文件,更多的IO用来上传状态变更。
  2. changelog带来TaskManager额外的内存使用。
  3. 容错恢复需要额外的重放 Changelog 带来的潜在的恢复时间的增加。

Changelog State Backend机制可以在代码中进行设置开启,方式如下:

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableChangelogStateBackend(true);

12.3.2.11 开启checkpoint压缩

Flink 为所有 checkpoints 和 savepoints 提供可选的压缩(默认:关闭)。 目前压缩仅支持使用 snappy 压缩算法,Flink未来版本会支持自定义压缩算法。 设置压缩后可以大大减少Flink存储状态的大小,压缩可以通过如下参数进行设置:

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
//开启checkpoint/savepoint 压缩
env.getConfig().setUseSnapshotCompression(true);

注意:压缩选项对增量快照没有影响,因为增量快照使用的是 RocksDB 的内部格式,该格式始终使用开箱即用的 snappy 压缩。

12.3.3 RocksDB优化

RocksDB是由Facebook开发并开源的高性能键值存储库,使用C++编写,适合低延迟快速、大规模数据存储。在生产环境中,许多大型的Flink流应用程序的状态存储主要是RocksDB进行存储,这种状态后端可以可靠地存储大状态。

RocksDB类似HBase,采用LSM树(Log-Structured Merge Tree)的方式组织数据,在RocksDB3.0版本后引入列族概念(Column Families),RocksDB中的键值都与列族相关联,默认列族为“default”。

RocksDB架构中主要有三种结构:memtable、sstfile和logfile。memtable存在于内存中,当向RocksDB写入数据时,数据首先被写入内存中的memtable中,具有较高的写入性能,当memtable超过一定阈值时,会将数据写入磁盘中的sstfile中,整个数据写入操作还可被追加到logfile中形成WAL日志,以便保证数据的可靠性和一致性。RocksDB底层会定期将多个sstfile文件进行合并,以减少磁盘文件数量并优化读取性能。当从RocksDB中读取数据时,首先会在内存中的memtable中查找对应数据,找到直接返回,如果在memtable中没有找到则会在BlockCache块缓存中查找(BlockCache是用于读取数据的内存缓存,部分内存用于缓存SSTable文件的块,以提高读取性能),如果在BlockCache中未找到,则会在磁盘上sstfile文件中查找,最终返回数据。

12.3.3.1 RocksDB内存调优

在Flink中使用RocksDB状态后端进行状态存储时,RocksDB使用的内存由Flink管理(由state.backend.rocksdb.memory.managed参数决定,默认true,即Flink管理RocksDB使用内存),默认RocksDB可用的内存为TaskManager的托管内存(Managed Memory)大小,由参数taskmanager.memory.managed.fraction控制大小,默认为总Flink 内存的0.4倍。编写Flink程序时,大多数程序不需要调整RocksDB底层配置,只需要简单的增加Flink的托管内存大小即可改善内存相关性能问题。

Flink管理RocksDB内存使用中,对同一个task slot上所有的RocksDB实例共享内存,共享内存主要消耗有三个来源(blockcache、索引和bloom过滤器、memtable),一些情况下,如果RocksDB由于缺少写缓存内存而频繁刷新或者读缓存未命中而性能不佳时,我们还可以手动管理RocksDB内存分配,Flink提供了以下两个参数来控制memtable和索引过滤器、blockcache之间的内存分配。

  • state.backend.rocksdb.memory.write-buffer-ratio,默认值 0.5,即 50% 的给定内存会分配给写缓冲区使用。
  • state.backend.rocksdb.memory.high-prio-pool-ratio,默认值 0.1,即 10% 的 block cache 内存会优先分配给索引及布隆过滤器。 我们强烈建议不要将此值设置为零,以防止索引和过滤器被频繁踢出缓存而导致性能问题。

可以在提交Flink任务时通过命令进行以上参数设置。要想手动控制RocksDB内存需要首先将state.backend.rocksdb.memory.managed设置为false,此外,用户需要保证JVM有足够内存可供RocksDB使用。

下面通过命令方式演示如何设置以上参数,案例代码同上。

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

DataStreamSource<StationLog> source = env.addSource(new RichParallelSourceFunction<StationLog>() {
    Boolean flag = true;

    /**
     * 主要方法:启动一个Source,大部分情况下都需要在run方法中实现一个循环产生数据
     * 这里计划1s 产生1条基站数据,由于是并行,当前节点有几个core就会有几条数据
     */
    @Override
    public void run(SourceContext<StationLog> ctx) throws Exception {
        Random random = new Random();
        String[] callTypes = {"fail", "success", "busy", "barring"};
        while (flag) {
            String sid = "sid_" + random.nextInt(10);
            String callOut = "1811234" + (random.nextInt(9000) + 1000);
            String callIn = "1915678" + (random.nextInt(9000) + 1000);
            String callType = callTypes[random.nextInt(4)];
            Long callTime = System.currentTimeMillis();
            Long durations = Long.valueOf(random.nextInt(50) + "");
            ctx.collect(new StationLog(sid, callOut, callIn, callType, callTime, durations));
            Thread.sleep(1000);//1s 产生一个事件
        }

    }

    //当取消对应的Flink任务时被调用
    @Override
    public void cancel() {
        flag = false;
    }
});

//处理数据
SingleOutputStreamOperator<String> result =
        source.keyBy(stationLog -> stationLog.sid)
                .map(new RichMapFunction<StationLog, String>() {

                    private ListState<String> listState;

                    @Override
                    public void open(Configuration parameters) throws Exception {
                        ListStateDescriptor<String> stateDescriptor = new ListStateDescriptor<>("liststate", String.class);
                        listState = getRuntimeContext().getListState(stateDescriptor);
                    }

                    @Override
                    public String map(StationLog value) throws Exception {
                        //100倍状态存储
                        for(int i = 0 ;i<100 ;i++){
                            listState.add(value.toString());
                        }
                        return value.toString();
                    }

                });

result.print();

env.execute();

将以上代码打包,并提交到Yarn集群运行任务,提交命令中设置并行度为8,checkpoint 周期为1s,状态后端使用rocksdb并设置状态保存在HDFS路径中,checkpoint语义为exactly_once,checkpoint进行增量保存,手动管理RocksDB内存,并设置写缓存区使用内存比例为0.6。

[root@node4 ~]# cd /software/flink-1.17.1/bin/
[root@node4 bin]# ./flink run -m node1:8081 \
-Dparallelism.default=8 \
-Dexecution.checkpointing.interval=1000 \
-Dstate.backend.type=rocksdb \
-Dstate.checkpoints.dir=hdfs://mycluster/rockDBState-dir \
-Dexecution.checkpointing.mode=exactly_once \
-Dstate.backend.incremental=true \
-Dstate.backend.rocksdb.memory.managed=false \
-Dstate.backend.rocksdb.memory.write-buffer-ratio=0.6 \
-c com.mashibing.flinkjava.code.chapter12.RocksDBCKTest /root/flink-jar-test/FlinkJavaCode-1.0-SNAPSHOT-jar-with-dependencies.jar

以上命令提交后,在Flink WebUI中默认看不到RocksDB相关参数监控指标,关于Flink如何监控查看RocksDB相关指标在后续小节会进行介绍,这里只需要记住这种参数调节方式即可。

关于RocksDB内存调优相关问题,有如下2个建议:

  1. 尝试提高性能的第一步应该是增加托管内存的大小而不是通过调整 RocksDB 底层参数引入复杂性。在集群内存资源充足情况下,除非程序本身逻辑需要大量的JVM堆内存(对象多、计算逻辑复杂),否则大部分总内存通常都可以用于 RocksDB 。可以通过调节taskmanager.memory.managed.fraction(默认0.4)参数来增加托管内存大小。
  2. 如果Flink程序中状态非常多时,并且可以看到频繁的memtable刷新(参考后续小节监控RocksDB指标监控),意味着出现了写端瓶颈,此刻如果不能提供更多内存,可以手动管理RocksDB内存分配,设置 state.backend.rocksdb.memory.managed: false,并增加写缓冲区的内存比例(state.backend.rocksdb.memory.write-buffer-ratio,默认0.5)。

12.3.3.2 RocksDB优化参数

绝大多数情况下我们在进行RocksDB内存调节时,只需要增大TaskManager托管内存即可,如果不能满足生产,需要手动调节RocksDB相关参数可以参照如下参数进行调节。

  • state.backend.rocksdb.memory.managed,默认true

在Flink中设置使用RocksDB状态后端进行状态存储时,RocksDB使用的内存是否由Flink进行管理,默认true,表示Flink管理RocksDB使用内存,分配大小为TaskManager的托管内存0.4倍。

  • state.backend.rocksdb.memory.write-buffer-ratio,默认值 0.5

该参数表示将分配给RocksDB内存总量的 50% 分配给写缓冲区使用,限制了写缓存区可能占用的最大内存使用量。该参数需要设置state.backend.rocksdb.memory.managed为false。

  • state.backend.rocksdb.memory.high-prio-pool-ratio,默认值 0.1

该参数表示将 10% 的 block cache 内存会优先分配给索引及布隆过滤器。 我们强烈建议不要将此值设置为零,以防止索引和过滤器被频繁踢出缓存而导致性能问题。该参数需要设置state.backend.rocksdb.memory.managed为false。

  • state.backend.rocksdb.memory.partitioned-index-filters,默认false

RocksDB BlockCache中存的数据主要包含三部分:Data Block(真实数据)、Index Block(每条数据索引)、Filter Block(对文件的Bloom Filter)。这三部分占用内存最多的是Index Block 和 Filter Block,如果存入RocksDB中的文件数特别多,Index Block和Filter Block数据会不断的从内存和磁盘中进行替出换入,性能非常差。

这时我们可以设置state.backend.rocksdb.memory.partitioned-index-filters参数为true,表示针对Index Block 和 Filter Block设置多级索引,将这些顶级索引数据存入内存中,然后根据顶级索引按需查询Index Block 和 Filter Block数据加载到BlockCache中,加快速度。内存较小场景中,可以开启该参数,性能提升10倍左右。

以上参数是粗粒度设置RocksDB内存分配参数,这些参数使用时只需要在提交Flink任务时进行设置即可,命令如下所示:

./flink run-application -t yarn-application \
-Dparallelism.default=8 \
-Dexecution.checkpointing.interval=1000 \
-Dstate.backend.type=rocksdb \
-Dstate.checkpoints.dir=hdfs://mycluster/rockdb-state-dir \
-Dexecution.checkpointing.mode=exactly_once \
-Dstate.backend.incremental=true \
-Dstate.backend.rocksdb.memory.managed=false \
-Dstate.backend.rocksdb.memory.write-buffer-ratio=0.6 \
-Dstate.backend.rocksdb.memory.high-prio-pool-ratio=0.2 \
-Dstate.backend.rocksdb.memory.partitioned-index-filters=true \
-c com.wubaibao.flinkjava.code.chapter12.RocksDBCKTest /root/flink-jar-test/FlinkJavaCode-1.0-SNAPSHOT-jar-with-dependencies.jar

如果以上参数不能满足需要,我们还可以对RocksDB进行底层优化参数设置来满足生产需要,以下RocksDB底层配置参数仅限于高级调优场景。高级优化中RocksDB底层调优常见的参数如下:

  • state.backend.rocksdb.writebuffer.count,默认2

Flink向RocksDB中保存数据时,Flink中每个算子的每个state都对应一个列族(ColumnFamily),每个列族都有自己的MemTable。该参数表示允许在内存中保留的Memtable最大个数,超过这个个数后,数据会被Flush刷写到磁盘上形成SST文件。

如果Flink程序算子对应的状态多,并且内存足够大,可以将该参数适当调大,例如调大到5左右,降低Flush刷下磁盘次数。

  • state.backend.rocksdb.writebuffer.size,默认64M

向RocksDB中写入数据,首先写入内存中,该参数表示数据写入RocksDB时在内存中空间占用大小。默认64M,如果Flink集群内存充足,可以适当调大该参数提高写入RocksDB性能。

  • state.backend.rocksdb.writebuffer.number-to-merge,默认1

该参数表示writebuffer在写出到存储前进行合并的最小阈值。默认为1,即有数据写入形成writebuffer就会写出到存储中。根据经验,将该值可以设置为3,避免频繁的数据写出到存储。不要将该值设置太大以避免频繁的Merge操作造成写停顿。

可以在提交Flink任务时设置“state.backend.rocksdb.metrics.cur-size-all-mem-tables”参数为true来启动Flink监控RocksDB中memTable占用内存大小来评估该值设置大小。

  • state.backend.rocksdb.block.cache-size,默认8M

该参数是设置RocksDB中BlockCache大小的参数,默认为8M。增加BlockCache大小可以明显增加读性能,在内存充足情况下,可以将该值设置到64~256MB,以增加数据读取性能。

可以在提交Flink任务时设置“state.backend.rocksdb.metrics.block-cache-usage”参数为true来启动Flink监控RocksDB Blockcache使用情况指标,可以实时观察Block Cache的用量并做出相应优化。

  • state.backend.rocksdb.block.blocksize,默认4kb

RocksDB底层组成SST文件的基本单位是block块,block块中包含了一系列有序的key和value集合,该参数表示的就是block块大小,默认为4kb。

增大block大小,可以提高数据写入性能,因为不需要频繁切换block写入数据,在相同数据量下block总数量会减少,RocksDB中对应的索引所占用的内存会减少。但是如果增大block大小后,对应的BlockCache大小没有变,那么就意为着BlockCache中可存放的block数变少,如果BlockCache中还存索引和布隆过滤器,那么可存储的block块数目会更少,这时从RocksDB中读取数据时可能会需要更多的磁盘IO操作,以查找到想要的数据。

所以,减少该值可以提高读取数据性能,降低写入性能;增大该值,可以提高写入性能但是会降低读取性能。根据经验来看,如果Flink集群内存充足,可以提高block大小到128K,同时将BlockCache一并增大,这样就会有较好的读写性能,如果内存很吃紧,就不要调节该值和BlockCache大小。

  • state.backend.rocksdb.thread.num,默认2

RocksDB后台会对写入的数据定期进行Compaction和Flush操作,该参数可以设置后台Compaction和Flush操作的线程数,默认值为2,如果使用的是机械硬盘可以相应提高该值,建议设置为4,不要设置太大,否则当有很多写入请求时,后台线程都在做Compaction操作会导致写停顿问题。

12.3.3.3 RocksDB参数使用

在Flink中设置以上RocksDB高级参数可以在代码中设置,也可以在提交任务时进行设置,在代码中设置时需要实现ConfigurableRocksDBOptionsFactory接口进行参数设置,相当于是硬编码且比较麻烦,建议在提交任务时通过指定对应参数进行设置,这种方式比较灵活方便。

  • 代码中设置RocksDB高级参数
... ...
//设置状态后端为RocksDBStateBackend,并指定增量checkpoint
EmbeddedRocksDBStateBackend rocksDBStateBackend = new EmbeddedRocksDBStateBackend(true);
rocksDBStateBackend.setRocksDBOptions(new MyOptionsFactory());
env.setStateBackend(rocksDBStateBackend);
env.getCheckpointConfig().setCheckpointStorage("hdfs://mycluster/rockDBState-dir");
... ...
//MyOptionsFactory类实现
class MyOptionsFactory implements ConfigurableRocksDBOptionsFactory {

    @Override
    public DBOptions createDBOptions(DBOptions currentOptions,
                                     Collection<AutoCloseable> handlesToClose) {
        return currentOptions
                //数据自动刷盘
                .setAtomicFlush(true);
    }

    @Override
    public ColumnFamilyOptions createColumnOptions(ColumnFamilyOptions currentOptions,
                                                   Collection<AutoCloseable> handlesToClose) {
        return currentOptions
                //设置在内存中允许保留的 memtable 最大个数,默认2
                .setMaxWriteBufferNumber(2)
                //设置数据写入RocksDB时在内存中空间占用大小,默认64M
                .setWriteBufferSize(64*1024*1024L)
                //设置内存中writebuffer 进行合并的最小阈值,默认1
                .setMinWriteBufferNumberToMerge(1)
                .setTableFormatConfig(
                    new BlockBasedTableConfig()
                        //设置RocksDB中BlockCache大小的参数,默认8M
                        .setBlockCache(new LRUCache(8*1024*1024L))
                        //设置SST文件底层block大小,默认4kb
                        .setBlockSize(4*1024L)
                );
    }

    @Override
    public RocksDBOptionsFactory configure(ReadableConfig configuration) {
        //返回配置项
        return this;
    }
}

以上高级参数设置有一定门槛,并且在代码设置比较麻烦,Flink还提供通过setPredefinedOptions方法选择Flink中预定义选项,这些选项中自动设置好了底层的一些参数,大大降低了使用门槛。设置RocksDB使用Flink预定义选项方式如下:

//设置状态后端为RocksDBStateBackend,并指定增量checkpoint
EmbeddedRocksDBStateBackend rocksDBStateBackend = new EmbeddedRocksDBStateBackend(true);
rocksDBStateBackend.setPredefinedOptions(PredefinedOptions.DEFAULT);
env.setStateBackend(rocksDBStateBackend);
        env.getCheckpointConfig().setCheckpointStorage("hdfs://mycluster/rockDBState-dir");

PredefinedOptions默认有四种选项:

  1. DEFAULT:默认,没有额外参数优化,RocksDB只使用磁盘。
  2. SPINNING_DISK_OPTIMIZED:RocksDB使用磁盘并设置一些优化参数。
  3. SPINNING_DISK_OPTIMIZED_HIGH_MEM:RocksDB使用磁盘和内存并设置一些优化参数。
  4. FLASH_SSD_OPTIMIZED:RocksDB使用固态磁盘并设置优化参数。

以上四种选项中,每一种Flink都设置了对应的优化参数,我们可以直接选择使用,当内存充足时,建议选择SPINNING_DISK_OPTIMIZED_HIGH_MEM方式。

  • 提交任务时参数设置RocksDB高级参数

提交任务设置高级参数通过“-D”进行指定即可,代码中不设置任何RocksDB相关配置,提交命令中指定状态后端使用RocksDB及设置参数命令如下:

./flink run-application -t yarn-application \
-Dparallelism.default=8 \
-Dexecution.checkpointing.interval=1000 \
-Dstate.backend.type=rocksdb \
-Dstate.checkpoints.dir=hdfs://mycluster/rockdb-state-dir \
-Dexecution.checkpointing.mode=exactly_once \
-Dstate.backend.incremental=true \
-Dstate.backend.rocksdb.memory.managed=false \
-Dstate.backend.rocksdb.writebuffer.count=2 \
-Dstate.backend.rocksdb.writebuffer.size=64M \
-Dstate.backend.rocksdb.writebuffer.number-to-merge=1 \
-Dstate.backend.rocksdb.block.cache-size=8M \
-Dstate.backend.rocksdb.block.blocksize=4kb \
-Dstate.backend.rocksdb.thread.num=2 \
-c com.mashibing.flinkjava.code.chapter12.RocksDBCKTest /root/flink-jar-test/FlinkJavaCode-1.0-SNAPSHOT-jar-with-dependencies.jar

12.3.3.4 RocksDB指标监控

Flink的监控指标系统支持监控RocksDB的原生指标。默认情况下,RocksDB的原生指标监控是未启用的,我们可以通过设置特定参数来启用。需要注意,一旦启用了RocksDB的原生指标监控,可能会对Flink应用程序的性能产生一定的影响。

Flink中开启RocksDB的原生指标监控需要设置“state.backend.latency-track.keyed-state-enabled”参数为true,该参数表示是否开启State状态性能监控,默认false,不跟踪监控State相关信息,开启该参数后,我们可以通过Flink WebUI查看Flink 任务中的一些状态写入延迟信息。

例如,我们在执行“RocksDBCKTest”Flink任务时,默认是看不到该任务中用户定义的ListState相关延迟状态,如果想要跟踪这些状态需要设置“state.backend.latency-track.keyed-state-enabled”参数为true。提交命令如下:

./flink run-application -t yarn-application \
-Dparallelism.default=8 \
-Dexecution.checkpointing.interval=1000 \
-Dstate.backend.type=rocksdb \
-Dstate.checkpoints.dir=hdfs://mycluster/rockdb-state-dir \
-Dexecution.checkpointing.mode=exactly_once \
-Dstate.backend.incremental=true \
-Dstate.backend.latency-track.keyed-state-enabled=true \
-c com.mashibing.flinkjava.code.chapter12.RocksDBCKTest /root/flink-jar-test/FlinkJavaCode-1.0-SNAPSHOT-jar-with-dependencies.jar

任务提交后,我们可以通过WebUI查看Flink任务中设置状态的一些延迟指标,Flink针对每个并行度都有生成对应的监控指标。如下图所示。

以上指标监控的是State写入的延迟中位数信息,单位是ns,例如上图中48000ns表示的写入延迟为0.48ms(1ms毫秒=1000us微秒=1000000ns纳秒)。我们也可以在提交Flink任务时,指定监控RocksDB的原生指标,常见的RocksDB监控指标如下:

  • state.backend.rocksdb.metrics.actual-delayed-write-rate,默认false

监控RocksDB写延迟率,0表示没有延迟。

  • state.backend.rocksdb.metrics.block-cache-capacity,默认false

监控RocksDB BlockCache 容量。

  • state.backend.rocksdb.metrics.block-cache-pinned-usage,默认false

监控RocksDB BlockCache 中固定内存占用大小。

  • state.backend.rocksdb.metrics.block-cache-usage,默认false

监控RocksDB BlockCache 内存占用情况。

  • state.backend.rocksdb.metrics.bytes-written,默认false

监控RocksDB中Put()、Delete()、 Merge()、Write()操作写入的未压缩字节数,不包括压缩写入字节数。

  • state.backend.rocksdb.metrics.cur-size-active-mem-table,默认false

监视活动memtable的大致大小(以字节为单位)。

  • state.backend.rocksdb.metrics.cur-size-all-mem-tables,默认false

监视活动和未刷新的不可变memtable的大致大小(以字节为单位)。

  • state.backend.rocksdb.metrics.estimate-num-keys,默认false

估计RocksDB中的键数。

  • state.backend.rocksdb.metrics.is-write-stopped,默认false

跟踪RocksDB中写操作是否停止。如果写已停止返回1,否则返回0。

  • state.backend.rocksdb.metrics.num-entries-active-mem-table,默认false

监视活动memtable中的条目总数。

  • state.backend.rocksdb.metrics.num-running-flushes,默认false

监视当前正在运行的flush次数。

例如,我们在提交Flink任务时除了指定“state.backend.latency-track.keyed-state-enabled”参数外,还可以开启监控RocksDB指标,以提交Flink “RocksDBCKTest”任务为例,提交命令如下:

./flink run-application -t yarn-application \
-Dparallelism.default=8 \
-Dexecution.checkpointing.interval=1000 \
-Dstate.backend.type=rocksdb \
-Dstate.checkpoints.dir=hdfs://mycluster/rockdb-state-dir \
-Dexecution.checkpointing.mode=exactly_once \
-Dstate.backend.incremental=true \
-Dstate.backend.latency-track.keyed-state-enabled=true \
-Dstate.backend.rocksdb.metrics.block-cache-usage=true \
-Dstate.backend.rocksdb.metrics.estimate-num-keys \
-c com.wubaibao.flinkjava.code.chapter12.RocksDBCKTest /root/flink-jar-test/FlinkJavaCode-1.0-SNAPSHOT-jar-with-dependencies.jar

任务提交后,我们可以通过WebUI查看RocksDB监控指标。如下图所示。

12.3.4 Timer状态保存

Flink中,当选择RocksDB作为State Backend时,Timer定时器会默认存储在RocksDB中,在 RocksDB 中维护计时器会有一定的成本,因此 Flink 也提供了将定时器存储在 JVM 堆上而使用 RocksDB 存储其他状态的选项。当Timer定时器数量较少时,基于堆的定时器可以有更好的性能。可以通过设置state.backend.rocksdb.timer-service.factory 参数为 heap来将Timer定时器存储在Jvm堆上提高性能。

12.3.5 设置状态生存时间TTL

此部分参考资源及代码优化小节内容。

12.3.6 Task本地恢复

在Flink Checkpoint机制中,每个task都生成其状态快照,然后将其写入分布式存储。每个task通过发送一个描述状态在分布式存储中位置的句柄来确认将状态成功写入JobManager。JobManager随后收集所有任务的句柄,并将它们捆绑到一个Checkpoint对象中。

在状态恢复的情况下,JobManager打开最新的Checkpoint对象,并将句柄发送回相应的tasks,tasks可以从分布式存储中恢复其状态。使用分布式存储来存储状态的优势是:状态具备容错及所有节点都可以从分布式存储中访问获取状态信息。

然而,使用远程分布式存储也有一个很大的缺点:所有tasks都必须通过网络从远程位置读取其状态。在许多情况下,恢复可以将失败的task重新调度到与先前运行相同的TaskManager中,但我们仍然必须从远程读取其状态,这可能导致大状态的长时间恢复。

12.3.6.1 Task本地恢复原理

为了解决Task恢复时从远程获取状态导致大状态恢复时间长的问题,Flink支持Task本地恢复。其思想如下:对于每个 checkpoint ,每个 task 不仅将 task 状态写入分布式存储中, 而且还在 task 本地存储(例如本地磁盘或内存)中保存状态快照的次要副本(secondary copy)。请注意,快照的主存储仍然必须是分布式存储,因为本地存储不能确保节点故障下的持久性,也不能为其他节点提供重新分发状态的访问,所以这个功能仍然需要保存状态快照的主副本(primary copy)。

这样,对于可以重新调度到以前的TaskManager节点进行恢复的 task ,我们可以从本地状态进行状态恢复,并避免远程读取状态的成本。考虑到许多故障不是节点故障,即使节点故障通常一次只影响一个或非常少的节点, 在恢复过程中,大多数 task 很可能会重新部署到它们以前的TaskManager节点,并从本地状态进行状态恢复,这就是 task 本地恢复有效地减少大状态恢复时间的原因。

注意:在大多数情况下,本地次要副本实现只是简单地将对分布式存储的写操作复制到本地文件,性能上可能会有一些额外的成本。

分布式快照主副本(primary copy)和本地次要副本(secondary copy)有如下几点需要注意:

  1. 对于Checkpointing,主副本必须成功,并且生成的本地次要副本的失败不会使checkpoint失败。如果主副本创建失败,即使存在本地次要副本,checkpoint也是失败的。
  2. Checkpoint 主副本和次要副本生命周期互不影响。主副本由JobManager进行管理,次要副本由TaskManager进行管理。
  3. 关于task恢复状态时,如果匹配的次要副本可用,Flink将始终首选从本地恢复task状态,如果在次要副本恢复过程中出现问题,Flink可以从主副本恢复task状态。仅当从主副本和次要副本恢复task都失败时,恢复才会失败。
  4. Task本地副本可能仅包含完整task状态的一部分(例如:写入一个本地文件时出现异常)。这种情况下,Flink会首先尝试在本地恢复本地部分,非本地状态从主副本恢复。主状态必须始终是完整的,并且是task本地状态的超集。
  5. Task本地状态可以存储在堆内存或者磁盘中。不必与Task 主状态保持一致。
  6. 如果TaskManager丢失,则所有task的本地状态都会丢失。
  7. task本地恢复仅涵盖Keyed State ,不久的将来会支持算子状态和定时器状态。
  8. unaligned checkpoint 目前不支持task本地恢复。

12.3.6.2 配置Task本地恢复

Task 本地恢复默认是禁用的,可以在提交Flink任务时通过设置state.backend.local-recovery参数为true 来启用。

--- 本文结束 The End ---