flink1.11非对齐检查点unalignedcheckpoint简介

@SmartSi @SmartSi     2022-12-02     121

关键词:

传送门:Flink 系统性学习笔记


1. 前言

在阅读本文之前,建议先阅读这两篇文章:Chandy-Lamport分布式快照算法小记分布式数据流的轻量级异步快照

2. Barrier 对齐的风险

在 Flink 的检查点机制中,Checkpoint Barrier 是划分 Checkpoint 的边界。在启用 Exactly Once 语义的条件下,当一个算子有多个输入流时,需要等待所有输入流中当前检查点 N 的 Barrier 都到达其输入缓冲区,才能安全地触发检查点,否则检查点 N 的快照数据和检查点 N + 1 的快照数据就会混在一起。图示如下。


Barrier 对齐不仅保证了状态的准确性,还巧妙地消去了原生 Chandy-Lamport 算法中记录输入流状态的步骤(之前说过,即使作业执行计划是有环图,也只需要记录回边流的状态),十分轻量级。

但是,Barrier 对齐是阻塞式的,在作业出现反压时可能会成为不定时炸弹。我们知道,检查点 Barrier 是从 Source 端产生并源源不断地向下游流动的。如果作业出现反压(哪怕整个DAG中的一条链路反压),数据流动的速度减慢,Barrier 到达下游算子的延迟就会变大,进而影响到检查点完成的延时(变大甚至超时失败)。如果反压长久不能得到解决,快照数据与实际数据之间的差距就越来越明显,一旦作业 Failover,势必丢失较多的处理进度。另一方面,作业恢复后需要重新处理的数据又会积压,加重反压,造成恶性循环。

为了规避风险,Flink 1.11版本中通过 FLIP-76 引入了非对齐检查点 Unaligned Checkpoint,下面简要介绍之。

3. 非对齐检查点

顾名思义,非对齐检查点取消了屏障对齐操作。其流程图示如下。

简单解说:

a) 当算子的所有输入流中的第一个 Barrier 到达算子的输入缓冲区时,立即将这个 Barrier 发往下游(输出缓冲区)。

b) 由于第一个 Barrier 没有被阻塞,它的步调会比较快,超过一部分缓冲区中的数据。算子会标记两部分数据:一是 Barrier 首先到达的那条流中被超过的数据,二是其他流中位于当前检查点 Barrier 之前的所有数据(当然也包括进入了输入缓冲区的数据),如下图中标黄的部分所示。


c) 将上述两部分数据连同算子的状态一起做异步快照。

由此可见,非对齐检查点的机制与原生 Chandy-Lamport 算法更为相似一些(即需要由算子来记录输入流的状态)。它与对齐检查点的区别主要有三:

  • 对齐检查点在最后一个 Barrier 到达算子时触发,非对齐检查点在第一个 Barrier 到达算子时就触发。
  • 对齐检查点在第一个 Barrier 到最后一个 Barrier 到达的区间内是阻塞的,而非对齐检查点不需要阻塞。

显然,即使再考虑反压的情况,Barrier 也不会因为输入流速度变慢而堵在各个算子的入口处,而是能比较顺畅地由 Source 端直达 Sink 端,从而缓解检查点失败超时的现象。

  • 对齐检查点能够保持快照N~N + 1之间的边界,但非对齐检查点模糊了这个边界。

既然不同检查点的数据都混在一起了,非对齐检查点还能保证 Exactly Once语义吗?答案是肯定的。当任务从非对齐检查点恢复时,除了对齐检查点也会涉及到的 Source 端重放和算子的计算状态恢复之外,未对齐的流数据也会被恢复到各个链路,三者合并起来就是能够保证 Exactly Once 的完整现场了。

非对齐检查点目前仍然作为试验性的功能存在,并且它也不是十全十美的(所谓优秀的implementation往往都要考虑trade-off),主要缺点有二:

  • 需要额外保存数据流的现场,总的状态大小可能会有比较明显的膨胀(文档中说可能会达到a couple of GB per task),磁盘压力大。当集群本身就具有I/O bound的特点时,该缺点的影响更明显。
  • 从状态恢复时也需要额外恢复数据流的现场,作业重新拉起的耗时可能会很长。特别地,如果第一次恢复失败,有可能触发death spiral(死亡螺旋)使得作业永远无法恢复。

所以,官方当前推荐仅将它应用于那些容易产生反压且I/O压力较小(比如原始状态不太大)的作业中。随着后续版本的打磨,非对齐检查点肯定会更加好用。

原文:Flink新特性之非对齐检查点(unaligned checkpoint)简介

AArch64 是不是支持非对齐访问?

】AArch64是不是支持非对齐访问?【英文标题】:DoesAArch64supportunalignedaccess?AArch64是否支持非对齐访问?【发布时间】:2016-11-2621:44:32【问题描述】:AArch64是否原生支持非对齐访问?我问是因为目前ocamlopt假设“否”。【问题讨论... 查看详情

flink1.11+版本如何生成watermark(代码片段)

传送门:Flink系统性学习笔记Flink1.11在Flink1.11版本之前,Flink提供了两种生成Watermark的策略,分别是AssignerWithPunctuatedWatermarks和AssignerWithPeriodicWatermarks,这两个接口都继承自TimestampAssigner接口。用户想使用不同的Watermark生成方式,... 查看详情

flink1.11+版本如何生成watermark(代码片段)

传送门:Flink系统性学习笔记Flink1.11在Flink1.11版本之前,Flink提供了两种生成Watermark的策略,分别是AssignerWithPunctuatedWatermarks和AssignerWithPeriodicWatermarks,这两个接口都继承自TimestampAssigner接口。用户想使用不同的Watermark生成方式,... 查看详情

flinkflink源码编译:flink1.11+版本编译及部署

1.概述转载:Flink源码编译:Flink1.11+版本编译及部署 查看详情

flink1.11之后的新版数据源(代码片段)

数据源注意:当前文档所描述的为新的数据源API,在Flink1.11中作为FLIP-27中的一部分引入。该新API仍处于BETA阶段。(从Flink1.11开始)大多数现有的source连接器尚未使用此新API实现,仍旧使用之前的API,也就是基... 查看详情

armv8-a非对齐数据访问支持(alignmentsupport)

目录1,对齐传输和非对齐传输2,AArch32 Alignmentsupport2.1Instructionalignment指令对齐2.2Unaligneddataaccess非对齐数据访问 2.3 SCTLR.A Alignmentcheckenable3,AArch64 Alignmentsupport3.1Instructionalignment指令对齐3.2Alignmentofdataaccesses对齐数据... 查看详情

flink1.11unalignedcheckpoint解析

传送门:Flink系统性学习笔记作为Flink最基础也是最关键的容错机制,Checkpoint快照机制很好地保证了Flink应用从异常状态恢复后的数据准确性。同时Checkpoint相关的metrics也是诊断Flink应用健康状态最为重要的指标,成功... 查看详情

flink1.11报错:java.lang.illegalstateexception:noexecutorfactoryfoundtoexecutetheapplication(代码片(代码片段)

一.报错信息Exceptioninthread"main"java.lang.IllegalStateException:NoExecutorFactoryfoundtoexecutetheapplication.atorg.apache.flink.core.execution.DefaultExecutorServiceLoader.getExecutorFactory(DefaultExec 查看详情

[问题踩坑]flink1.11.1sqlview中udtf调用异常column‘xxx‘notfoundinanytable(代码片段)

在Flink1.11.1版本中,执行下面的FlinkSQL,会抛出异常:"org.apache.calcite.sql.validate.SqlValidatorException:Column'message'notfoundinanytable"--创建Kafka数据源表test_tableCREATETABLEtest 查看详情

flink监控检查点checkpoint

...1.11Flink的Web页面中提供了一些页面标签,用于监控作业的检查点Checkpoint。这些监控统计信息即使在作业终止后也可以看到。Checkpoints监控页面共有四个不同的Tab页签:Overview、History、Summary和Configuration,它们分别从不同角度进行... 查看详情

flinksql1.11流批一体hive数仓(代码片段)

Flink1.11features已经冻结,流批一体在新版中是浓墨重彩的一笔,在此提前对Flink1.11中流批一体方面的改善进行深度解读,大家可期待正式版本的发布。Flink1.11中流计算结合Hive批处理数仓,给离线数仓带来Flink流处... 查看详情

非对齐访问和alignmentfault

什么是对齐异常?简单来说,当CPU访问内存地址时,如果发现访问的地址是不对齐的,硬件(部分)就会自动触发对齐异常。对齐即要求被访问的地址满足其数据类型的位宽要求,比如要访问一个4字节int型的数据,但是提供的地址... 查看详情

非对象字段错误对齐或重叠

】非对象字段错误对齐或重叠【英文标题】:Incorrectlyalignedoroverlappedbyanon-objectfielderror【发布时间】:2010-11-1411:23:31【问题描述】:我正在尝试创建以下结构:[StructLayout(LayoutKind.Explicit,Size=14)]publicstructMessage[FieldOffset(0)]publicushortX... 查看详情

[问题踩坑]flink1.11.1sqlview中udtf调用异常column‘xxx‘notfoundinanytable(代码片段)

在Flink1.11.1版本中,执行下面的FlinkSQL,会抛出异常:"org.apache.calcite.sql.validate.SqlValidatorException:Column'message'notfoundinanytable"--创建Kafka数据源表test_tableCREATETABLEtest_table(usernameSTRING,messageSTRING)WITH('connector'... 查看详情

带你深入了解对齐与非对齐访问(arm指令集)(代码片段)

首先你需要知道在什么情况下你才需要用到对齐与非对齐这个概念typedefstruct char*pCmd; //4个字节 char*pCmdPara;//4个字节 charisFree; //4个字节 AT_QUEUE_ITEM_T;AT_QUEUE_ITEM_Tq[500]=0;//12*500=6000个 查看详情

一文带你深入了解《c语言对齐与非对齐访问》(arm指令集)(代码片段)

首先你需要知道在什么情况下你才需要用到对齐与非对齐这个概念typedefstruct char*pCmd; //4个字节 char*pCmdPara;//4个字节 charisFree; //4个字节 AT_QUEUE_ITEM_T;AT_QUEUE_ITEM_Tq[500]=0;//12*500=6000个字节上面 查看详情

一文带你深入了解《c语言对齐与非对齐访问》(arm指令集)(代码片段)

首先你需要知道在什么情况下你才需要用到对齐与非对齐这个概念typedefstruct char*pCmd; //4个字节 char*pCmdPara;//4个字节 charisFree; //4个字节 AT_QUEUE_ITEM_T;AT_QUEUE_ITEM_Tq[500]=0;//12*500=6000个字节上面 查看详情

[问题踩坑]flink1.11.1onyarnapplication模式时,jobid始终为00000000000000000000000000000000(代码片段)

使用YarnApplication模式启动Flink任务的命令如下./bin/flinkrun-application-tyarn-applicationhdfs://hadoop-master:9000/shadow/FlinkSQLTest-1.0-SNAPSHOT.jar当Flink任务启动后,发现JobId异常始终为00000000000000000000000000000000,而正常的JobId应该为一个随机... 查看详情