暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Flink源码解析2-JobManager初始化(42)

beenrun 2023-01-14
545

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 logging
    EnvironmentInformation.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.ymal
    Configuration 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 configuration
          configuration.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 DispatcherResourceManagerComponentFactory
          dispatcherResourceManagerComponentFactory =
          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 nothing
          shutDownAsync(
          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 services
            configuration.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!

            一个今天胜似两个明天!

            感谢阅读。期待点赞、分享、关注。

            文章转载自beenrun,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

            评论