Flink中异步AsyncIO的实现 (源码分析)

2024-05-13 21:58

本文主要是介绍Flink中异步AsyncIO的实现 (源码分析),希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

先上张图整体了解Flink中的异步io

阿里贡献给flink的,优点就不说了嘛,官网上都有,就是写库不会柱塞性能更好

然后来看一下, Flink 中异步io主要分为两种

  一种是有序Ordered

  一种是无序UNordered

主要区别是往下游output的顺序(注意这里顺序不是写库的顺序既然都异步了写库的顺序自然是无法保证的),有序的会按接收的顺序继续往下游output发送,无序就是谁先处理完谁就先往下游发送

两张图了解这两种模式的实现

有序:record数据会通过异步线程写库,Emitter是一个守护进程,会不停的拉取queue头部的数据,如果头部的数据异步写库完成,Emitter将头数据往下游发送,如果头元素还没有异步写库完成,柱塞 www.wityx.com     

无序:record数据会通过异步线程写库,这里有两个queue,一开始放在uncompleteedQueue,当哪个record异步写库成功后就直接放到completedQueue中,Emitter是一个守护进程,completedQueue只要有数据,会不停的拉取queue数据往下游发送 

可以看到原理还是很简单的,两句话就总结完了,就是利用queue和java的异步线程,现在来看下源码

这里AsyncIO在Flink中被设计成operator中的一种,自然去OneInputStreamOperator的实现类中去找

于是来看一下AsyncWaitOperator.java

看到它的open方法(open方法会在taskmanager启动job的时候全部统一调用,可以翻一下以前的文章)

这里启动了一个守护线程Emitter,来看下线程具体做了什么

 1处拉取数据,2处就是常规的将拉取到的数据往下游emit,Emitter拉取数据,这里先不讲因为分为有序的和无序的

 这里已经知道了这个Emitter的作用是循环的拉取数据往下游发送

 回到AsyncWaitOperator.java在它的open方法初始化了Emitter,那它是如何处理接收到的数据的呢,看它的ProcessElement()方法

 其实主要就是三个个方法

先是!!!将record封装成了一个包装类StreamRecordQueueEntry,主要是这个包装类的构造方法中,创建了一个CompleteableFuture(这个的complete方法其实会等到用户代码执行的时候用户自己决定什么时候完成)

1处主要就是讲元素加入到了对应的queue,这里也分为两种有序和无序的

这里也先不讲这两种模式加入数据的区别

接着2处就是调用用户的代码了,来看看官网的异步io的例子

 给了一个Future作为参数,用户自己起了一个线程(这里思考一下就知道了为什么要新起一个异步线程去执行,因为如果不起线程的话,那processElement方法就柱塞了,无法异步了)去写库读库等,然后调用了这个参数的complete方法(也就是前面那个包装类中的CompleteableFuture)并且传入了一个结果

看下complete方法源码

 这个resultFuture是每个record的包装类StreamRecordQueueEntry的其中一个属性是一个CompletableFuture

 那现在就清楚了,用户代码在自己新起的线程中当自己的逻辑执行完以后会使这个异步线程结束,并输入一个结果

 那这个干嘛用的呢

最开始的图中看到有序和无序实现原理,有序用一个queue,无序用两个queue分别就对应了

OrderedStreamElementQueue类中

 UnorderedStreamElementQueue类中

回到前面有两个地方没有细讲,一是两种模式的Emitter是如何拉取数据的,二是两种模式下数据是如何加入OrderedStreamElementQueue的

有序模式:

1.先来看一下有序模式的,Emitter的数据拉取,和数据的加入

    其tryPut()方法

     onComplete方法

       onCompleteHandler方法

  这里比较绕,先将接收的数据加入queue中,然后onComplete()中当上一个异步线程getFuture() 其实就是每个元素包装类里面的那个CompletableFuture,当他结束时(会在用户方法用户调用complete时结束)异步调用传入的对象的 accept方法,accept方法中调用了onCompleteHandler()方法,onCompleteHandler方法中会判断queue是否为空,以及queue的头元素是否完成了用户的异步方法,当完成的时候,就会将headIsCompleted这个对象signalAll()唤醒

2.接着看有序模式Emitter的拉取数据

   这里有序方式拉取数据的逻辑很清晰,如果为空或者头元素没有完成用户的异步方法,headIsCompleted这个对象会wait住(上面可以知道,当加入元素的到queue且头元素完成异步方法的时候会signalAll())然后将头数据返回,往下游发送

这样就实现了有序发送,因为Emitter只拉取头元素且已经完成用户异步方法的头元素

无序模式: 

  这里和有序模式就大同小异了,只是变成了,接收数据后直接加入uncompletedQueue,当数据完成异步方法的时候就,放到completedQueue里面去并signalAll(),只要completedqueue里面有数据,Emitter就拉取往下发

这样就实现了无序模式,也就是异步写入谁先处理完就直接放到完成队列里面去,然后往下发,不用管接收数据的顺序

这篇关于Flink中异步AsyncIO的实现 (源码分析)的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/986939

相关文章

使用Python实现IP地址和端口状态检测与监控

《使用Python实现IP地址和端口状态检测与监控》在网络运维和服务器管理中,IP地址和端口的可用性监控是保障业务连续性的基础需求,本文将带你用Python从零打造一个高可用IP监控系统,感兴趣的小伙... 目录概述:为什么需要IP监控系统使用步骤说明1. 环境准备2. 系统部署3. 核心功能配置系统效果展

Python实现微信自动锁定工具

《Python实现微信自动锁定工具》在数字化办公时代,微信已成为职场沟通的重要工具,但临时离开时忘记锁屏可能导致敏感信息泄露,下面我们就来看看如何使用Python打造一个微信自动锁定工具吧... 目录引言:当微信隐私遇到自动化守护效果展示核心功能全景图技术亮点深度解析1. 无操作检测引擎2. 微信路径智能获

Python中pywin32 常用窗口操作的实现

《Python中pywin32常用窗口操作的实现》本文主要介绍了Python中pywin32常用窗口操作的实现,pywin32主要的作用是供Python开发者快速调用WindowsAPI的一个... 目录获取窗口句柄获取最前端窗口句柄获取指定坐标处的窗口根据窗口的完整标题匹配获取句柄根据窗口的类别匹配获取句

在 Spring Boot 中实现异常处理最佳实践

《在SpringBoot中实现异常处理最佳实践》本文介绍如何在SpringBoot中实现异常处理,涵盖核心概念、实现方法、与先前查询的集成、性能分析、常见问题和最佳实践,感兴趣的朋友一起看看吧... 目录一、Spring Boot 异常处理的背景与核心概念1.1 为什么需要异常处理?1.2 Spring B

Python中的Walrus运算符分析示例详解

《Python中的Walrus运算符分析示例详解》Python中的Walrus运算符(:=)是Python3.8引入的一个新特性,允许在表达式中同时赋值和返回值,它的核心作用是减少重复计算,提升代码简... 目录1. 在循环中避免重复计算2. 在条件判断中同时赋值变量3. 在列表推导式或字典推导式中简化逻辑

Python位移操作和位运算的实现示例

《Python位移操作和位运算的实现示例》本文主要介绍了Python位移操作和位运算的实现示例,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一... 目录1. 位移操作1.1 左移操作 (<<)1.2 右移操作 (>>)注意事项:2. 位运算2.1

如何在 Spring Boot 中实现 FreeMarker 模板

《如何在SpringBoot中实现FreeMarker模板》FreeMarker是一种功能强大、轻量级的模板引擎,用于在Java应用中生成动态文本输出(如HTML、XML、邮件内容等),本文... 目录什么是 FreeMarker 模板?在 Spring Boot 中实现 FreeMarker 模板1. 环

Qt实现网络数据解析的方法总结

《Qt实现网络数据解析的方法总结》在Qt中解析网络数据通常涉及接收原始字节流,并将其转换为有意义的应用层数据,这篇文章为大家介绍了详细步骤和示例,感兴趣的小伙伴可以了解下... 目录1. 网络数据接收2. 缓冲区管理(处理粘包/拆包)3. 常见数据格式解析3.1 jsON解析3.2 XML解析3.3 自定义

SpringMVC 通过ajax 前后端数据交互的实现方法

《SpringMVC通过ajax前后端数据交互的实现方法》:本文主要介绍SpringMVC通过ajax前后端数据交互的实现方法,本文给大家介绍的非常详细,对大家的学习或工作具有一定的参考借鉴价... 在前端的开发过程中,经常在html页面通过AJAX进行前后端数据的交互,SpringMVC的controll

Spring Security自定义身份认证的实现方法

《SpringSecurity自定义身份认证的实现方法》:本文主要介绍SpringSecurity自定义身份认证的实现方法,下面对SpringSecurity的这三种自定义身份认证进行详细讲解,... 目录1.内存身份认证(1)创建配置类(2)验证内存身份认证2.JDBC身份认证(1)数据准备 (2)配置依