详解Flink组件通信——RPC协议 flink cep or
ztj100 2024-12-25 16:49 23 浏览 0 评论
Flink组件通讯过程
RPC(本地/远程)调用,底层是通过Akka提供的tell/ask方法进行通信。
Flink中RPC框架中涉及的主要类:
1、RpcGateway
Flink的RPC协议通过RpcGateway来定义,主要定义通信行为;用于远程调用RpcEndpoint的某些方法,可以理解为客服端代理。
若想与远端Actor通信,则必须提供地址(ip和port),如在Flink-on-Yarn模式下,JobMaster会先启动ActorSystem,此时TaskExecutor的Container还未分配,后面与TaskExecutor通信时,必须让其提供对应地址。
从类继承图可以看到基本上所有组件都实现了RpcGateway接口,其代码如下:
public interface RpcGateway {
/**
* Returns the fully qualified address underwhich the associated rpc endpoint is reachable.
*
* @return Fully qualified (RPC) address underwhich the associated rpc endpoint is reachable
*/
String getAddress();
/**
* Returns the fully qualified hostname underwhich the associated rpc endpoint is reachable.
*
* @return Fully qualified hostname under whichthe associated rpc endpoint is reachable
*/
String getHostname();
}
2、RpcEndpoint
RpcEndpoint是通信终端,提供RPC服务组件的生命周期管理(start、stop)。每个RpcEndpoint对应了一个路径(endpointId和actorSystem共同确定),每个路径对应一个Actor,其实现了RpcGateway接口,其构造函数如下:
protected RpcEndpoint(final RpcService rpcService, final String endpointId) {
// 保存rpcService和endpointId
this.rpcService =checkNotNull(rpcService, "rpcService");
this.endpointId =checkNotNull(endpointId, "endpointId");
// 通过RpcService启动RpcServer
this.rpcServer =rpcService.startServer(this);
// 主线程执行器,所有调用在主线程中串行执行
this.mainThreadExecutor= new MainThreadExecutor(rpcServer, this::validateRunsInMainThread);
}
构造的时候调用rpcService.startServer()启动RpcServer,进入可以接收处理请求的状态,最后将RpcServer绑定到主线程上真正执行起来。
在RpcEndpoint中还定义了一些方法如runAsync(Runnable)、callAsync(Callable, Time)方法来执行Rpc调用,值得注意的是在Flink的设计中,对于同一个Endpoint,所有的调用都运行在主线程,因此不会有并发问题,当启动RpcEndpoint/进行Rpc调用时,其会委托RcpServer进行处理。
3、RpcService和RpcServer
Akka有两种核心的异步通信方式:tell和ask。
RpcService 和 RpcServer是RpcEndPoint的成员变量。
1)RpcService 是Rpc服务的接口,其主要作用如下:
- 根据提供的RpcEndpoint来启动和停止RpcServer(Actor);
- 根据提供的地址连接到RpcServer,并返回一个RpcGateway;
- 延迟/立刻调度Runnable、Callable;
在Flink中实现类为AkkaRpcService,是 Akka 的 ActorSystem 的封装,基本可以理解成 ActorSystem 的一个适配器。在ClusterEntrypoint(JobMaster)和TaskManagerRunner(TaskExecutor)启动的过程中初始化并启动。
AkkaRpcService中封装了ActorSystem,并保存了ActorRef到RpcEndpoint的映射关系。RpcService跟RpcGateway类似,也提供了获取地址和端口的方法。
在构造RpcEndpoint时会启动指定rpcEndpoint上的RpcServer,其会根据RpcEndpoint类型(FencedRpcEndpoint或其他)来创建不同的AkkaRpcActor(FencedAkkaRpcActor或AkkaRpcActor),并将RpcEndpoint和AkkaRpcActor对应的ActorRef保存起来,AkkaRpcActor是底层Akka调用的实际接收者,RPC的请求在客户端被封装成RpcInvocation对象,以Akka消息的形式发送。
最终使用动态代理将所有的消息转发到InvocationHandler,具体代码如下:
public <Cextends RpcEndpoint & RpcGateway> RpcServer startServer(C rpcEndpoint) {
... ...
// 生成RpcServer对象,而后对该server的调用都会进入Handler的invoke方法处理,handler实现了多个接口的方法
// 生成一个包含这些接口的代理,将调用转发到InvocationHandler
@SuppressWarnings("unchecked")
RpcServerserver = (RpcServer) Proxy.newProxyInstance(
classLoader,
implementedRpcGateways.toArray(newClass<?>[implementedRpcGateways.size()]),
akkaInvocationHandler);
return server;
}
2)RpcServer负责接收响应远端RPC消息请求。有两个实现:
- AkkaInvocationHandler
- FencedAkkaInvocationHandler
RpcServer的启动是通知底层的AkkaRpcActor切换为START状态,开始处理远程调用请求:
class AkkaInvocationHandler implements InvocationHandler, AkkaBasedEndpoint,RpcServer {
@Override
public void start() {
rpcEndpoint.tell(ControlMessages.START,ActorRef.noSender());
}
}
4、AkkaRpcActor
AkkaRpcActor是Akka的具体实现,主要负责处理如下类型消息:
1)本地Rpc调用LocalRpcInvocation
会指派给RpcEndpoint进行处理,如果有响应结果,则将响应结果返还给Sender。
2)RunAsync & CallAsync
这类消息带有可执行的代码,直接在Actor的线程中执行。
3)控制消息ControlMessages
用来控制Actor行为,START启动,STOP停止,停止后收到的消息会丢弃掉。
5、RPC交互过程
RPC通信过程分为请求和响应。
1、 RPC请求发送
在RpcService中调用connect()方法与对端的RpcEndpoint建立连接,connect()方法根据给的地址返回InvocationHandler(AkkaInvocationHandler或FencedAkkaInvocationHandler,也就是RpcServer的实现)。
前面分析到客户端提供代理对象RpcServer,代理对象会调用AkkaInvocationHandler的invoke方法并传入RPC调用的方法和参数信息,代码如下:
AkkaInvocationHandler.java
public Object invoke(Object proxy, Method method, Object[] args) throws Throwable {
Class<?> declaringClass =method.getDeclaringClass();
Object result;
// 判断方法所属的class
if(declaringClass.equals(AkkaBasedEndpoint.class) ||
declaringClass.equals(Object.class) ||
declaringClass.equals(RpcGateway.class) ||
declaringClass.equals(StartStoppable.class) ||
declaringClass.equals(MainThreadExecutable.class)||
declaringClass.equals(RpcServer.class)) {
result = method.invoke(this, args);
} else if(declaringClass.equals(FencedRpcGateway.class)) {
throw new UnsupportedOperationException("AkkaInvocationHandler does not support thecall FencedRpcGateway#" +
method.getName() + ". This indicates thatyou retrieved a FencedRpcGateway without specifying a " +
"fencingtoken. Please use RpcService#connect(RpcService, F, Time) with F being thefencing token to " +
"retrieve aproperly FencedRpcGateway.");
} else {
// rpc调用
result = invokeRpc(method, args);
}
return result;
}
代码中判断所属的类,如果是RPC方法,则调用invokeRpc方法。将方法调用封装为RPCInvocation消息。如果是本地则生成LocalRPCInvocation,本地消息不需要序列化,如果是远程调用则创建RemoteRPCInvocation。
判断远程方法调用是否需要等待结果,如果无需等待(void),则使用向Actor发送tell类型的消息,如果需要返回结果,则向Acrot发送ask类型的消息,代码如下:
AkkaInvocationHandler.java
private Object invokeRpc(Method method, Object[]args) throws Exception {
String methodName = method.getName();
Class<?>[] parameterTypes =method.getParameterTypes();
Annotation[][] parameterAnnotations =method.getParameterAnnotations();
Time futureTimeout =extractRpcTimeout(parameterAnnotations, args, timeout);
final RpcInvocation rpcInvocation =createRpcInvocationMessage(methodName, parameterTypes, args);
Class<?> returnType =method.getReturnType();
final Object result;
if (Objects.equals(returnType, Void.TYPE)) {
tell(rpcInvocation);
result = null;
} else {
// Capture the call stack. Itis significantly faster to do that via an exception than
// via Thread.getStackTrace(),because exceptions lazily initialize the stack trace, initially only
// capture a lightweightnative pointer, and convert that into the stack trace lazily when needed.
final ThrowablecallStackCapture = captureAskCallStack ? new Throwable() : null;
// execute an asynchronouscall
// 异步调用等待返回
final CompletableFuture<?> resultFuture = ask(rpcInvocation, futureTimeout);
final CompletableFuture<Object> completableFuture = newCompletableFuture<>();
resultFuture.whenComplete((resultValue,failure) -> {
if (failure != null) {
completableFuture.completeExceptionally(resolveTimeoutException(failure,callStackCapture, method));
} else {
completableFuture.complete(deserializeValueIfNeeded(resultValue,method));
}
});
if (Objects.equals(returnType,CompletableFuture.class)) {
// 如果返回值是CompletableFuture类型,不用阻塞等待返回,直接返回Future对象
result =completableFuture;
} else {
try {
// 如果返回值不是CompletableFuture类型,阻塞等待返回
result = completableFuture.get(futureTimeout.getSize(),futureTimeout.getUnit());
} catch(ExecutionException ee) {
throw new RpcException("Failure while obtaining synchronous RPC result.",ExceptionUtils.stripExecutionException(ee));
}
}
}
return result;
}
2、 RPC请求响应
RPC消息通过RpcEndpoint所绑定的Actor的ActorRef发送的,AkkaRpcActor是消息接收的入口,AkkaRpcActor在RpcEndpoint中构造生成,负责将消息交给不同的方法进行处理。
AkkaRpcActor.java
public Receive createReceive() {
return ReceiveBuilder.create()
.match(RemoteHandshakeMessage.class,this::handleHandshakeMessage)
.match(ControlMessages.class,this::handleControlMessage)
.matchAny(this::handleMessage)
.build();
}
接收的消息有3种:
1)握手消息
在客户端构造时会通过ActorSelection发送过来。收到消息后检查接口、版本是否匹配。
AkkaRpcActor.java
private void handleHandshakeMessage(RemoteHandshakeMessagehandshakeMessage) {
if(!isCompatibleVersion(handshakeMessage.getVersion())) {
// 版本不兼容异常处理
sendErrorIfSender(newAkkaHandshakeException(
String.format(
"Versionmismatch between source (%s) and target (%s) rpc component. Please verify thatall components have the same version.",
handshakeMessage.getVersion(),
getVersion())));
} else if(!isGatewaySupported(handshakeMessage.getRpcGateway())) {
// RpcGateway不匹配异常处理
sendErrorIfSender(newAkkaHandshakeException(
String.format(
"The rpcendpoint does not support the gateway %s.",
handshakeMessage.getRpcGateway().getSimpleName())));
} else {
getSender().tell(new Status.Success(HandshakeSuccessMessage.INSTANCE),getSelf());
}
}
2)控制消息
在RpcEndpoint调用start方法后,会向自身发送一条Processing.START消息来转换当前Actor的状态为STARTED,STOP也类似,并且只有在Actor状态为STARTED时才会处理RPC请求。
AkkaRpcActor.java
private void handleControlMessage(ControlMessages controlMessage) {
try {
switch (controlMessage) {
case START:
state =state.start(this);
break;
case STOP:
state =state.stop();
break;
case TERMINATE:
state =state.terminate(this);
break;
default:
handleUnknownControlMessage(controlMessage);
}
} catch (Exception e) {
this.rpcEndpointTerminationResult= RpcEndpointTerminationResult.failure(e);
throw e;
}
}
3)RPC消息
通过解析RpcInvocation获取方法名和参数类型,并从RpcEndpoint类中找到Method对象,通过反射调用该方法。如果有返回结果,会以Akka消息的形式发送回发送者。
AkkaRpcActor.java
private void handleMessage(final Object message) {
if (state.isRunning()) {
mainThreadValidator.enterMainThread();
try {
handleRpcMessage(message);
} finally {
mainThreadValidator.exitMainThread();
}
} else {
log.info("The rpcendpoint {} has not been started yet. Discarding message {} until processing isstarted.",
rpcEndpoint.getClass().getName(),
message.getClass().getName());
sendErrorIfSender(newAkkaRpcException(
String.format("Discardmessage, because the rpc endpoint %s has not been started yet.",rpcEndpoint.getAddress())));
}
}
更多精彩内容:
相关推荐
- sharding-jdbc实现`分库分表`与`读写分离`
-
一、前言本文将基于以下环境整合...
- 三分钟了解mysql中主键、外键、非空、唯一、默认约束是什么
-
在数据库中,数据表是数据库中最重要、最基本的操作对象,是数据存储的基本单位。数据表被定义为列的集合,数据在表中是按照行和列的格式来存储的。每一行代表一条唯一的记录,每一列代表记录中的一个域。...
- MySQL8行级锁_mysql如何加行级锁
-
MySQL8行级锁版本:8.0.34基本概念...
- mysql使用小技巧_mysql使用入门
-
1、MySQL中有许多很实用的函数,好好利用它们可以省去很多时间:group_concat()将取到的值用逗号连接,可以这么用:selectgroup_concat(distinctid)fr...
- MySQL/MariaDB中如何支持全部的Unicode?
-
永远不要在MySQL中使用utf8,并且始终使用utf8mb4。utf8mb4介绍MySQL/MariaDB中,utf8字符集并不是对Unicode的真正实现,即不是真正的UTF-8编码,因...
- 聊聊 MySQL Server 可执行注释,你懂了吗?
-
前言MySQLServer当前支持如下3种注释风格:...
- MySQL系列-源码编译安装(v5.7.34)
-
一、系统环境要求...
- MySQL的锁就锁住我啦!与腾讯大佬的技术交谈,是我小看它了
-
对酒当歌,人生几何!朝朝暮暮,唯有己脱。苦苦寻觅找工作之间,殊不知今日之事乃我心之痛,难道是我不配拥有工作嘛。自面试后他所谓的等待都过去一段时日,可惜在下京东上的小金库都要见低啦。每每想到不由心中一...
- MySQL字符问题_mysql中字符串的位置
-
中文写入乱码问题:我输入的中文编码是urf8的,建的库是urf8的,但是插入mysql总是乱码,一堆"???????????????????????"我用的是ibatis,终于找到原因了,我是这么解决...
- 深圳尚学堂:mysql基本sql语句大全(三)
-
数据开发-经典1.按姓氏笔画排序:Select*FromTableNameOrderByCustomerNameCollateChinese_PRC_Stroke_ci_as//从少...
- MySQL进行行级锁的?一会next-key锁,一会间隙锁,一会记录锁?
-
大家好,是不是很多人都对MySQL加行级锁的规则搞的迷迷糊糊,一会是next-key锁,一会是间隙锁,一会又是记录锁。坦白说,确实还挺复杂的,但是好在我找点了点规律,也知道如何如何用命令分析加...
- 一文讲清怎么利用Python Django实现Excel数据表的导入导出功能
-
摘要:Python作为一门简单易学且功能强大的编程语言,广受程序员、数据分析师和AI工程师的青睐。本文系统讲解了如何使用Python的Django框架结合openpyxl库实现Excel...
- 用DataX实现两个MySQL实例间的数据同步
-
DataXDataX使用Java实现。如果可以实现数据库实例之间准实时的...
- MySQL数据库知识_mysql数据库基础知识
-
MySQL是一种关系型数据库管理系统;那废话不多说,直接上自己以前学习整理文档:查看数据库命令:(1).查看存储过程状态:showprocedurestatus;(2).显示系统变量:show...
- 如何为MySQL中的JSON字段设置索引
-
背景MySQL在2015年中发布的5.7.8版本中首次引入了JSON数据类型。自此,它成了一种逃离严格列定义的方式,可以存储各种形状和大小的JSON文档,例如审计日志、配置信息、第三方数据包、用户自定...
你 发表评论:
欢迎- 一周热门
-
-
MySQL中这14个小玩意,让人眼前一亮!
-
旗舰机新标杆 OPPO Find X2系列正式发布 售价5499元起
-
【VueTorrent】一款吊炸天的qBittorrent主题,人人都可用
-
面试官:使用int类型做加减操作,是线程安全吗
-
C++编程知识:ToString()字符串转换你用正确了吗?
-
【Spring Boot】WebSocket 的 6 种集成方式
-
PyTorch 深度学习实战(26):多目标强化学习Multi-Objective RL
-
pytorch中的 scatter_()函数使用和详解
-
与 Java 17 相比,Java 21 究竟有多快?
-
基于TensorRT_LLM的大模型推理加速与OpenAI兼容服务优化
-
- 最近发表
- 标签列表
-
- idea eval reset (50)
- vue dispatch (70)
- update canceled (42)
- order by asc (53)
- spring gateway (67)
- 简单代码编程 贪吃蛇 (40)
- transforms.resize (33)
- redisson trylock (35)
- 卸载node (35)
- np.reshape (33)
- torch.arange (34)
- npm 源 (35)
- vue3 deep (35)
- win10 ssh (35)
- vue foreach (34)
- idea设置编码为utf8 (35)
- vue 数组添加元素 (34)
- std find (34)
- tablefield注解用途 (35)
- python str转json (34)
- java websocket客户端 (34)
- tensor.view (34)
- java jackson (34)
- vmware17pro最新密钥 (34)
- mysql单表最大数据量 (35)