爆款云主机2核4G限时秒杀,88元/年起!
查看详情

活动

天翼云最新优惠活动,涵盖免费试用,产品折扣等,助您降本增效!
热门活动
  • 618智算钜惠季 爆款云主机2核4G限时秒杀,88元/年起!
  • 免费体验DeepSeek,上天翼云息壤 NEW 新老用户均可免费体验2500万Tokens,限时两周
  • 云上钜惠 HOT 爆款云主机全场特惠,更有万元锦鲤券等你来领!
  • 算力套餐 HOT 让算力触手可及
  • 天翼云脑AOne NEW 连接、保护、办公,All-in-One!
  • 中小企业应用上云专场 产品组合下单即享折上9折起,助力企业快速上云
  • 息壤高校钜惠活动 NEW 天翼云息壤杯高校AI大赛,数款产品享受线上订购超值特惠
  • 天翼云电脑专场 HOT 移动办公新选择,爆款4核8G畅享1年3.5折起,快来抢购!
  • 天翼云奖励推广计划 加入成为云推官,推荐新用户注册下单得现金奖励
免费活动
  • 免费试用中心 HOT 多款云产品免费试用,快来开启云上之旅
  • 天翼云用户体验官 NEW 您的洞察,重塑科技边界

智算服务

打造统一的产品能力,实现算网调度、训练推理、技术架构、资源管理一体化智算服务
智算云(DeepSeek专区)
科研助手
  • 算力商城
  • 应用商城
  • 开发机
  • 并行计算
算力互联调度平台
  • 应用市场
  • 算力市场
  • 算力调度推荐
一站式智算服务平台
  • 模型广场
  • 体验中心
  • 服务接入
智算一体机
  • 智算一体机
大模型
  • DeepSeek-R1-昇腾版(671B)
  • DeepSeek-R1-英伟达版(671B)
  • DeepSeek-V3-昇腾版(671B)
  • DeepSeek-R1-Distill-Llama-70B
  • DeepSeek-R1-Distill-Qwen-32B
  • Qwen2-72B-Instruct
  • StableDiffusion-V2.1
  • TeleChat-12B

应用商城

天翼云精选行业优秀合作伙伴及千余款商品,提供一站式云上应用服务
进入甄选商城进入云市场创新解决方案
办公协同
  • WPS云文档
  • 安全邮箱
  • EMM手机管家
  • 智能商业平台
财务管理
  • 工资条
  • 税务风控云
企业应用
  • 翼信息化运维服务
  • 翼视频云归档解决方案
工业能源
  • 智慧工厂_生产流程管理解决方案
  • 智慧工地
建站工具
  • SSL证书
  • 新域名服务
网络工具
  • 翼云加速
灾备迁移
  • 云管家2.0
  • 翼备份
资源管理
  • 全栈混合云敏捷版(软件)
  • 全栈混合云敏捷版(一体机)
行业应用
  • 翼电子教室
  • 翼智慧显示一体化解决方案

合作伙伴

天翼云携手合作伙伴,共创云上生态,合作共赢
天翼云生态合作中心
  • 天翼云生态合作中心
天翼云渠道合作伙伴
  • 天翼云代理渠道合作伙伴
天翼云服务合作伙伴
  • 天翼云集成商交付能力认证
天翼云应用合作伙伴
  • 天翼云云市场合作伙伴
  • 天翼云甄选商城合作伙伴
天翼云技术合作伙伴
  • 天翼云OpenAPI中心
  • 天翼云EasyCoding平台
天翼云培训认证
  • 天翼云学堂
  • 天翼云市场商学院
天翼云合作计划
  • 云汇计划
天翼云东升计划
  • 适配中心
  • 东升计划
  • 适配互认证

开发者

开发者相关功能入口汇聚
技术社区
  • 专栏文章
  • 互动问答
  • 技术视频
资源与工具
  • OpenAPI中心
开放能力
  • EasyCoding敏捷开发平台
培训与认证
  • 天翼云学堂
  • 天翼云认证
魔乐社区
  • 魔乐社区

支持与服务

为您提供全方位支持与服务,全流程技术保障,助您轻松上云,安全无忧
文档与工具
  • 文档中心
  • 新手上云
  • 自助服务
  • OpenAPI中心
定价
  • 价格计算器
  • 定价策略
基础服务
  • 售前咨询
  • 在线支持
  • 在线支持
  • 工单服务
  • 建议与反馈
  • 用户体验官
  • 服务保障
  • 客户公告
  • 会员中心
增值服务
  • 红心服务
  • 首保服务
  • 客户支持计划
  • 专家技术服务
  • 备案管家

了解天翼云

天翼云秉承央企使命,致力于成为数字经济主力军,投身科技强国伟大事业,为用户提供安全、普惠云服务
品牌介绍
  • 关于天翼云
  • 智算云
  • 天翼云4.0
  • 新闻资讯
  • 天翼云APP
基础设施
  • 全球基础设施
  • 信任中心
最佳实践
  • 精选案例
  • 超级探访
  • 云杂志
  • 分析师和白皮书
  • 天翼云·创新直播间
市场活动
  • 2025智能云生态大会
  • 2024智算云生态大会
  • 2023云生态大会
  • 2022云生态大会
  • 天翼云中国行
天翼云
  • 活动
  • 智算服务
  • 产品
  • 解决方案
  • 应用商城
  • 合作伙伴
  • 开发者
  • 支持与服务
  • 了解天翼云
      • 文档
      • 控制中心
      • 备案
      • 管理中心

      【Flink状态管理(七)】Checkpoint的触发:2. CheckpointBarrier触发算子Checkpoint操作之CheckpointBarrier的对齐操作

      首页 知识中心 其他 文章详情页

      【Flink状态管理(七)】Checkpoint的触发:2. CheckpointBarrier触发算子Checkpoint操作之CheckpointBarrier的对齐操作

      2025-03-06 09:17:42 阅读次数:10

      对齐,操作

      我们已经知道,在CheckpointCoordinator中会触发数据源算子的Checkpoint操作,同时向下游节点发送CheckpointBarrier事件。当下游Task实例接收到上游节点发送的CheckpointBarrier事件消息,且接收到所有InputChannel中的CheckpointBarrier事件消息时,当前Task实例才会触发本节点的Checkpoint操作。

      这样设定的目的是让下游节点将所有InputChannel中属于当前Checkpoint的数据全部接入本节点,然后再对数据元素进行处理,以保证数据的一致性,一旦出现异常也能从上一次Checkpoint持久化结果中恢复当前Task实例的状态数据。

      一. CheckpointBarrier对齐过程

      在CheckpointedInputGate.pollNext()方法中先从上游节点获取Buffer或Event类型数据,然后分别做相应的数据处理,pollNext()方法的主要逻辑如下。

      1. 获取cp的buffer和Event:

      从Optional<BufferOrEvent> next中获取BufferOrEvent类型数据,BufferOrEvent既定义了Buffer数据也定义了Event。

      1. 锁定等待Barrier事件:

      如果当前InputChannel被barrierHandler对象锁定,则将所有的BufferOrEvent数据本地缓存,直到InputChannel的锁被打开。barrierHandler会等所有InputChannel的CheckpointBarrier事件消息全部到达节点后,才继续处理该Task实例的Buffer数据,保证数据计算结果的正确性。

      1. Barrier Reset

      在Buffer数据的处理过程中,如果Buffer缓冲区被填满,会进行清理操作和Barrier Reset操作。

      1. 消息类型处理:
        • 如果BufferOrEvent的消息类型为Buffer,则直接返回next;

        • 如果是CheckpointBarrier类型,则接入的CheckpointBarrier事件,最终根据CheckpointBarrier对齐情况选择是否触发当前节点的Checkpoint操作。

        • 如果接收到的是CancelCheckpointMarker事件,则取消本次Checkpoint操作。

        • 如果接收到的是EndOfPartitionEvent事件,表示上游Partition中的数据已经消费完毕,此时调用barrierHandler.processEndOfPartition()方法进行处理,最后清理缓冲区中的Buffer数据。

      
      // 获取BufferOrEvent数据
      BufferOrEvent bufferOrEvent = next.get();
      if (barrierHandler.isBlocked(offsetChannelIndex(bufferOrEvent.
         getChannelIndex()))) {
         // 如果当前channel被barrierHandler对象锁定,则将BufferOrEvent数据先缓存下来
         bufferStorage.add(bufferOrEvent);
         // 如果缓冲区被填满,则进行清理操作和Barrier Reset操作
         if (bufferStorage.isFull()) {
             barrierHandler.checkpointSizeLimitExceeded(bufferStorage.getMaxBufferedBytes());
             bufferStorage.rollOver();
         }
      }else if (bufferOrEvent.isBuffer()) {
         // 如果是业务数据则直接返回,留给算子处理
         return next;
      }else if (bufferOrEvent.getEvent().getClass() == CheckpointBarrier.class) {
         // 如果是CheckpointBarrier类型的事件,则对接入的Barrier进行处理
         CheckpointBarrier checkpointBarrier = (CheckpointBarrier) bufferOrEvent.getEvent();
         if (!endOfInputGate) {
            // 根据算子的对齐情况选择是否需要进行Checkpoint操作
            if (barrierHandler.processBarrier(checkpointBarrier,
                                              offsetChannelIndex(bufferOrEvent.
                                              getChannelIndex()), 
                                              bufferStorage.getPendingBytes())) {
               bufferStorage.rollOver();
            }
         }
      }else if (bufferOrEvent.getEvent().getClass() == CancelCheckpointMarker.class) {
         // 如果是CancelCheckpointMarker类型事件,则调用processCancellationBarrier()方
        //   法进行处理
         if (barrierHandler.processCancellationBarrier(
             (CancelCheckpointMarker) bufferOrEvent.getEvent())) {
            bufferStorage.rollOver();
         }
      }
      

      通过以上步骤我们可以看出,在CheckpointedInputGate中实现了CheckpointBarrier对齐的全部过程,并通过bufferStorage对接入的Buffer数据进行缓存,直到CheckpointBarrier事件全部对齐,才会对接入的数据进行处理。

       

      接下来我们重点看CheckpointBarrierHandler的具体实现

      二. CheckpointBarrierHandler的具体实现

      1. CheckpointBarrierHandler分类

      CheckpointBarrier事件的对齐过程主要借助CheckpointBarrierHandler实现。
      【Flink状态管理(七)】Checkpoint的触发:2. CheckpointBarrier触发算子Checkpoint操作之CheckpointBarrier的对齐操作

      如图,CheckpointBarrierHandler主要有CheckpointBarrierAligner和CheckpointBarrierTracker两种子类实现。

      1. CheckpointBarrierAligner用于实现Exactly-Once数据的一致性保障,对所有InputChannel中的CheckpointBarrier进行严格的对齐控制,并决定Task实例中InputChannel的block和unblock时间点。
      2. CheckpointBarrierTracker则实现了At-Least-Once语义处理保障,并没有对CheckpointBarrier进行非常严格的控制。

       

      2. CheckpointBarrierAligner实现Barrier对齐操作

      在CheckpointInputGate中通过CheckpointBarrierHandler.processBarrier()方法处理接收到的CheckpointBarrier事件。

      如代码所示,实现类CheckpointBarrierAligner.processBarrier()方法主要逻辑如下。

      • 是否进行Barrier对齐操作
      1. 从receivedBarrier中获取barrierId,判断totalNumberOfInputChannels是否为1,如果InputChannel数量为1,则触发Checkpoint操作,不需要进行CheckpointBarrier对齐操作。
      2. 如果InputChannel数量不为1,则判断numBarriersReceived是否大于0,即是否已经开始接收CheckpointBarrier事件,并进行Barrier对齐操作。
      • barrier具体操作:
      1. 如果barrierId == currentCheckpointId条件为True,则调用onBarrier()方法进行处理。
      2. 如果barrierId > currentCheckpointId,表明已经有新的Barrier事件发出,超过了当前的CheckpointId,这种情况就会忽略当前的Checkpoint,并调用beginNewAlignment()方法开启新的Checkpoint。
      3. 如果以上条件都不满足,表明当前的Checkpoint操作已经被取消或Barrier信息属于先前的Checkpoint操作,此时直接返回false。
      1. 满足numBarriersReceived + numClosedChannels ==
        totalNumberOfInputChannels条件后,触发该节点的Checkpoint操作。实际上会调用notifyCheckpoint()方法触发该Task实例的Checkpoint操作。
      public boolean processBarrier(CheckpointBarrier receivedBarrier, 
                                    int channelIndex, 
                                    long bufferedBytes) throws Exception {
         // 首先获取barrierId
         final long barrierId = receivedBarrier.getId();
         // 如果InputChannels为1,直接触发Checkpoint操作,不需要对齐处理
         if (totalNumberOfInputChannels == 1) {
            if (barrierId > currentCheckpointId) {
               // 提交新的Checkpoint操作
               currentCheckpointId = barrierId;
               notifyCheckpoint(receivedBarrier, bufferedBytes, 
                  latestAlignmentDurationNanos);
            }
            return false;
         }
         boolean checkpointAborted = false;
         if (numBarriersReceived > 0) {
            // 继续进行对齐操作
            if (barrierId == currentCheckpointId) {
               onBarrier(channelIndex);
            }else if (barrierId > currentCheckpointId) {
               LOG.warn("{}: Received checkpoint barrier for checkpoint {} before 
                  completing current checkpoint {}. " + "Skipping current checkpoint.",
                  taskName,
                  barrierId,
                  currentCheckpointId);
               // 通知Task当前Checkpoint没有完成
               notifyAbort(currentCheckpointId,
                  new CheckpointException(
                     "Barrier id: " + barrierId,
                     CheckpointFailureReason.CHECKPOINT_DECLINED_SUBSUMED));
               // 终止当前的Checkpoint操作
               releaseBlocksAndResetBarriers();
               checkpointAborted = true;
               // 开启新的Checkpoint操作
               beginNewAlignment(barrierId, channelIndex);
            }else {
               return false;
            }
         }else if (barrierId > currentCheckpointId) {
            // 创建新的Checkpoint
            beginNewAlignment(barrierId, channelIndex);
         }else {
            return false;
         }
         // 当Barrier接收的数量加上Channel关闭的数量等于整个InputChannels的数量时触发
            Checkpoint操作
         if (numBarriersReceived + numClosedChannels == totalNumberOfInputChannels) {
            if (LOG.isDebugEnabled()) {
               LOG.debug("{}: Received all barriers, triggering checkpoint {} at {}.",
                  taskName,
                  receivedBarrier.getId(),
                  receivedBarrier.getTimestamp());
            }
            // 释放Block并重置Barrier
            releaseBlocksAndResetBarriers();
           // 开始触发Checkpoint操作
            notifyCheckpoint(receivedBarrier, bufferedBytes, 
               latestAlignmentDurationNanos);
            return true;
         }
         return checkpointAborted;
      }
      

      经过以上步骤,基本上完成了CheckpointBarrier的对齐操作,当CheckpointBarrier完成对齐操作后,接下来就是通过notifyCheckpoint()方法触发StreamTask节点的Checkpoint操作。

      版权声明:本文内容来自第三方投稿或授权转载,原文地址:https://blog.csdn.net/hiliang521/article/details/136190480,作者:roman_日积跬步-终至千里,版权归原作者所有。本网站转在其作品的目的在于传递更多信息,不拥有版权,亦不承担相应法律责任。如因作品内容、版权等问题需要同本网站联系,请发邮件至ctyunbbs@chinatelecom.cn沟通。

      上一篇:【小白到大牛之路】交换机后台管理之重复输入用户名和密码

      下一篇:C 语分支初启航,循环开篇韵悠长--关系,条件操作符

      相关文章

      2025-05-14 10:33:31

      计算机初级选手的成长历程——操作符详解(2)

      计算机初级选手的成长历程——操作符详解(2)

      2025-05-14 10:33:31
      对象 , 操作 , 操作符 , 表达式 , 运算 , 逗号 , 逻辑
      2025-05-14 10:02:48

      MongoDB常用管理命令(1)

      MongoDB常用管理命令(1)

      2025-05-14 10:02:48
      会话 , 命令 , 操作 , 节点
      2025-05-13 09:49:12

      JDBC事务管理、四大特征(ACID)、事务提交与回滚、MySQL事务管理

      JDBC(Java Database Connectivity)事务是指一系列作为单个逻辑工作单元执行的数据库操作,这些操作要么全部成功——>提交,要么全部失败——>回滚,从而确保数据的一致性和完整性。

      2025-05-13 09:49:12
      MySQL , 事务 , 执行 , 提交 , 操作 , 数据库
      2025-05-07 09:10:01

      C语言:自定义类型——结构体

      数组是一组相同类型元素的集合,而结构体同样也是一些值的集合,不同的是,在结构体中,这些值被称为成员变量,而结构体的每个成员变量可以是不同类型的变量:如: 标量、数组、指针,甚⾄是其他结构体。

      2025-05-07 09:10:01
      位段 , 内存 , 字节 , 对齐 , 成员 , 类型 , 结构
      2025-05-07 09:10:01

      C语言:自定义类型——联合和枚举

      像结构体⼀样,联合体也是由⼀个或者多个成员构成,这些成员可以是不同的类型。

      2025-05-07 09:10:01
      define , 大小 , 对齐 , 成员 , 枚举 , 类型 , 联合体
      2025-05-06 09:19:39

      Linux下学【MySQL】表中修改和删除的进阶操作(配实操图和SQL语句通俗易懂)

      Linux下学【MySQL】表中修改和删除的进阶操作(配实操图和SQL语句通俗易懂)

      2025-05-06 09:19:39
      MySQL , update , 删除 , 成绩 , 操作
      2025-05-06 09:19:39

      【C/C++】如何求出类对象的大小----类结构中的内存对齐

      【C/C++】如何求出类对象的大小----类结构中的内存对齐

      2025-05-06 09:19:39
      byte , int , 内存 , 对齐 , 成员
      2025-04-23 08:18:38

      基础—SQL—图形化界面工具的DataGrip使用(2)

      基础—SQL—图形化界面工具的DataGrip使用(2)

      2025-04-23 08:18:38
      创建 , 操作 , 数据库 , 界面 , 语句
      2025-04-23 08:18:38

      基础—SQL—DQL(数据查询语言)聚合函数

      聚合函数指的是讲一列数据作为一个整体,进行纵向的计算。

      2025-04-23 08:18:38
      函数 , 员工 , 操作 , 查询 , 统计 , 聚合
      2025-04-22 09:44:09

      【C语言:自定义类型(结构体、位段、共用体、枚举)】

      【C语言:自定义类型(结构体、位段、共用体、枚举)】

      2025-04-22 09:44:09
      位段 , 对齐 , 成员 , 枚举 , 结构
      查看更多
      推荐标签

      作者介绍

      天翼云小翼
      天翼云用户

      文章

      33561

      阅读量

      5254045

      查看更多

      最新文章

      MongoDB常用管理命令(1)

      2025-05-14 10:02:48

      使数组中位数等于 K 的最少操作数

      2025-04-18 07:11:32

      Vim dom 比Real dom哪个渲染更快?

      2025-03-28 06:55:00

      【设计模式】命令模式

      2025-03-14 09:05:42

      【Flink状态管理五】Checkpoint的设计与实现

      2025-03-06 09:17:42

      Golang中import 导入包的几种方式:点,别名与下划线

      2025-03-04 09:16:53

      查看更多

      热门文章

      命令行 cmd 操作方式

      2023-03-15 09:21:53

      给你一个长度为 n、下标从 0 开始的整数数组 nums,nums[i] 表示收集位于下标 i 处的巧克力成本。每个巧克力都对应一个不同的类型,最初,位于下标 i 的巧克力就对应第 i 个类型。

      2024-11-07 07:57:04

      OpenCV从入门到精通——形态学操作

      2024-11-06 07:16:52

      Tensorflow入门(1.0)

      2024-11-21 09:55:25

      【设计模式】命令模式

      2025-03-14 09:05:42

      使用 OpenCV 进行图像的形态学操作

      2024-12-18 08:27:46

      查看更多

      热门标签

      linux java python javascript 数组 前端 docker Linux vue 函数 shell git 节点 容器 示例
      查看更多

      相关产品

      弹性云主机

      随时自助获取、弹性伸缩的云服务器资源

      天翼云电脑(公众版)

      便捷、安全、高效的云电脑服务

      对象存储

      高品质、低成本的云上存储服务

      云硬盘

      为云上计算资源提供持久性块存储

      查看更多

      随机文章

      Golang中import 导入包的几种方式:点,别名与下划线

      【Flink状态管理五】Checkpoint的设计与实现

      【设计模式】命令模式

      命令行 cmd 操作方式

      Vim dom 比Real dom哪个渲染更快?

      MYSQL中的增删改查操作(如果想知道MYSQL中有关增删改查操作的知识,那么只看这一篇就足够了!)

      • 7*24小时售后
      • 无忧退款
      • 免费备案
      • 专家服务
      售前咨询热线
      400-810-9889转1
      关注天翼云
      • 旗舰店
      • 天翼云APP
      • 天翼云微信公众号
      服务与支持
      • 备案中心
      • 售前咨询
      • 智能客服
      • 自助服务
      • 工单管理
      • 客户公告
      • 涉诈举报
      账户管理
      • 管理中心
      • 订单管理
      • 余额管理
      • 发票管理
      • 充值汇款
      • 续费管理
      快速入口
      • 天翼云旗舰店
      • 文档中心
      • 最新活动
      • 免费试用
      • 信任中心
      • 天翼云学堂
      云网生态
      • 甄选商城
      • 渠道合作
      • 云市场合作
      了解天翼云
      • 关于天翼云
      • 天翼云APP
      • 服务案例
      • 新闻资讯
      • 联系我们
      热门产品
      • 云电脑
      • 弹性云主机
      • 云电脑政企版
      • 天翼云手机
      • 云数据库
      • 对象存储
      • 云硬盘
      • Web应用防火墙
      • 服务器安全卫士
      • CDN加速
      热门推荐
      • 云服务备份
      • 边缘安全加速平台
      • 全站加速
      • 安全加速
      • 云服务器
      • 云主机
      • 智能边缘云
      • 应用编排服务
      • 微服务引擎
      • 共享流量包
      更多推荐
      • web应用防火墙
      • 密钥管理
      • 等保咨询
      • 安全专区
      • 应用运维管理
      • 云日志服务
      • 文档数据库服务
      • 云搜索服务
      • 数据湖探索
      • 数据仓库服务
      友情链接
      • 中国电信集团
      • 189邮箱
      • 天翼企业云盘
      • 天翼云盘
      ©2025 天翼云科技有限公司版权所有 增值电信业务经营许可证A2.B1.B2-20090001
      公司地址:北京市东城区青龙胡同甲1号、3号2幢2层205-32室
      • 用户协议
      • 隐私政策
      • 个人信息保护
      • 法律声明
      备案 京公网安备11010802043424号 京ICP备 2021034386号