爆款云主机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网络数据传输】OperatorChain的设计与实现

      首页 知识中心 软件开发 文章详情页

      【Flink网络数据传输】OperatorChain的设计与实现

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

      接口,算子,输出

       

      1.OperatorChain的设计与实现

      OperatorChain的大致逻辑

      在JobGraph对象的创建过程中,将链化可以连在一起的算子,常见的有StreamMap、StreamFilter等类型的算子。OperatorChain中的所有算子都会被运行在同一个Task实例中。StreamTaskNetworkOutput会将接入的数据元素写入算子链的HeadOperator中,从而开启整个OperatorChain的数据处理。

       
      OperatorChain的Output组件:将数据发送到下游

      如图所示,在OperatorChain中通过Output组件将上下游算子相连,当上游算子数据处理完毕后,会通过Output组件发送到下游的算子中继续处理。

       
      OperatorChain的collect():收集处理完的数据

      如图所示,OperatorChain内部定义了WatermarkGaugeExposingOutput接口,且该接口分别继承了Output和Collector接口。Collector接口提供了collect()方法,用于收集处理完的数据。

       
      OperatorChain的Output接口:也能输出Watermark和LatencyMarker等事件

      Output接口提供了emitWatermark()、emitLatencyMarker()等方法,用于对Collector接口进行拓展,使得Output接口实现类可以输出Watermark和LatencyMarker等事件。WatermarkGaugeExposingOutput接口则提供了获取WatermarkGauge的方法,用于监控最新的Watermark。

       
      OperatorChain内部定义了不同的WatermarkGaugeExposingOutput接口实现类。

      1. RecordWriterOutput:用于输出OperatorChain中尾端算子处理完成的数据,借助RecordWriter组件将数据元素写入网络。
      2. ChainingOutput/CopyingChainingOutput:适用于上下游算子连接在一起且上游算子属于单输出类型的情况。
      3. BroadcastingOutputCollector/CopyingBroadcastingOutputCollector:上游算子是多输出类型但上下游算子之间的Selector为空时,创建广播类型的BroadcastingOutputCollector。
      4. DirectedOutput/CopyingDirectedOutput:上游算子是多输出类型且Selector不为空时,创建DirectedOutput或CopyingDirectedOutput连接上下游算子。

      【Flink网络数据传输】OperatorChain的设计与实现

      例子:收集数据并通过Output发数据数据到下游

      例如在WordCount的程序中定义flatMap()方法时,会调用Collector.collect()方法收集数据元素,每个算子在定义的函数或使用Output接口的实现类中,完成了上游算子向下游算子发送数据元素的操作。

       

      2.OperatorChain的创建和初始化

      接下来我们看OperatorChain的初始化过程,如下代码,OperatorChain的构造器包含如下逻辑。

      1. 创建StreamOperator(即算子)实例,这里StreamOperator会封装为StreamOperatorFactory并存储在StreamGraph结构中。
      2. 获取算子之间的链接配置。chainedConfigs的配置决定了算子之间Output接口的具体实现。
      3. 遍历当前作业所有节点的输出边,并构建RecordWriterOutput组件,最终通过RecordWriterOutput组件将数据元素输出到网络中。
      4. 创建OperatorChain内部算子之间的上下游连接,完成OperatorChain内部上下游算子之间的数据传输。
      5. 单独创建headOperator。headOperator是OperatorChain的头部节点,创建完成后将headOperator暴露到StreamTask实例,供DataOutput接口实现类调用。
      6. 如果OperatorChain构建失败,则关闭实例,防止出现内存泄漏。
      public OperatorChain(
            StreamTask<OUT, OP> containingTask,
            RecordWriterDelegate<SerializationDelegate<StreamRecord<OUT>>> 
               recordWriterDelegate
            ) {
         // 获取当前StreamTask的userCodeClassloader
         final ClassLoader userCodeClassloader = containingTask.getUserCodeClassLoader();
         // 获取StreamConfig
         final StreamConfig configuration = containingTask.getConfiguration();
         // 获取StreamOperatorFactory
         StreamOperatorFactory<OUT> operatorFactory = 
         configuration.getStreamOperatorFactory(userCodeClassloader);
         // 读取chainedConfigs
         Map<Integer, StreamConfig> chainedConfigs = 
      configuration.getTransitiveChainedTaskConfigsWithSelf(userCodeClassloader);
         // 根据StreamEdge创建RecordWriterOutput组件
         List<StreamEdge> outEdgesInOrder = 
         configuration.getOutEdgesInOrder(userCodeClassloader);
         Map<StreamEdge, RecordWriterOutput<?>> streamOutputMap = 
         new HashMap<>(outEdgesInOrder.size());
         this.streamOutputs = new RecordWriterOutput<?>[outEdgesInOrder.size()];
         boolean success = false;
         try {
            for (int i = 0; i < outEdgesInOrder.size(); i++) {
               StreamEdge outEdge = outEdgesInOrder.get(i);
               // 为每个输出边创建RecordWriterOutput
               RecordWriterOutput<?> streamOutput = createStreamOutput(
                  recordWriterDelegate.getRecordWriter(i),
                  outEdge,
                  chainedConfigs.get(outEdge.getSourceId()),
                  containingTask.getEnvironment());
               this.streamOutputs[i] = streamOutput;
               streamOutputMap.put(outEdge, streamOutput);
            }
            // 创建OperatorChain内部算子之间的连接
            List<StreamOperator<?>> allOps = new ArrayList<>(chainedConfigs.size());
            this.chainEntryPoint = createOutputCollector(
               containingTask,
               configuration,
               chainedConfigs,
               userCodeClassloader,
               streamOutputMap,
               allOps,
               containingTask.getMailboxExecutorFactory());
            if (operatorFactory != null) {
               WatermarkGaugeExposingOutput<StreamRecord<OUT>> output = 
                  getChainEntryPoint();
               // 创建headOperator
               headOperator = StreamOperatorFactoryUtil.createOperator(
                     operatorFactory,
                     containingTask,
                     configuration,
                     output);
               headOperator.getMetricGroup().gauge(MetricNames.IO_CURRENT_OUTPUT_WATERMARK,
               output.getWatermarkGauge());
            } else {
               headOperator = null;
            }
            allOps.add(headOperator);
            this.allOperators = allOps.toArray(new StreamOperator<?>[allOps.size()]);
            success = true;
         }
         finally {
            // 如果创建不成功,则关闭StreamOutputs中的RecordWriterOutput
            // 这里防止内存泄漏
            if (!success) {
               for (RecordWriterOutput<?> output : this.streamOutputs) {
                  if (output != null) {
                     output.close();
                  }
               }
            }
         }
      }
      

      OperatorChain作用小结

      当OperatorChain创建完成后,就能正常接收StreamTaskInput中的数据元素了。在OperatorChain内部算子之间进行数据传递和处理,最终通过RecordWriterOutput组件将处理完成的数据元素发送到网络中,供下游的Task实例使用。

      对于OperatorChain内部Output接口的实现,这里暂不展开。

       

      3.创建RecordWriterOutput

      RecordWriterOutput用于将数据输出到网络指定位置。

      OperatorChain.createStreamOutput()逻辑如下:

      1. 获取输出边的OutputTag标签,判断当前Stream节点输出边是否为旁路输出,即在DataStream API中是否使用了旁路输出的相关方法。
      2. 返回RecordWriterOutput。RecordWriterOutput中包含RecordWriter组件,最终会通过RecordWriter将算子链处理完成的数据写入网络。
      private RecordWriterOutput<OUT> createStreamOutput(
            RecordWriter<SerializationDelegate<StreamRecord<OUT>>> recordWriter,
            StreamEdge edge,
            StreamConfig upStreamConfig,
            Environment taskEnvironment) {
         // 获取OutputTag
         OutputTag sideOutputTag = edge.getOutputTag(); 
         // 获取数据序列化器TypeSerializer
         TypeSerializer outSerializer = null;
         // 如果StreamEdge指定了OutputTag
         if (edge.getOutputTag() != null) {
            // 则进行边路输出
            outSerializer = upStreamConfig.getTypeSerializerSideOut(
                  edge.getOutputTag(), taskEnvironment.getUserClassLoader());
         } else {
            // 正常输出
            outSerializer = 
            upStreamConfig.getTypeSerializerOut(taskEnvironment.getUserClassLoader());
         }
         // 返回创建的RecordWriterOutput实例
         return new RecordWriterOutput<>(recordWriter, outSerializer, sideOutputTag, this);
      }
      

       

      StreamRecord将数据输出的逻辑

      在RecordWriterOutput.collect()方法中定义了StreamRecord数据的输出逻辑,实际上是调用pushToRecordWriter()方法将数据写入RecordWriter,最终通过RecordWriter组件进行数据元素的网络输出。

      public void collect(StreamRecord<OUT> record) {
         if (this.outputTag != null) {
            return;
         }
         pushToRecordWriter(record);
      }
      

       

      pushToRecordWriter发送数据

      1. 调用serializationDelegate.setInstance()方法,对接入的数据元素进行序列化操作,将数据元素转换成二进制格式。
      2. 调用recordWriter.emit()方法通过RecordWriter组件将serializationDelegate中序列化后的二进制数据输出到下游网络中。
      //RecordWriterOutput.pushToRecordWriter()
      private <X> void pushToRecordWriter(StreamRecord<X> record) {
         serializationDelegate.setInstance(record);
         try {
            recordWriter.emit(serializationDelegate);
         }
         catch (Exception e) {
            throw new RuntimeException(e.getMessage(), e);
         }
      }
       
      版权声明:本文内容来自第三方投稿或授权转载,原文地址:https://blog.csdn.net/hiliang521/article/details/136499192,作者:roman_日积跬步-终至千里,版权归原作者所有。本网站转在其作品的目的在于传递更多信息,不拥有版权,亦不承担相应法律责任。如因作品内容、版权等问题需要同本网站联系,请发邮件至ctyunbbs@chinatelecom.cn沟通。

      上一篇:C语言刷题 | 用%f控制符输出6位小数(17)

      下一篇:Java零基础入门之IO流详解(二)

      相关文章

      2025-05-16 09:15:17

      Linux系统基础-文件系统

      Linux系统基础-文件系统

      2025-05-16 09:15:17
      hello , 写入 , 文件 , 输出
      2025-05-14 10:33:31

      【数据结构】第一章——绪论(2)

      【数据结构】第一章——绪论(2)

      2025-05-14 10:33:31
      函数 , 实现 , 打印 , 理解 , 算法 , 输入 , 输出
      2025-05-14 10:33:25

      超级好用的C++实用库之网络

      在网络相关的项目中,我们经常需要去获取和设置设备的IP地址、子网掩码、网关地址、MAC地址等信息。这些信息一般与操作系统相关,在Windows系统和Linux系统上调用的接口是不一样的。

      2025-05-14 10:33:25
      Linux , 参数 , 地址 , 接口 , 网卡 , 返回值
      2025-05-14 09:51:15

      JAVA 两个类同时实现同一个接口

      在Java中,两个类同时实现同一个接口是非常常见的。接口定义了一组方法,实现接口的类必须提供这些方法的具体实现。

      2025-05-14 09:51:15
      Lambda , 函数 , 实现 , 接口 , 方法 , 表达式
      2025-05-13 09:49:12

      Java学习(动态代理的思想详细分析与案例准备)(1)

      Java学习(动态代理的思想详细分析与案例准备)(1)

      2025-05-13 09:49:12
      java , 代理 , 代码 , 对象 , 接口 , 方法 , 需要
      2025-05-12 08:58:16

      包裹分类

      包裹分类

      2025-05-12 08:58:16
      ID , 包裹 , 字母 , 数字 , 输出
      2025-05-09 09:30:05

      WebAPi接口安全之公钥私钥加密

      WebAPi接口安全之公钥私钥加密

      2025-05-09 09:30:05
      加密 , 参数 , 接口 , 请求 , 重写
      2025-05-09 08:50:35

      springboot实战学习(11)(更新用户基本信息接口主逻辑)

      springboot实战学习(11)(更新用户基本信息接口主逻辑)

      2025-05-09 08:50:35
      接口 , 方法 , 更新 , 用户 , 请求
      2025-05-09 08:50:35

      springboot实战学习(1)(开发模式与环境)

      springboot实战学习(1)(开发模式与环境)

      2025-05-09 08:50:35
      依赖 , 前端 , 后端 , 开发 , 接口 , 数据库 , 文档
      2025-05-08 09:03:38

      A+B问题

      A+B问题

      2025-05-08 09:03:38
      样例 , 格式 , 程序 , 输入 , 输出
      查看更多
      推荐标签

      作者介绍

      天翼云小翼
      天翼云用户

      文章

      33561

      阅读量

      5253616

      查看更多

      最新文章

      超级好用的C++实用库之网络

      2025-05-14 10:33:25

      JAVA 两个类同时实现同一个接口

      2025-05-14 09:51:15

      Java学习(动态代理的思想详细分析与案例准备)(1)

      2025-05-13 09:49:12

      WebAPi接口安全之公钥私钥加密

      2025-05-09 09:30:05

      springboot实战学习(11)(更新用户基本信息接口主逻辑)

      2025-05-09 08:50:35

      springboot实战学习(1)(开发模式与环境)

      2025-05-09 08:50:35

      查看更多

      热门文章

      JAVA__接口的作用

      2023-04-18 14:14:13

      什么是api接口

      2023-03-22 09:03:21

      kotlin匿名内部类与接口实现

      2023-04-18 14:15:13

      Go 语言入门很简单 -- 13. Go 接口 #私藏项目实操分享#

      2023-04-21 03:11:48

      SpringBoot写的后端API接口如何写得更优雅

      2023-06-15 06:37:47

      ts重点学习47-接口与类型别名得异同笔记

      2023-03-16 06:47:13

      查看更多

      热门标签

      java Java python 编程开发 代码 开发语言 算法 线程 Python html 数组 C++ 元素 javascript c++
      查看更多

      相关产品

      弹性云主机

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

      天翼云电脑(公众版)

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

      对象存储

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

      云硬盘

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

      查看更多

      随机文章

      shell编程-数组与运算符详解(超详细)

      从零做软件开发项目系列之三——系统设计

      信息学奥赛一本通(C++)在线评测系统——基础(一)C++语言—— 1053:最大数输出

      一次Http Get请求健壮性问题的排查过程

      初学Java,LinkedList功能最全的集合类(二十九)

      springboot实现图片或者其他文件回显功能

      • 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号