Spark2.4.0源码——RpcEnv

参考《Spark内核设计的艺术:架构设计与实现——耿嘉安》

NettyRpcEnv概述

 Spark的NettyRpc环境的一些重要组件:

1private[netty] val transportConf = SparkTransportConf.fromSparkConf(...) 2 3private val dispatcher: Dispatcher = new Dispatcher(this, numUsableCores) 4 5private val streamManager = new NettyStreamManager(this) 6 7private val transportContext = new TransportContext(transportConf, 8  new NettyRpcHandler(dispatcher, this, streamManager)) 9 10//用于创建TransportClient的工厂类 11private val clientFactory = transportContext.createClientFactory(createClientBootstraps()) 12 13//volatile 关键字保证变量在多线程之间的可见性 14@volatile private var server: TransportServer = _

绪:RpcEndpoint&RpcEndpointRef

RpcEndpoint

RpcEndpoint是对能够处理RPC请求,给某一特定服务提供本地及跨节点调用的RPC组件的抽象,所有运行于RPC框架上的实体都应该继承trait RPCEndpoint。

1package org.apache.spark.rpc 2 3import org.apache.spark.SparkException 4 5//创建RpcEnv的工厂类,必须有一个空构造函数才能通过反射创建 6private[spark] trait RpcEnvFactory { 7 8 def create(config: RpcEnvConfig): RpcEnv 9} 10 11private[spark] trait RpcEndpoint { 12 13 //当前RpcEndpoint所属的RpcEnv 14 val rpcEnv: RpcEnv 15 16 //获取RpcEndpoint相关联的RpcEndpointRef 17 final def self: RpcEndpointRef = { 18 require(rpcEnv != null, "rpcEnv has not been initialized") 19 rpcEnv.endpointRef(this) 20 } 21 22 //接收消息并处理,不回复客户端 23 def receive: PartialFunction[Any, Unit] = { 24 case _ => throw new SparkException(self + " does not implement 'receive'") 25 } 26 27 //接收消息并处理,通过RpcCallContext回复客户端 28 def receiveAndReply(context: RpcCallContext): PartialFunction[Any, Unit] = { 29 case _ => context.sendFailure(new SparkException(self + " won't reply anything")) 30 } 31 32 //onError、onConnected、onDisconnected、onNetworkError、onStart、onStop顾名思义 33 34 //用于停止当前RpcEndpoint,注意onStop是trait定义的抽象方法,在停止RpcEndpoint时调用,做一些收尾工作 35 final def stop(): Unit = { 36 val _self = self 37 if (_self != null) { 38 rpcEnv.stop(_self) 39 } 40 } 41} 42 43//线程安全的、串行处理消息的ThreadSafeRpcEndpoint 44private[spark] trait ThreadSafeRpcEndpoint extends RpcEndpoint

trait ThreadSafeRpcEndpoint/... extends RpcEndpoint
ThreadSafeRpcEndpoint主要用于消息的串行处理,必须是线程安全的
Master/Worker/HeartbeatReceiver/... extends ThreadSafeRpcEndpoint

RpcEndpointRef

要向一个远端RpcEndpoint发送请求,就必须持有这个RpcEndpoint的远程引用RpcEndpointRef,它是线程安全的。

1private[spark] abstract class RpcEndpointRef(conf: SparkConf) 2 extends Serializable with Logging { 3 4 //rpc最大重连次数,默认3,可使用spark.rpc.numRetries属性配置 5 private[this] val maxRetries = RpcUtils.numRetries(conf) 6 //rpc每次重连等待的毫秒数,默认3s,可使用spark.rpc.retry.wait属性配置 7 private[this] val retryWaitMs = RpcUtils.retryWaitMs(conf) 8 //rpc的ask操作默认超时时间,默认120s,可使用spark.rpc.askTimeout(优先级高)/spark.network.timeout属性配置 9 private[this] val defaultAskTimeout = RpcUtils.askRpcTimeout(conf) 10 11 //返回当前RpcEndpointRef对应的RpcEndpoint的RPC地址 12 def address: RpcAddress 13 14 //返回当前RpcEndpointRef对应的RpcEndpoint的名称 15 def name: String 16 17 //发送单向异步的消息到相应的RpcEndpoint.receive。 18 def send(message: Any): Unit 19 20 //发送一条消息到相应的RpcEndpoint.receiveAndReply,并在指定的超时内接收处理结果。此方法只发送消息一次,从不重试。 21 def ask[T: ClassTag](message: Any, timeout: RpcTimeout): Future[T] 22 def ask[T: ClassTag](message: Any): Future[T] = ask(message, defaultAskTimeout) 23 24 //发送同步请求到相应的RpcEndpoint.receiveAndReply,并在超时时间内等待处理结果,当抛出异常时会请求重试次数以内的重连。 25 def askSync[T: ClassTag](message: Any): T = askSync(message, defaultAskTimeout) 26 def askSync[T: ClassTag](message: Any, timeout: RpcTimeout): T = { 27 val future = ask[T](message, timeout) 28 timeout.awaitResult(future) 29 } 30 31}

消息投递规则:
at-most-once:投递0或1此,消息可能会丢失
at-least-once:潜在地多次投递并保证至少成功一次,消息可能会重复
exactly-once:准确发送一次,消息不会丢失也不会重复

1 TransportConf

RPC传输上下文配置类,用于创建TransportClientFactory和TransportServer。

1//通过SparkTransportConf的fromSparkConf方法来构建TransportConf需要三个参数:sparkConf、模块名module和可用内核数 2private[netty] val transportConf = SparkTransportConf.fromSparkConf( 3 //先克隆SparkConf并设置节点间取数据的连接数 4 conf.clone.set("spark.rpc.io.numConnectionsPerPeer", "1"), 5 //设置模块名 6 "rpc", 7 //Netty传输线程数,如果小于或等于0,线程数就是系统可用处理器的数量,最多为8线程。 8 conf.getInt("spark.rpc.io.threads", numUsableCores))

2 Dispatcher

Dispatcher负责将消息路由到应该对此消息处理的RpcEndpoint,可以提高NettyRpcEnv对消息的异步处理和并行处理能力。

private val dispatcher: Dispatcher = new Dispatcher(this, numUsableCores)

  

基本概念:

InboxMessage:Inbox盒子内的消息,是一个trait,所有类型的RPC消息都要继承自InboxMessage。

Inbox:端点内的盒子,每个RpcEndpoint都有一个对应的盒子,这个盒子有存储InboxMessage的列表messages,所有的消息都缓存在messages并由RpcEndpoint异步处理。

EndpointData:RPC端点数据,包括RpcEndpoint、NettyRpcEndpointRef和Inbox等属于同一个端点的实例。

endpoints:端点实例RpcEndpoint与EndpointData之间映射关系的缓存。

endpointRefs:端点实例RpcEndpoint与RpcEndpointRef之间映射关系的缓存.

receivers:存储EndpointData的阻塞队列,只有Inbox中有消息的EndpointData才会被放入此队列。

stopped:Dispatcher是否停止的状态。

threadPool:用于对消息进行调度的线程池,里面运行的任务都是MessageLoop。

2.1 Dispatcher注册RpcEndpoint

1def registerRpcEndpoint(name: String, endpoint: RpcEndpoint): NettyRpcEndpointRef = { 2 //使用当前RpcEndpoint所在的NettyRpcEnv的地址和RpcEndpoint的名称创建RpcEndpointAddress对象 3 val addr = RpcEndpointAddress(nettyEnv.address, name) 4 //创建RpcEndpoint的引用对象 5 val endpointRef = new NettyRpcEndpointRef(nettyEnv.conf, addr, nettyEnv) 6 synchronized { 7 if (stopped) { 8 throw new IllegalStateException("RpcEnv has been stopped") 9 } 10 if (endpoints.putIfAbsent(name, new EndpointData(name, endpoint, endpointRef)) != null) { 11 throw new IllegalArgumentException(s"There is already an RpcEndpoint called $name") 12 } 13 //创建EndpointData并放入endpoints缓存 14 val data = endpoints.get(name) 15 //将RpcEndpoint与NettyRpcEndpointRef映射关系放入endpointRefs缓存 16 endpointRefs.put(data.endpoint, data.ref) 17 //将EndpointData放入阻塞队列receivers,由于EndpointData是新建的,内部会新建Inbox并执行Inbox的主构造函数, 18 //向Inbox自身的messages列表中放入OnStart消息,MessageLoop线程会取出此EndpointData并调用当前Inbox的process方法 19 //处理OnStart消息,启动与此Inbox相关联的Endpoint。 20 receivers.offer(data) // for the OnStart message 21 } 22 endpointRef 23 }

   2.2 Dispatcher的调度原理

1private val threadpool: ThreadPoolExecutor = { 2 //获取可用处理器数,numUsableCores是NettyRpcEnv的入参,如果大于0则等于numUsableCores,否则为当前系统可用处理器 3 val availableCores = 4 if (numUsableCores > 0) numUsableCores else Runtime.getRuntime.availableProcessors() 5 //获取当前线程池的大小,默认为2和可用处理器之间的最大值,可用spark.rpc.netty.dispatcher.numThreads属性配置 6 val numThreads = nettyEnv.conf.getInt("spark.rpc.netty.dispatcher.numThreads", 7 math.max(2, availableCores)) 8 //创建线程池 9 val pool = ThreadUtils.newDaemonFixedThreadPool(numThreads, "dispatcher-event-loop") 10 //启动多个线程运行MessageLoop任务 11 for (i <- 0 until numThreads) { 12 pool.execute(new MessageLoop) 13 } 14 //返回线程池引用 15 pool 16 } 17 18 /** Message loop used for dispatching messages. */ 19 private class MessageLoop extends Runnable { 20 override def run(): Unit = { 21 try { 22 while (true) { 23 try { 24 //在阻塞队列中获取EndpointData 25 val data = receivers.take() 26 //如果EndpointData是空数据,则将它重新放回队列并直接返回,这样可以让其他MessageLoop获取到这个空EndpointData并结束线程        //private val PoisonPill = new EndpointData(null,null,null) 27 if (data == PoisonPill) { 28 // Put PoisonPill back so that other MessageLoops can see it. 29 receivers.offer(PoisonPill) 30 return 31 } 32 //调用inbox的process方法对消息进行处理 33 data.inbox.process(Dispatcher.this) 34 } catch { 35 case NonFatal(e) => logError(e.getMessage, e) 36 } 37 } 38 } catch { 39 case ie: InterruptedException => // exit 40 } 41 } 42 }

   Inbox的process方法:

1def process(dispatcher: Dispatcher): Unit = { 2 var message: InboxMessage = null 3 //线程并发检查,如果不允许多线程执行且当前激活线程不为0,直接返回 4 inbox.synchronized { 5 if (!enableConcurrent && numActiveThreads != 0) { 6 return 7 } 8 //获取消息,如果消息不为空,则当前激活线程+1,否则return返回 9 message = messages.poll() 10 if (message != null) { 11 numActiveThreads += 1 12 } else { 13 return 14 } 15 } 16 17 while (true) { 18 //对匹配执行时可能发生的错误,使用Endpoint的onError方法处理 19 safelyCall(endpoint) { 20 //匹配不同类型的消息进行处理 21 message match{...} 22 } 23 24 //对激活进程数量的控制,如果不允许多线程处理且当前激活进程不为1,当前线程退出,numActiveThreads - 1 25 //如果message为空,没有消息需要处理,当前线程退出,numActiveThreads - 1 26 inbox.synchronized { 27 // "enableConcurrent" will be set to false after `onStop` is called, so we should check it 28 // every time. 29 if (!enableConcurrent && numActiveThreads != 1) { 30 // If we are not the only one worker, exit 31 numActiveThreads -= 1 32 return 33 } 34 message = messages.poll() 35 if (message == null) { 36 numActiveThreads -= 1 37 return 38 } 39 } 40 } 41 }

2.3 Dispatcher对RpcEndpoint去注册

1def stop(rpcEndpointRef: RpcEndpointRef): Unit = { 2 synchronized { 3 if (stopped) { 4 // This endpoint will be stopped by Dispatcher.stop() method. 5 return 6 } 7 unregisterRpcEndpoint(rpcEndpointRef.name) 8 } 9 } 10 11private def unregisterRpcEndpoint(name: String): Unit = { 12 //取出EndpointData 13 val data = endpoints.remove(name) 14 if (data != null) { 15 //调用Inbox的stop方法 16 data.inbox.stop() 17 //将EndpointData重新放入receivers队列,让MessageLoop线程能读取到Stop状态,进行相应的处理 18 receivers.offer(data) // for the OnStop message 19 } 20 // Don't clean `endpointRefs` here because it's possible that some messages are being processed 21 // now and they can use `getRpcEndpointRef`. So `endpointRefs` will be cleaned in Inbox via 22 // `removeRpcEndpointRef`. 23 } 24 25/* 26 * 当要移除一个EndpointData时,其Inbox可能正在对消息进行处理,所以调用Inbox的stop方法平滑过渡处理; 27 * 将允许并发运行设置为false,并设置当前Inbox为stopped状态,将当前Inbox所属的EndpointData重新放入receivers, 28 * Inbox.process方法会匹配执行相应的处理,调用Dispatcher.removeRpcEndpointRef方法从endpointRefs缓存中移除当前RpcEndpointRef的映射;  * 在匹配执行OnStop消息的最后,会调用RpcEndpoint的OnStop方法停止RpcEndpoint。 29 */ 30 def stop(): Unit = inbox.synchronized { 31 // The following codes should be in `synchronized` so that we can make sure "OnStop" is the last 32 // message 33 if (!stopped) { 34 // We should disable concurrent here. Then when RpcEndpoint.onStop is called, it's the only 35 // thread that is processing messages. So `RpcEndpoint.onStop` can release its resources 36 // safely. 37 enableConcurrent = false 38 stopped = true 39 messages.add(OnStop) 40 // Note: The concurrent events in messages will be processed one by one. 41 } 42 }

Dispatcher.stop()方法用来停止Dispatcher,之前的stop(rpcEndpointRef:RpcEndpointRef)用于对RpcEndpoint的去注册。

1def stop(): Unit = { 2 synchronized { 3 if (stopped) { 4 return 5 } 6 stopped = true 7 } 8 // Stop all endpoints. This will queue all endpoints for processing by the message loops. 9 //调用unregisterRpcEndpoint方法,对Dispatcher中的所有EndpointData进行去注册,会向endpoints中每个EndpointData中的Inbox中放置 10 //OnStop消息;最后向receivers中投放PoisonPill,即空EndpointData,以停止所有的MessageLoop线程 11 endpoints.keySet().asScala.foreach(unregisterRpcEndpoint) 12 // Enqueue a message that tells the message loops to stop. 13 receivers.offer(PoisonPill) 14 threadpool.shutdown() 15 }

  2.4 Dispatcher提交消息

1/** 2 * 将消息提交给指定的RpcEndpoint 3 * @param endpointName endpoint名称 4 * @param message 消息类型 5 * @param callbackIfStopped endpoint为stop状态时的回调函数 6 */ 7 private def postMessage( 8 endpointName: String, 9 message: InboxMessage, 10 callbackIfStopped: (Exception) => Unit): Unit = { 11 val error = synchronized { 12 //从endpoints缓存获取EndpointData 13 val data = endpoints.get(endpointName) 14 if (stopped) { 15 Some(new RpcEnvStoppedException()) 16 } else if (data == null) { 17 Some(new SparkException(s"Could not find $endpointName.")) 18 } else { 19 //如果endpointData不是停止状态且endpoints缓存中确实有这个EndpointData 20 //调用对应的Inbox.post将消息加入Inbox的消息列表中 21 data.inbox.post(message) 22 //将EndpointData加入receivers队列,以便MessageLoop线程处理此Inbox中的消息 23 receivers.offer(data) 24 None 25 } 26 } 27 // We don't need to call `onStop` in the `synchronized` block 28 error.foreach(callbackIfStopped) 29 } 30 31//在Inbox未停止时,将message加入messages缓存 32def post(message: InboxMessage): Unit = inbox.synchronized { 33 if (stopped) { 34 // We already put "OnStop" into "messages", so we should drop further messages 35 onDrop(message) 36 } else { 37 messages.add(message) 38 false 39 } 40}

3 NettyStreamManager

基于ConcurrentHashMap提供NettyRpcEnv的文件流服务,支持普通文件、jar文件及目录的添加缓存和文件流读取,各个Excutor节点可以使用Driver端提供的NettyStreamManager从Driver端下载jar包或文件支持任务的运行。

4 TransportContext

TransportContext内部包含TransportConf和RpcHandler,封装了用于创建TransportClientFactory和TransportServer的上下文信息;TransportClientFactory是创建TransportClient的工厂类,用于创建RPC框架的客户端,transportServer是RPC框架的服务端。

1private val transportContext = new TransportContext(transportConf, 2 new NettyRpcHandler(dispatcher, this, streamManager))

  创建TransportContext需要两个参数:transportConf和NettyRpcHandler,主要看一下NettyRpcHandler

1private[netty] class NettyRpcHandler( 2 dispatcher: Dispatcher, 3 nettyEnv: NettyRpcEnv, 4 streamManager: StreamManager) extends RpcHandler with Logging { 5 6 // A variable to track the remote RpcEnv addresses of all clients 7 private val remoteAddresses = new ConcurrentHashMap[RpcAddress, RpcAddress]() 8 9 //带回调函数的receive方法,调用internalReceive方法将将ByteBuffer类型的消息转化为RequestMessage 10 //最后调用dispatcher.postRemoteMessage将消息投递到Inbox,由RpcEndpoint处理消息并回复客户端 11 override def receive( 12 client: TransportClient, 13 message: ByteBuffer, 14 callback: RpcResponseCallback): Unit = { 15 val messageToDispatch = internalReceive(client, message) 16 dispatcher.postRemoteMessage(messageToDispatch, callback) 17 } 18 19 //方法重载,RpcEndpoint处理完消息不会回复客户端 20 override def receive( 21 client: TransportClient, 22 message: ByteBuffer): Unit = { 23 val messageToDispatch = internalReceive(client, message) 24 dispatcher.postOneWayMessage(messageToDispatch) 25 } 26 27 //将ByteBuffer类型的消息转化为RequestMessage 28 private def internalReceive(client: TransportClient, message: ByteBuffer): RequestMessage = { 29 //从TransportClient中获取远端地址RpcAddress 30 val addr = client.getChannel().remoteAddress().asInstanceOf[InetSocketAddress] 31 assert(addr != null) 32 val clientAddr = RpcAddress(addr.getHostString, addr.getPort) 33 //封装消息 34 val requestMessage = RequestMessage(nettyEnv, client, message) 35 //如果没有发送者地址信息,使用从TransportClient获取的远端地址RpcAddress、消息的接收者(RpcEndpoint)、消息内容构造新的消息 36 if (requestMessage.senderAddress == null) { 37 // Create a new message with the socket address of the client as the sender. 38 new RequestMessage(clientAddr, requestMessage.receiver, requestMessage.content) 39 } else { 40 // The remote RpcEnv listens to some port, we should also fire a RemoteProcessConnected for 41 // the listening address 42 //获取发送者地址信息,将远端地址RpcAddress和发送者地址信息映射关系放入缓存remoteAddresses 43 val remoteEnvAddress = requestMessage.senderAddress 44 if (remoteAddresses.putIfAbsent(clientAddr, remoteEnvAddress) == null) { 45 //向endpoints缓存中的所有EndpointData的Inbox中放入RemoteProcessConnected类型的消息 46 dispatcher.postToAll(RemoteProcessConnected(remoteEnvAddress)) 47 } 48 requestMessage 49 } 50 } 51 52 ... //其他类型消息的处理,与receive类似 53 54}

5 客户端发送请求

1//用于处理请求超时的调度器 2val timeoutScheduler = ThreadUtils.newDaemonSingleThreadScheduledExecutor("netty-rpc-env-timeout") 3 4//用于异步处理客户端创建的线程池 5private[netty] val clientConnectionExecutor = ThreadUtils.newDaemonCachedThreadPool( 6 "netty-rpc-connection", 7 conf.getInt("spark.rpc.connect.threads", 64)) 8 9/** 10 * 缓存远端RPC地址与OutBox的关系 11 * OutBox与之前的Inbox类似,Outbox是在客户端使用,通过OutboxMessage封装对外发送的消息 12 * Inbox在服务端使用,通过InboxMessage封装接收的消息。 13 * outbox内部有messgaes列表存放消息,通过drainOutbox方法循环取出消息并调用sendWith方法处理 14 * 15 */ 16private val outboxes = new ConcurrentHashMap[RpcAddress, Outbox]()

  篇幅原因到此为止,很多东西还停留在代码层面,有点云里雾里,后面研究其他组件的时候有机会再重读RPC环境的代码吧==!

  请求的发送与接收处理流程

1、通过NettyRpcEndpointRef的send/ask方法向远端节点的RpcEndpoint发送消息,消息会先被封装为OutboxMessage,然后放入远端RpcEndpoint的地址对应的Outbox的messages列表中。

2、Outbox的drainOutbox方法不断从messages列表取出OutboxMessage,并使用内部的TransportClient向远端NettyRpcEnv发送OutboxMessage。

3、发送的请求与在远端RpcEndpoint的TransportServer建立连接,请求先经过RPC管道的处理后由NettyRpcHandler处理,NettyRpcHandler的receive方法会调用Dispatcher的post...方法将消息放入EndpointData内部的Inbox的messges中,最后MessageLoop线程会读取消息并将消息发送给对应的RpcEndpoint处理。

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

swap空间的增减方法

(1)增大swap空间去激活swap交换区:swapoff v /dev/vg00/lvswap扩展交换lv:lvextend L 10G /dev/vg00/lvswap重新生成swap交换区:mkswap /dev/vg00/lvswap激活新生成的交换区:swapon v /dev/vg00/lvswap

Spark2.4.0源码——RpcEnv - HelloWorld