JobManager源码分析
JobManager是Flink集群的主节点,协调Flink应用程序执行有关的职责:JobManager决定何时调度下一个 task(或一组 task)、对完成的 task 或执行失败做出反应、协调 checkpoint、并且协调从失败中恢复等等,这个进程由四个不同的组件构成。
(1)ResourceManager:负责 Flink 集群中的资源提供、回收、分配 - 它管理 task slots。
(2)Dispatcher:Dispatcher 提供了一个 REST 接口,用来提交 Flink 应用程序执行,并为每个提交的作业启动一个新的 JobMaster。它还运行 Flink WebUI 用来提供作业执行信息。
(3)JobMaster:JobMaster 负责管理单个JobGraph的执行。Flink 集群中可以同时运行多个作业,每个作业都有自己的 JobMaster。
(4)WebMonitorEndpoint:里面有很多Handler,如果客户端提交flink run的方式类提交job到flink集群,最终是WebMonitorEndpoint接收,并决定使用那个Handler进行处理。
JobManager启动类Main
由上一个文章Flink源码解析1-脚本入口(41)知道入口类为:org.apache.flink.runtime.entrypoint.StandaloneSessionClusterEntrypoint
ClusterEntrypoint.runClusterEntrypoint(entrypoint);运行集群
public static void main(String[] args) {// startup checks and loggingEnvironmentInformation.logEnvironmentInfo(LOG, StandaloneSessionClusterEntrypoint.class.getSimpleName(), args);SignalHandler.register(LOG);JvmShutdownSafeguard.installAsShutdownHook(LOG);final EntrypointClusterConfiguration entrypointClusterConfiguration =ClusterEntrypointUtils.parseParametersOrExit(args,new EntrypointClusterConfigurationParserFactory(),StandaloneSessionClusterEntrypoint.class);//解析Flink的配置文件,flink-conf.ymalConfiguration configuration = loadConfiguration(entrypointClusterConfiguration);//创建一个对象StandaloneSessionClusterEntrypoint entrypoint =new StandaloneSessionClusterEntrypoint(configuration);//运行集群ClusterEntrypoint.runClusterEntrypoint(entrypoint);}
JobManager启动类
进行启动clusterEntrypoint.startCluster();
public static void runClusterEntrypoint(ClusterEntrypoint clusterEntrypoint) {final String clusterEntrypointName = clusterEntrypoint.getClass().getSimpleName();try {//调用传入参数的启动集群方法startCluster()clusterEntrypoint.startCluster();} catch (ClusterEntrypointException e) {LOG.error(String.format("Could not start cluster entrypoint %s.", clusterEntrypointName),e);System.exit(STARTUP_FAILURE_RETURN_CODE);}int returnCode;Throwable throwable = null;try {returnCode = clusterEntrypoint.getTerminationFuture().get().processExitCode();} catch (Throwable e) {throwable = ExceptionUtils.stripExecutionException(e);returnCode = RUNTIME_FAILURE_RETURN_CODE;}LOG.info("Terminating cluster entrypoint process {} with exit code {}.",clusterEntrypointName,returnCode,throwable);System.exit(returnCode);}
初始化文件并启动集群
(1)启动插件管理
(2)初始化文件系统
(3)启动集群
public void startCluster() throws ClusterEntrypointException {LOG.info("Starting {}.", getClass().getSimpleName());try {FlinkSecurityManager.setFromConfiguration(configuration);//启动插件管理器PluginManager pluginManager =PluginUtils.createPluginManagerFromRootFolder(configuration);//初始化文件系统/*** 本地Local 客户端会用 JobGraph ==>JobGraphFile* HDFS:FileSystem* 封装对象:Flink中用的HadoopFileSystem,包装了HFDS的FileSystem对象*/configureFileSystems(configuration, pluginManager);SecurityContext securityContext = installSecurityContext(configuration);ClusterEntrypointUtils.configureUncaughtExceptionHandler(configuration);//启动集群securityContext.runSecured((Callable<Void>)() -> {runCluster(configuration, pluginManager);return null;});}}
创建需要的工厂
首先服务初始化
然后创建需要的各种工厂在DefaultDispatcherResourceManagerComponentFactory类中实现
(1)创建dispatcher的Factory
(2) 创建ResourceManager的Factory
(3) 创建RestEndpointFactory(WebMonitorEndpoint的Factory)
private void runCluster(Configuration configuration, PluginManager pluginManager)throws Exception {synchronized (lock) {//服务的初始化initializeServices(configuration, pluginManager);// write host information into configurationconfiguration.setString(JobManagerOptions.ADDRESS, commonRpcService.getAddress());configuration.setInteger(JobManagerOptions.PORT, commonRpcService.getPort());//创建dispatcher的Factory//创建ResourceManager的Factory//创建RestEndpointFactory(WebMonitorEndpoint的Factory)/*** this.dispatcherRunnerFactory = dispatcherRunnerFactory;* this.resourceManagerFactory = resourceManagerFactory;* this.restEndpointFactory = restEndpointFactory;*/final DispatcherResourceManagerComponentFactorydispatcherResourceManagerComponentFactory =createDispatcherResourceManagerComponentFactory(configuration);/*** 启动关键组件 Dispatcher 和 ResourceManager 和WebMonitorEndpoint*/clusterComponent =dispatcherResourceManagerComponentFactory.create(configuration,resourceId.unwrap(),ioExecutor,commonRpcService,haServices,blobServer,heartbeatServices,delegationTokenManager,metricRegistry,executionGraphInfoStore,new RpcMetricQueryServiceRetriever(metricRegistry.getMetricQueryServiceRpcService()),this);//集群进行关闭的时候,回调这里clusterComponent.getShutDownFuture().whenComplete((ApplicationStatus applicationStatus, Throwable throwable) -> {if (throwable != null) {shutDownAsync(ApplicationStatus.UNKNOWN,ShutdownBehaviour.GRACEFUL_SHUTDOWN,ExceptionUtils.stringifyException(throwable),false);} else {// This is the general shutdown path. If a separate more// specific shutdown was// already triggered, this will do nothingshutDownAsync(applicationStatus,ShutdownBehaviour.GRACEFUL_SHUTDOWN,null,true);}});}}
初始化服务
执行 initializeServices(configuration, pluginManager);
protected void initializeServices(Configuration configuration, PluginManager pluginManager)throws Exception {LOG.info("Initializing cluster services.");synchronized (lock) {resourceId =configuration.getOptional(JobManagerOptions.JOB_MANAGER_RESOURCE_ID).map(value ->DeterminismEnvelope.deterministicValue(new ResourceID(value))).orElseGet(() ->DeterminismEnvelope.nondeterministicValue(ResourceID.generate()));LOG.debug("Initialize cluster entrypoint {} with resource id {}.",getClass().getSimpleName(),resourceId);workingDirectory =ClusterEntrypointUtils.createJobManagerWorkingDirectory(configuration, resourceId);LOG.info("Using working directory: {}.", workingDirectory);rpcSystem = RpcSystem.load(configuration);//第一步// 创建Akka RPC服务 commonRpcService 基于Akka的RpcService实现//commonRpcService 是一个基于Akka的actorSystem 其实就是一个tcp的RPC服务,端口为6123//1.初始化ActorSystem//2.启动Actor//启动好了后,会自己给自己发送一个消息,就是利用这里进行发送消息的commonRpcService =RpcUtils.createRemoteRpcService(rpcSystem,configuration,configuration.getString(JobManagerOptions.ADDRESS),getRPCPortRange(configuration),configuration.getString(JobManagerOptions.BIND_HOST),configuration.getOptional(JobManagerOptions.RPC_BIND_PORT));JMXService.startInstance(configuration.getString(JMXServerOptions.JMX_SERVER_PORT));// update the configuration used to create the high availability servicesconfiguration.setString(JobManagerOptions.ADDRESS, commonRpcService.getAddress());configuration.setInteger(JobManagerOptions.PORT, commonRpcService.getPort());/*** 第二步:初始化一个用来做线程的线程池,服务器CPU的数量的4倍,如果当前CPU为8那么线程数量为4*8=32,专门用来做IO*/ioExecutor =Executors.newFixedThreadPool(ClusterEntrypointUtils.getPoolSize(configuration),new ExecutorThreadFactory("cluster-io"));/*** 第三步 HA的service* 如:ResourceManager的leader选举,JobManager leader的选举等* haServices==ZookeeperHAService*/haServices = createHaServices(configuration, ioExecutor, rpcSystem);/*** 第四步骤* 初始化blobServer,大文件的上传等,比如:jar包* 为了完成大文件的传输*/blobServer =BlobUtils.createBlobServer(configuration,Reference.borrowed(workingDirectory.unwrap().getBlobStorageDirectory()),haServices.createBlobStore());blobServer.start();configuration.setString(BlobServerOptions.PORT, String.valueOf(blobServer.getPort()));/*** 第五步* 初始化一个心跳服务* 在主节点中有很多心跳服务,所有的心跳服务都是由heartbeatServices来进行提供的,* 谁需要就由heartbeatServices创建一个实例HeartBeatImpl,用来实现心跳*/heartbeatServices = createHeartbeatServices(configuration);delegationTokenManager =KerberosDelegationTokenManagerFactory.create(getClass().getClassLoader(),configuration,commonRpcService.getScheduledExecutor(),ioExecutor);/*** 第六步,性能监控服务* 跟踪已经注册的Metrics*/metricRegistry = createMetricRegistry(configuration, pluginManager, rpcSystem);final RpcService metricQueryServiceRpcService =MetricUtils.startRemoteMetricsRpcService(configuration,commonRpcService.getAddress(),configuration.getString(JobManagerOptions.BIND_HOST),rpcSystem);metricRegistry.startQueryService(metricQueryServiceRpcService, null);final String hostname = RpcUtils.getHostname(commonRpcService);processMetricGroup =MetricUtils.instantiateProcessMetricGroup(metricRegistry,hostname,ConfigurationUtils.getSystemResourceMetricsProbingInterval(configuration));/*** 第七步executionGraphInfoStore:存储execution graph的服务,默认有两种实现* 1.FileExecutionGraphInfoStore:会持久化到文件系统,也会在内存中缓存* 2.MemoryExecutionGraphInfoStore:在内存中进行缓存* 默认实现是基于文件的*/executionGraphInfoStore =createSerializableExecutionGraphStore(configuration, commonRpcService.getScheduledExecutor());}}
initializeServices初始化流程
JobManager初始化总结

One today is worth two tomorrows!
一个今天胜似两个明天!
感谢阅读。期待点赞、分享、关注。




