• 较真儿学源码系列-PowerJob启动流程源码分析


            PowerJob版本:4.3.2-main。


    1 简介

            PowerJob是全新一代的分布式任务调度与计算框架,官网地址:http://www.powerjob.tech/。其中介绍了PowerJob的功能特点,以及与其他调度框架的对比,这里就不再赘述了。

            以上是PowerJob的架构图,取自官网。可以看出,PowerJob是典型的客户端/服务端交互的架构(但是在PowerJob中却没有一般分布式中间件会有的注册中心)。本文就从启动流程出发,来一起探究下PowerJob在启动阶段中都做了些什么动作。


    2 服务端

            既然要看启动流程源码,那么首先就来看下启动类。PowerJob依赖于Spring Boot,启动类为PowerJobServerApplication:

    1. /**
    2. * powerjob-server entry
    3. *
    4. * @author tjq
    5. * @since 2020/3/29
    6. */
    7. @Slf4j
    8. @EnableScheduling
    9. @SpringBootApplication
    10. public class PowerJobServerApplication {
    11. private static final String TIPS = "\n\n" +
    12. "******************* PowerJob Tips *******************\n" +
    13. "如果应用无法启动,我们建议您仔细阅读以下文档来解决:\n" +
    14. "if server can't startup, we recommend that you read the documentation to find a solution:\n" +
    15. "https://www.yuque.com/powerjob/guidence/problem\n" +
    16. "******************* PowerJob Tips *******************\n\n";
    17. public static void main(String[] args) {
    18. pre();
    19. // Start SpringBoot application.
    20. try {
    21. SpringApplication.run(PowerJobServerApplication.class, args);
    22. } catch (Throwable t) {
    23. log.error(TIPS);
    24. throw t;
    25. }
    26. }
    27. private static void pre() {
    28. log.info(TIPS);
    29. PropertyUtils.init();
    30. }
    31. }
    32. /**
    33. * 加载配置文件
    34. *
    35. * @author tjq
    36. * @since 2020/5/18
    37. */
    38. @Slf4j
    39. public class PropertyUtils {
    40. private static final Properties PROPERTIES = new Properties();
    41. public static Properties getProperties() {
    42. return PROPERTIES;
    43. }
    44. public static void init() {
    45. URL propertiesURL = PropertyUtils.class.getClassLoader().getResource("application.properties");
    46. Objects.requireNonNull(propertiesURL);
    47. try (InputStream is = propertiesURL.openStream()) {
    48. PROPERTIES.load(is);
    49. } catch (Exception e) {
    50. ExceptionUtils.rethrow(e);
    51. }
    52. }
    53. }

            其中pre方法只是将配置文件内容加载进Properties缓存中。

            既然启动类看不出什么逻辑,那么接下来就来看下服务端启动时有没有什么初始化的逻辑:

    2.1 PowerTransportService

            PowerTransportService类是用作数据传输服务的,也就是客户端和服务端之间的通信。其实现了InitializingBean接口,所以查看下其afterPropertiesSet方法的实现:

    1. /**
    2. * PowerTransportService:
    3. */
    4. @Override
    5. public void afterPropertiesSet() throws Exception {
    6. log.info("[PowerTransportService] start to initialize whole PowerTransportService!");
    7. log.info("[PowerTransportService] activeProtocols: {}", activeProtocols);
    8. if (StringUtils.isEmpty(activeProtocols)) {
    9. throw new IllegalArgumentException("activeProtocols can't be empty!");
    10. }
    11. for (String protocol : activeProtocols.split(OmsConstant.COMMA)) {
    12. try {
    13. final int port = parseProtocolPort(protocol);
    14. //初始化网络通讯
    15. initRemoteFrameWork(protocol, port);
    16. } catch (Throwable t) {
    17. log.error("[PowerTransportService] initialize protocol[{}] failed. If you don't need to use this protocol, you can turn it off by 'oms.transporter.active.protocols'", protocol);
    18. ExceptionUtils.rethrow(t);
    19. }
    20. }
    21. //选择默认的通信协议,默认为HTTP
    22. choseDefault();
    23. log.info("[PowerTransportService] initialize successfully!");
    24. log.info("[PowerTransportService] ALL_PROTOCOLS: {}", protocolName2Info);
    25. }
    26. /**
    27. * 第18行代码处:
    28. */
    29. private void initRemoteFrameWork(String protocol, int port) {
    30. // 从构造器注入改为从 applicationContext 获取来避免循环依赖
    31. //获取所有注解了@Actor的bean
    32. final Map beansWithAnnotation = applicationContext.getBeansWithAnnotation(Actor.class);
    33. log.info("[PowerTransportService] find Actor num={},names={}", beansWithAnnotation.size(), beansWithAnnotation.keySet());
    34. Address address = new Address()
    35. .setHost(NetUtils.getLocalHost())
    36. .setPort(port);
    37. EngineConfig engineConfig = new EngineConfig()
    38. .setServerType(ServerType.SERVER)
    39. .setType(protocol.toUpperCase())
    40. .setBindAddress(address)
    41. .setActorList(Lists.newArrayList(beansWithAnnotation.values()));
    42. log.info("[PowerTransportService] start to initialize RemoteEngine[type={},address={}]", protocol, address);
    43. RemoteEngine re = new PowerJobRemoteEngine();
    44. //初始化网络层
    45. final EngineOutput engineOutput = re.start(engineConfig);
    46. log.info("[PowerTransportService] start RemoteEngine[type={},address={}] successfully", protocol, address);
    47. //放入相关缓存中
    48. this.engines.add(re);
    49. this.protocolName2Info.put(protocol, new ProtocolInfo(protocol, address.toFullAddress(), engineOutput.getTransporter()));
    50. }

            其中需要说明的一点是:PowerJob网络层使用的协议是Akka或Vert.x:

    • Akka是一个在JVM上构建高并发、分布式和弹性消息驱动的应用程序。其是用Scala写的,使用到了Actor模型;
    • 而Vert.x是一个在JVM上构建响应式应用程序的工具包,其底层基于Netty(之前我对Netty的源码也进行过分析,感兴趣的话可以查看《较真儿学源码系列-Netty核心流程源码分析》)。

            PowerJob屏蔽了底层的实现,用两个自定义的注解@Actor和@Handler进行了统一的封装。客户端传过来的请求会自动跳转到@Actor注解的类、@Handler注解的方法上。拿客户端发送心跳给服务端的逻辑为例,流程如下所示:

            下面就来继续看下,在上面的第53行代码处。PowerJob是如何完成这个绑定的:

    1. /**
    2. * PowerJobRemoteEngine:
    3. */
    4. @Override
    5. public EngineOutput start(EngineConfig engineConfig) {
    6. final String engineType = engineConfig.getType();
    7. EngineOutput engineOutput = new EngineOutput();
    8. log.info("[PowerJobRemoteEngine] [{}] start remote engine with config: {}", engineType, engineConfig);
    9. //获取所有@Actor的类,和其中注解了@Handler的方法
    10. List actorInfos = ActorFactory.load(engineConfig.getActorList());
    11. //遍历获取指定协议的CSInitializer
    12. csInitializer = CSInitializerFactory.build(engineType);
    13. String type = csInitializer.type();
    14. Stopwatch sw = Stopwatch.createStarted();
    15. log.info("[PowerJobRemoteEngine] [{}] try to startup CSInitializer[type={}]", engineType, type);
    16. //CsInitializer初始化,这里以Vert.x为例,查看下其实现
    17. csInitializer.init(new CSInitializerConfig()
    18. .setBindAddress(engineConfig.getBindAddress())
    19. .setServerType(engineConfig.getServerType())
    20. );
    21. // 构建通讯器
    22. Transporter transporter = csInitializer.buildTransporter();
    23. engineOutput.setTransporter(transporter);
    24. log.info("[PowerJobRemoteEngine] [{}] start to bind Handler", engineType);
    25. actorInfos.forEach(actor -> actor.getHandlerInfos().forEach(handlerInfo -> log.info("[PowerJobRemoteEngine] [{}] PATH={}, handler={}", engineType, handlerInfo.getLocation().toPath(), handlerInfo.getMethod())));
    26. // 绑定 handler
    27. csInitializer.bindHandlers(actorInfos);
    28. log.info("[PowerJobRemoteEngine] [{}] startup successfully, cost: {}", engineType, sw);
    29. return engineOutput;
    30. }
    31. /**
    32. * ActorFactory:
    33. * 第11行代码处:
    34. */
    35. static List load(List actorList) {
    36. List actorInfos = Lists.newArrayList();
    37. actorList.forEach(actor -> {
    38. final Class clz = actor.getClass();
    39. try {
    40. final Actor anno = clz.getAnnotation(Actor.class);
    41. ActorInfo actorInfo = new ActorInfo().setActor(actor).setAnno(anno);
    42. //获取所有注解了@Handler的方法,并缓存进handlerInfos里
    43. actorInfo.setHandlerInfos(loadHandlerInfos4Actor(actorInfo));
    44. actorInfos.add(actorInfo);
    45. } catch (Throwable t) {
    46. log.error("[ActorFactory] process Actor[{}] failed!", clz);
    47. ExceptionUtils.rethrow(t);
    48. }
    49. });
    50. return actorInfos;
    51. }
    52. /**
    53. * CSInitializerFactory:
    54. * 第13行代码处:
    55. */
    56. static CSInitializer build(String targetType) {
    57. Reflections reflections = new Reflections(OmsConstant.PACKAGE);
    58. //使用Reflections反射获取CSInitializer的实现类,即AkkaCSInitializer和HttpVertxCSInitializer。Reflections的介绍和简单使用请看https://blog.csdn.net/weixin_30342639/article/details/124521467
    59. Setextends CSInitializer>> cSInitializerClzSet = reflections.getSubTypesOf(CSInitializer.class);
    60. log.info("[CSInitializerFactory] scan subTypeOf CSInitializer: {}", cSInitializerClzSet);
    61. for (Classextends CSInitializer> clz : cSInitializerClzSet) {
    62. try {
    63. CSInitializer csInitializer = clz.getDeclaredConstructor().newInstance();
    64. //获取类型:AKKA/HTTP
    65. String type = csInitializer.type();
    66. log.info("[CSInitializerFactory] new instance for CSInitializer[{}] successfully, type={}, object: {}", clz, type, csInitializer);
    67. //遍历获取指定协议的CSInitializer
    68. if (targetType.equalsIgnoreCase(type)) {
    69. return csInitializer;
    70. }
    71. } catch (Exception e) {
    72. log.error("[CSInitializerFactory] new instance for CSInitializer[{}] failed, maybe you should provide a non-parameter constructor", clz);
    73. ExceptionUtils.rethrow(e);
    74. }
    75. }
    76. throw new PowerJobException(String.format("can't load CSInitializer[%s], ensure your package name start with 'tech.powerjob' and import the dependencies!", targetType));
    77. }
    78. /**
    79. * HttpVertxCSInitializer:
    80. * 第21行代码处:
    81. * 这里也就是在做Vert.x的初始化工作,不再继续深入了
    82. */
    83. @Override
    84. public void init(CSInitializerConfig config) {
    85. this.config = config;
    86. vertx = VertxInitializer.buildVertx();
    87. httpServer = VertxInitializer.buildHttpServer(vertx);
    88. httpClient = VertxInitializer.buildHttpClient(vertx);
    89. }
    90. /**
    91. * 第34行代码处:
    92. */
    93. @Override
    94. @SneakyThrows
    95. public void bindHandlers(List actorInfos) {
    96. Router router = Router.router(vertx);
    97. // 处理请求响应
    98. router.route().handler(BodyHandler.create());
    99. actorInfos.forEach(actorInfo -> {
    100. Optional.ofNullable(actorInfo.getHandlerInfos()).orElse(Collections.emptyList()).forEach(handlerInfo -> {
    101. //获取目标地址,即@Actor和@Handler注解上path的拼接
    102. String handlerHttpPath = handlerInfo.getLocation().toPath();
    103. ProcessType processType = handlerInfo.getAnno().processType();
    104. Handler routingContextHandler = buildRequestHandler(actorInfo, handlerInfo);
    105. //添加路由绑定
    106. Route route = router.post(handlerHttpPath);
    107. if (processType == ProcessType.BLOCKING) {
    108. //绑定阻塞调用handler
    109. route.blockingHandler(routingContextHandler, false);
    110. } else {
    111. //绑定非阻塞调用handler
    112. route.handler(routingContextHandler);
    113. }
    114. });
    115. });
    116. // 启动 vertx http server
    117. final int port = config.getBindAddress().getPort();
    118. final String host = config.getBindAddress().getHost();
    119. httpServer.requestHandler(router)
    120. .exceptionHandler(e -> log.error("[PowerJob] unknown exception in Actor communication!", e))
    121. .listen(port, host)
    122. .toCompletionStage()
    123. .toCompletableFuture()
    124. .get(1, TimeUnit.MINUTES);
    125. log.info("[PowerJobRemoteEngine] startup vertx HttpServer successfully!");
    126. }
    127. /**
    128. * 第127行代码处:
    129. */
    130. private Handler buildRequestHandler(ActorInfo actorInfo, HandlerInfo handlerInfo) {
    131. Method method = handlerInfo.getMethod();
    132. Optional> powerSerializeClz = RemoteUtils.findPowerSerialize(method.getParameterTypes());
    133. // 内部框架,严格模式,绑定失败直接报错
    134. if (!powerSerializeClz.isPresent()) {
    135. throw new PowerJobException("can't find any 'PowerSerialize' object in handler args: " + handlerInfo.getLocation());
    136. }
    137. //这里实际上是注册了一个事件驱动的回调函数(Netty的玩法),当有请求过来的时候,会走到下面的代码里
    138. return ctx -> {
    139. final RequestBody body = ctx.body();
    140. final Object convertResult = body.asPojo(powerSerializeClz.get());
    141. try {
    142. //这里通过反射调用相关的@Handler注解的方法
    143. Object response = method.invoke(actorInfo.getActor(), convertResult);
    144. if (response != null) {
    145. if (response instanceof String) {
    146. ctx.end((String) response);
    147. } else {
    148. ctx.json(JsonObject.mapFrom(response));
    149. }
    150. return;
    151. }
    152. ctx.end();
    153. } catch (Throwable t) {
    154. // 注意这里是框架实际运行时,日志输出用标准 PowerJob 格式
    155. log.error("[PowerJob] invoke Handler[{}] failed!", handlerInfo.getLocation(), t);
    156. ctx.fail(HttpResponseStatus.INTERNAL_SERVER_ERROR.code(), t);
    157. }
    158. };
    159. }
    160.         通过上面的buildRequestHandler方法可知,当Vert.x接收到客户端的请求时,会调用到一个回调函数。而这个回调函数会最终通过反射的方式调用到自定义的@Handler注解的方法中来。

              上面只是完成了初始化和绑定的操作,还缺少具体调用时候的逻辑(这里属于一个整体流程,所以就一起分析了)。继续拿发送心跳为例,调用的方法是:TransportUtils.reportWorkerHeartbeat:

      1. /**
      2. * TransportUtils:
      3. */
      4. public static void reportWorkerHeartbeat(WorkerHeartbeat req, String address, Transporter transporter) {
      5. //绑定url调用信息
      6. final URL url = easyBuildUrl(ServerType.SERVER, S4W_PATH, S4W_HANDLER_WORKER_HEARTBEAT, address);
      7. //这里依旧是拿Vert.x的实现来分析
      8. transporter.tell(url, req);
      9. }
      10. /**
      11. * VertxTransporter:
      12. * 第8行代码处:
      13. */
      14. @Override
      15. public void tell(URL url, PowerSerializable request) {
      16. post(url, request, null);
      17. }
      18. /**
      19. * 这里就涉及到Vert.x的具体调用细节了,不再深入
      20. */
      21. @SuppressWarnings("unchecked")
      22. private CompletionStage post(URL url, PowerSerializable request, Class clz) {
      23. final String host = url.getAddress().getHost();
      24. final int port = url.getAddress().getPort();
      25. final String path = url.getLocation().toPath();
      26. RequestOptions requestOptions = new RequestOptions()
      27. .setMethod(HttpMethod.POST)
      28. .setHost(host)
      29. .setPort(port)
      30. .setURI(path);
      31. // 获取远程服务器的HTTP连接
      32. Future httpClientRequestFuture = httpClient.request(requestOptions);
      33. // 转换 -> 发送请求获取响应
      34. Future responseFuture = httpClientRequestFuture.compose(httpClientRequest ->
      35. httpClientRequest
      36. .putHeader(HttpHeaderNames.CONTENT_TYPE, HttpHeaderValues.APPLICATION_JSON)
      37. .send(JsonObject.mapFrom(request).toBuffer())
      38. );
      39. return responseFuture.compose(httpClientResponse -> {
      40. // throw exception
      41. final int statusCode = httpClientResponse.statusCode();
      42. if (statusCode != HttpResponseStatus.OK.code()) {
      43. // CompletableFuture.get() 时会传递抛出该异常
      44. throw new RemotingException(String.format("request [host:%s,port:%s,url:%s] failed, status: %d, msg: %s",
      45. host, port, path, statusCode, httpClientResponse.statusMessage()
      46. ));
      47. }
      48. return httpClientResponse.body().compose(x -> {
      49. if (clz == null) {
      50. return Future.succeededFuture(null);
      51. }
      52. if (clz.equals(String.class)) {
      53. return Future.succeededFuture((T) x.toString());
      54. }
      55. return Future.succeededFuture(x.toJsonObject().mapTo(clz));
      56. });
      57. }).toCompletionStage();
      58. }

              当客户端往服务端发送完请求后,服务端接受到相关的请求,会调用到相应的回调函数,从而最终调用到@Handler注解的方法。整个流程就串起来了。

      2.2. InstanceMetadataService

              同PowerTransportService类一样,InstanceMetadataService也实现了InitializingBean接口,所以查看下其afterPropertiesSet方法的实现:

      1. /**
      2. * InstanceMetadataService:
      3. */
      4. @Override
      5. public void afterPropertiesSet() throws Exception {
      6. instanceId2JobInfoCache = CacheBuilder.newBuilder()
      7. .concurrencyLevel(CACHE_CONCURRENCY_LEVEL)
      8. .maximumSize(instanceMetadataCacheSize)
      9. .softValues()
      10. .build();
      11. }

              其中只是初始化了一个本地缓存,没有多余的逻辑。

      2.3 CoreScheduleTaskManager

              CoreScheduleTaskManager也实现了InitializingBean接口,查看下其afterPropertiesSet方法的实现:

      1. /**
      2. * CoreScheduleTaskManager:
      3. */
      4. @SuppressWarnings("AlibabaAvoidManuallyCreateThread")
      5. @Override
      6. public void afterPropertiesSet() {
      7. // 定时调度
      8. coreThreadContainer.add(new Thread(new LoopRunnable("ScheduleCronJob", PowerScheduleService.SCHEDULE_RATE, () -> powerScheduleService.scheduleNormalJob(TimeExpressionType.CRON)), "Thread-ScheduleCronJob"));
      9. coreThreadContainer.add(new Thread(new LoopRunnable("ScheduleDailyTimeIntervalJob", PowerScheduleService.SCHEDULE_RATE, () -> powerScheduleService.scheduleNormalJob(TimeExpressionType.DAILY_TIME_INTERVAL)), "Thread-ScheduleDailyTimeIntervalJob"));
      10. coreThreadContainer.add(new Thread(new LoopRunnable("ScheduleCronWorkflow", PowerScheduleService.SCHEDULE_RATE, powerScheduleService::scheduleCronWorkflow), "Thread-ScheduleCronWorkflow"));
      11. coreThreadContainer.add(new Thread(new LoopRunnable("ScheduleFrequentJob", PowerScheduleService.SCHEDULE_RATE, powerScheduleService::scheduleFrequentJob), "Thread-ScheduleFrequentJob"));
      12. // 数据清理
      13. coreThreadContainer.add(new Thread(new LoopRunnable("CleanWorkerData", PowerScheduleService.SCHEDULE_RATE, powerScheduleService::cleanData), "Thread-CleanWorkerData"));
      14. // 状态检查
      15. coreThreadContainer.add(new Thread(new LoopRunnable("CheckRunningInstance", InstanceStatusCheckService.CHECK_INTERVAL, instanceStatusCheckService::checkRunningInstance), "Thread-CheckRunningInstance"));
      16. coreThreadContainer.add(new Thread(new LoopRunnable("CheckWaitingDispatchInstance", InstanceStatusCheckService.CHECK_INTERVAL, instanceStatusCheckService::checkWaitingDispatchInstance), "Thread-CheckWaitingDispatchInstance"));
      17. coreThreadContainer.add(new Thread(new LoopRunnable("CheckWaitingWorkerReceiveInstance", InstanceStatusCheckService.CHECK_INTERVAL, instanceStatusCheckService::checkWaitingWorkerReceiveInstance), "Thread-CheckWaitingWorkerReceiveInstance"));
      18. coreThreadContainer.add(new Thread(new LoopRunnable("CheckWorkflowInstance", InstanceStatusCheckService.CHECK_INTERVAL, instanceStatusCheckService::checkWorkflowInstance), "Thread-CheckWorkflowInstance"));
      19. coreThreadContainer.forEach(Thread::start);
      20. }

              可以看到,其中初始化了一堆的定时任务。从中挑几个定时任务来分析下:

      2.3.1 ScheduleCronJob/ScheduleDailyTimeIntervalJob

      1. /**
      2. * PowerScheduleService:
      3. */
      4. public void scheduleNormalJob(TimeExpressionType timeExpressionType) {
      5. long start = System.currentTimeMillis();
      6. // 调度 CRON 表达式 JOB
      7. try {
      8. //获取在PowerJob控制台配置的appId,也就是服务id(PowerJob会使用数据库来存储服务、实例等相关数据。在进行一些操作前,会先落库,然后再执行。以此避免相关操作丢失的情况出现)
      9. final List allAppIds = appInfoRepository.listAppIdByCurrentServer(transportService.defaultProtocol().getAddress());
      10. if (CollectionUtils.isEmpty(allAppIds)) {
      11. log.info("[NormalScheduler] current server has no app's job to schedule.");
      12. return;
      13. }
      14. scheduleNormalJob0(timeExpressionType, allAppIds);
      15. } catch (Exception e) {
      16. log.error("[NormalScheduler] schedule cron job failed.", e);
      17. }
      18. long cost = System.currentTimeMillis() - start;
      19. log.info("[NormalScheduler] {} job schedule use {} ms.", timeExpressionType, cost);
      20. if (cost > SCHEDULE_RATE) {
      21. log.warn("[NormalScheduler] The database query is using too much time({}ms), please check if the database load is too high!", cost);
      22. }
      23. }
      24. /**
      25. * 第14行代码处:
      26. * 调度普通服务端计算表达式类型(CRON、DAILY_TIME_INTERVAL)的任务
      27. *
      28. * @param timeExpressionType 表达式类型
      29. * @param appIds appIds
      30. */
      31. private void scheduleNormalJob0(TimeExpressionType timeExpressionType, List appIds) {
      32. long nowTime = System.currentTimeMillis();
      33. long timeThreshold = nowTime + 2 * SCHEDULE_RATE;
      34. //分组执行
      35. Lists.partition(appIds, MAX_APP_NUM).forEach(partAppIds -> {
      36. try {
      37. // 查询条件:任务开启 + 使用CRON表达调度时间 + 指定appId + 即将需要调度执行
      38. //获取在PowerJob控制台配置的执行任务
      39. List jobInfos = jobInfoRepository.findByAppIdInAndStatusAndTimeExpressionTypeAndNextTriggerTimeLessThanEqual(partAppIds, SwitchableStatus.ENABLE.getV(), timeExpressionType.getV(), timeThreshold);
      40. if (CollectionUtils.isEmpty(jobInfos)) {
      41. return;
      42. }
      43. // 1. 批量写日志表
      44. Map jobId2InstanceId = Maps.newHashMap();
      45. log.info("[NormalScheduler] These {} jobs will be scheduled: {}.", timeExpressionType.name(), jobInfos);
      46. jobInfos.forEach(jobInfo -> {
      47. //实例表进行落库(任务和实例的关系为:每执行一次任务会生成一个实例)
      48. Long instanceId = instanceService.create(jobInfo.getId(), jobInfo.getAppId(), jobInfo.getJobParams(), null, null, jobInfo.getNextTriggerTime()).getInstanceId();
      49. jobId2InstanceId.put(jobInfo.getId(), instanceId);
      50. });
      51. instanceInfoRepository.flush();
      52. // 2. 推入时间轮中等待调度执行(对时间轮代码的分析请查看https://blog.csdn.net/weixin_30342639/article/details/132732836)
      53. jobInfos.forEach(jobInfoDO -> {
      54. Long instanceId = jobId2InstanceId.get(jobInfoDO.getId());
      55. //获取下次执行时间
      56. long targetTriggerTime = jobInfoDO.getNextTriggerTime();
      57. long delay = 0;
      58. if (targetTriggerTime < nowTime) {
      59. log.warn("[Job-{}] schedule delay, expect: {}, current: {}", jobInfoDO.getId(), targetTriggerTime, System.currentTimeMillis());
      60. } else {
      61. //计算距离下次执行时间所需要延迟的时间
      62. delay = targetTriggerTime - nowTime;
      63. }
      64. //任务实例推入时间轮,等待被执行
      65. InstanceTimeWheelService.schedule(instanceId, delay, () -> dispatchService.dispatch(jobInfoDO, instanceId, Optional.empty(), Optional.empty()));
      66. });
      67. // 3. 计算下一次调度时间(忽略5S内的重复执行,即CRON模式下最小的连续执行间隔为 SCHEDULE_RATE ms)
      68. jobInfos.forEach(jobInfoDO -> {
      69. try {
      70. //重新计算下次执行时间,并落库
      71. //这里会用到策略模式,不同类型的任务(Cron/固定频率/每日固定间隔)会有不同的计算方式,具体不再深入分析,感兴趣可自行查看
      72. refreshJob(timeExpressionType, jobInfoDO);
      73. } catch (Exception e) {
      74. log.error("[Job-{}] refresh job failed.", jobInfoDO.getId(), e);
      75. }
      76. });
      77. jobInfoRepository.flush();
      78. } catch (Exception e) {
      79. log.error("[NormalScheduler] schedule {} job failed.", timeExpressionType.name(), e);
      80. }
      81. });
      82. }

              在上面第76行代码处,当任务推入到时间轮之后(对时间轮代码的分析请查看《较真儿学源码系列-PowerJob时间轮源码分析》),等到需要被执行的时候,会调用到DispatchService.dispatch方法来派发任务:

      1. /**
      2. * DispatchService:
      3. * 将任务从Server派发到Worker(TaskTracker)
      4. * 只会派发当前状态为等待派发的任务实例
      5. * **************************************************
      6. * 2021-02-03 modify by Echo009
      7. * 1、移除参数 当前运行次数、工作流实例ID、实例参数
      8. * 更改为从当前任务实例中获取获取以上信息
      9. * 2、移除运行次数相关的(runningTimes)处理逻辑
      10. * 迁移至 {@link InstanceManager#updateStatus} 中处理
      11. * **************************************************
      12. *
      13. * @param jobInfo 任务的元信息
      14. * @param instanceId 任务实例ID
      15. * @param instanceInfoOptional 任务实例信息,可选
      16. * @param overloadOptional 超载信息,可选
      17. */
      18. @UseCacheLock(type = "processJobInstance", key = "#jobInfo.getMaxInstanceNum() > 0 || T(tech.powerjob.common.enums.TimeExpressionType).FREQUENT_TYPES.contains(#jobInfo.getTimeExpressionType()) ? #jobInfo.getId() : #instanceId", concurrencyLevel = 1024)
      19. public void dispatch(JobInfoDO jobInfo, Long instanceId, Optional instanceInfoOptional, Optional> overloadOptional) {
      20. // 允许从外部传入实例信息,减少 io 次数
      21. // 检查当前任务是否被取消
      22. //获取实例
      23. InstanceInfoDO instanceInfo = instanceInfoOptional.orElseGet(() -> instanceInfoRepository.findByInstanceId(instanceId));
      24. Long jobId = instanceInfo.getJobId();
      25. //...
      26. // 任务信息已经被删除
      27. if (jobInfo.getId() == null) {
      28. log.warn("[Dispatcher-{}|{}] cancel dispatch due to job(id={}) has been deleted!", jobId, instanceId, jobId);
      29. instanceManager.processFinishedInstance(instanceId, instanceInfo.getWfInstanceId(), FAILED, "can't find job by id " + jobId);
      30. return;
      31. }
      32. Date now = new Date();
      33. String dbInstanceParams = instanceInfo.getInstanceParams() == null ? "" : instanceInfo.getInstanceParams();
      34. log.info("[Dispatcher-{}|{}] start to dispatch job: {};instancePrams: {}.", jobId, instanceId, jobInfo, dbInstanceParams);
      35. // 查询当前运行的实例数
      36. long current = System.currentTimeMillis();
      37. Integer maxInstanceNum = jobInfo.getMaxInstanceNum();
      38. // 秒级任务只派发到一台机器,具体的 maxInstanceNum 由 TaskTracker 控制
      39. if (TimeExpressionType.FREQUENT_TYPES.contains(jobInfo.getTimeExpressionType())) {
      40. maxInstanceNum = 1;
      41. }
      42. //...
      43. // 获取当前最合适的 worker 列表
      44. List suitableWorkers = workerClusterQueryService.getSuitableWorkers(jobInfo);
      45. //...
      46. // 判断是否超载,在所有可用 worker 超载的情况下直接跳过当前任务
      47. suitableWorkers = filterOverloadWorker(suitableWorkers);
      48. //...
      49. //获取worker ip
      50. List workerIpList = suitableWorkers.stream().map(WorkerInfo::getAddress).collect(Collectors.toList());
      51. // 构造任务调度请求
      52. ServerScheduleJobReq req = constructServerScheduleJobReq(jobInfo, instanceInfo, workerIpList);
      53. // 发送请求(不可靠,需要一个后台线程定期轮询状态)
      54. //只取第一个worker
      55. WorkerInfo taskTracker = suitableWorkers.get(0);
      56. String taskTrackerAddress = taskTracker.getAddress();
      57. URL workerUrl = ServerURLFactory.dispatchJob2Worker(taskTrackerAddress);
      58. transportService.tell(taskTracker.getProtocol(), workerUrl, req);
      59. log.info("[Dispatcher-{}|{}] send schedule request to TaskTracker[protocol:{},address:{}] successfully: {}.", jobId, instanceId, taskTracker.getProtocol(), taskTrackerAddress, req);
      60. // 修改状态
      61. //修改实例表状态
      62. instanceInfoRepository.update4TriggerSucceed(instanceId, WAITING_WORKER_RECEIVE.getV(), current, taskTrackerAddress, now, instanceInfo.getStatus());
      63. // 装载缓存
      64. instanceMetadataService.loadJobInfo(instanceId, jobInfo);
      65. }
      66. /**
      67. * InstanceManager:
      68. * 第31行代码处:
      69. * 收尾完成的任务实例
      70. *
      71. * @param instanceId 任务实例ID
      72. * @param wfInstanceId 工作流实例ID,非必须
      73. * @param status 任务状态,有 成功/失败/手动停止
      74. * @param result 执行结果
      75. */
      76. public void processFinishedInstance(Long instanceId, Long wfInstanceId, InstanceStatus status, String result) {
      77. log.info("[Instance-{}] process finished, final status is {}.", instanceId, status.name());
      78. // 上报日志数据
      79. //时间轮延迟执行
      80. HashedWheelTimerHolder.INACCURATE_TIMER.schedule(() -> instanceLogService.sync(instanceId), 60, TimeUnit.SECONDS);
      81. // workflow 特殊处理
      82. if (wfInstanceId != null) {
      83. // 手动停止在工作流中也认为是失败(理论上不应该发生)
      84. workflowInstanceManager.move(wfInstanceId, instanceId, status, result);
      85. }
      86. // 告警
      87. if (status == InstanceStatus.FAILED) {
      88. alert(instanceId, result);
      89. }
      90. // 主动移除缓存,减小内存占用
      91. instanceMetadataService.invalidateJobInfo(instanceId);
      92. }
      93. /**
      94. * InstanceLogService:
      95. * 第96行代码处:
      96. * 将本地的任务实例运行日志同步到 mongoDB 存储,在任务执行结束后异步执行
      97. *
      98. * @param instanceId 任务实例ID
      99. */
      100. @Async(PJThreadPool.BACKGROUND_POOL)
      101. public void sync(Long instanceId) {
      102. Stopwatch sw = Stopwatch.createStarted();
      103. try {
      104. // 先持久化到本地文件
      105. File stableLogFile = genStableLogFile(instanceId);
      106. // 将文件推送到 MongoDB
      107. if (gridFsManager.available()) {
      108. try {
      109. gridFsManager.store(stableLogFile, GridFsManager.LOG_BUCKET, genMongoFileName(instanceId));
      110. log.info("[InstanceLog-{}] push local instanceLogs to mongoDB succeed, using: {}.", instanceId, sw.stop());
      111. } catch (Exception e) {
      112. log.warn("[InstanceLog-{}] push local instanceLogs to mongoDB failed.", instanceId, e);
      113. }
      114. }
      115. } catch (Exception e) {
      116. log.warn("[InstanceLog-{}] sync local instanceLogs failed.", instanceId, e);
      117. }
      118. // 删除本地数据库数据
      119. try {
      120. instanceId2LastReportTime.remove(instanceId);
      121. CommonUtils.executeWithRetry0(() -> localInstanceLogRepository.deleteByInstanceId(instanceId));
      122. log.info("[InstanceLog-{}] delete local instanceLog successfully.", instanceId);
      123. } catch (Exception e) {
      124. log.warn("[InstanceLog-{}] delete local instanceLog failed.", instanceId, e);
      125. }
      126. }
      127. /**
      128. * WorkerClusterQueryService:
      129. * 第50行代码处:
      130. * get worker for job
      131. *
      132. * @param jobInfo job
      133. * @return worker cluster info, sorted by metrics desc
      134. */
      135. public List getSuitableWorkers(JobInfoDO jobInfo) {
      136. //获取该集群所有的机器信息
      137. List workers = Lists.newLinkedList(getWorkerInfosByAppId(jobInfo.getAppId()).values());
      138. workers.removeIf(workerInfo -> filterWorker(workerInfo, jobInfo));
      139. DispatchStrategy dispatchStrategy = DispatchStrategy.of(jobInfo.getDispatchStrategy());
      140. switch (dispatchStrategy) {
      141. case RANDOM:
      142. //随机的方式就是打乱顺序
      143. Collections.shuffle(workers);
      144. break;
      145. case HEALTH_FIRST:
      146. //健康优先的方式需要计算下得分,按得分高低排序
      147. workers.sort((o1, o2) -> o2.getSystemMetrics().calculateScore() - o1.getSystemMetrics().calculateScore());
      148. break;
      149. default:
      150. // do nothing
      151. }
      152. // 限定集群大小(0代表不限制)
      153. if (!workers.isEmpty() && jobInfo.getMaxWorkerCount() > 0 && workers.size() > jobInfo.getMaxWorkerCount()) {
      154. workers = workers.subList(0, jobInfo.getMaxWorkerCount());
      155. }
      156. return workers;
      157. }
      158. /**
      159. * SystemMetrics:
      160. * 第171行代码处:
      161. * Calculate score, based on CPU and memory info.
      162. *
      163. * @return score
      164. */
      165. public int calculateScore() {
      166. if (score > 0) {
      167. return score;
      168. }
      169. // Memory is vital to TaskTracker, so we set the multiplier factor as 2.
      170. //未使用内存指标的权重为2
      171. double memScore = (jvmMaxMemory - jvmUsedMemory) * 2;
      172. // Calculate the remaining load of CPU. Multiplier is set as 1.
      173. //剩余可用cpu数的权重为1
      174. double cpuScore = cpuProcessors - cpuLoad;
      175. // Windows can not fetch CPU load, set cpuScore as 1.
      176. //Windows系统拿不到cpu使用数,所以cpu得分设置为1
      177. if (cpuScore > cpuProcessors) {
      178. cpuScore = 1;
      179. }
      180. //最终得分就是内存得分+cpu得分
      181. score = (int) (memScore + cpuScore);
      182. return score;
      183. }
      184. /**
      185. * ServerURLFactory:
      186. * 第69行代码处:
      187. */
      188. public static URL dispatchJob2Worker(String address) {
      189. return simileBuild(address, ServerType.WORKER, WORKER_PATH, WTT_HANDLER_RUN_JOB);
      190. }

              在上面第216行代码处会将请求跳转到WorkerActor.onReceiveServerScheduleJobReq方法:

      1. /**
      2. * WorkerActor:
      3. */
      4. @Handler(path = WTT_HANDLER_RUN_JOB)
      5. public void onReceiveServerScheduleJobReq(ServerScheduleJobReq req) {
      6. taskTrackerActor.onReceiveServerScheduleJobReq(req);
      7. }
      8. /**
      9. * TaskTrackerActor:
      10. * 服务器任务调度处理器
      11. */
      12. @Handler(path = WTT_HANDLER_RUN_JOB)
      13. public void onReceiveServerScheduleJobReq(ServerScheduleJobReq req) {
      14. log.debug("[TaskTrackerActor] server schedule job by request: {}.", req);
      15. Long instanceId = req.getInstanceId();
      16. // 区分轻量级任务模型以及重量级任务模型
      17. //单机执行的OpenApi/Corn/工作流是轻量级任务,其他的是重量级任务
      18. if (isLightweightTask(req)) {
      19. final LightTaskTracker taskTracker = LightTaskTrackerManager.getTaskTracker(instanceId);
      20. if (taskTracker != null) {
      21. log.warn("[TaskTrackerActor] LightTaskTracker({}) for instance(id={}) already exists.", taskTracker, instanceId);
      22. return;
      23. }
      24. // 判断是否已经 overload
      25. if (LightTaskTrackerManager.currentTaskTrackerSize() >= workerRuntime.getWorkerConfig().getMaxLightweightTaskNum() * LightTaskTrackerManager.OVERLOAD_FACTOR) {
      26. // ignore this request
      27. log.warn("[TaskTrackerActor] this worker is overload,ignore this request(instanceId={}),current size = {}!", instanceId, LightTaskTrackerManager.currentTaskTrackerSize());
      28. return;
      29. }
      30. if (LightTaskTrackerManager.currentTaskTrackerSize() >= workerRuntime.getWorkerConfig().getMaxLightweightTaskNum()) {
      31. log.warn("[TaskTrackerActor] this worker will be overload soon,current size = {}!", LightTaskTrackerManager.currentTaskTrackerSize());
      32. }
      33. // 创建轻量级任务
      34. //在轻量级本地缓存中添加TaskTracker
      35. LightTaskTrackerManager.atomicCreateTaskTracker(instanceId, ignore -> LightTaskTracker.create(req, workerRuntime));
      36. } else {
      37. HeavyTaskTracker taskTracker = HeavyTaskTrackerManager.getTaskTracker(instanceId);
      38. if (taskTracker != null) {
      39. log.warn("[TaskTrackerActor] HeavyTaskTracker({}) for instance(id={}) already exists.", taskTracker, instanceId);
      40. return;
      41. }
      42. // 判断是否已经 overload
      43. if (HeavyTaskTrackerManager.currentTaskTrackerSize() >= workerRuntime.getWorkerConfig().getMaxHeavyweightTaskNum()) {
      44. // ignore this request
      45. log.warn("[TaskTrackerActor] this worker is overload,ignore this request(instanceId={})! current size = {},", instanceId, HeavyTaskTrackerManager.currentTaskTrackerSize());
      46. return;
      47. }
      48. // 原子创建,防止多实例的存在
      49. //在重量级本地缓存中添加TaskTracker
      50. HeavyTaskTrackerManager.atomicCreateTaskTracker(instanceId, ignore -> HeavyTaskTracker.create(req, workerRuntime));
      51. }
      52. }

              由上可以看到,派发任务就是往TaskTracker里添加一条本地缓存。

              这里出现了TaskTracker,TaskTrackerProcessorTrackerProcessor的概念直接参考PowerJob作者写过的文章:

              这里我们继续分析下上面第51行代码处的HeavyTaskTracker.create方法:

      1. /**
      2. * HeavyTaskTracker:
      3. * 静态方法创建 TaskTracker
      4. *
      5. * @param req 服务端调度任务请求
      6. * @return API/CRON -> CommonTaskTracker, FIX_RATE/FIX_DELAY -> FrequentTaskTracker
      7. */
      8. public static HeavyTaskTracker create(ServerScheduleJobReq req, WorkerRuntime workerRuntime) {
      9. try {
      10. TimeExpressionType timeExpressionType = TimeExpressionType.valueOf(req.getTimeExpressionType());
      11. switch (timeExpressionType) {
      12. case FIXED_RATE:
      13. case FIXED_DELAY:
      14. return new FrequentTaskTracker(req, workerRuntime);
      15. default:
      16. //这里我们分析下CommonTaskTracker构造器的实现
      17. return new CommonTaskTracker(req, workerRuntime);
      18. }
      19. } catch (Exception e) {
      20. reportCreateErrorToServer(req, workerRuntime, e);
      21. }
      22. return null;
      23. }
      24. /**
      25. * CommonTaskTracker:
      26. * 第17行代码处:
      27. */
      28. protected CommonTaskTracker(ServerScheduleJobReq req, WorkerRuntime workerRuntime) {
      29. super(req, workerRuntime);
      30. }
      31. /**
      32. * HeavyTaskTracker:
      33. */
      34. protected HeavyTaskTracker(ServerScheduleJobReq req, WorkerRuntime workerRuntime) {
      35. // 初始化成员变量
      36. super(req, workerRuntime);
      37. // 赋予时间表达式类型
      38. instanceInfo.setTimeExpressionType(TimeExpressionType.valueOf(req.getTimeExpressionType()).getV());
      39. // 保护性操作
      40. instanceInfo.setThreadConcurrency(Math.max(1, instanceInfo.getThreadConcurrency()));
      41. this.ptStatusHolder = new ProcessorTrackerStatusHolder(instanceId, req.getMaxWorkerCount(), req.getAllWorkerAddress());
      42. this.taskPersistenceService = workerRuntime.getTaskPersistenceService();
      43. // 构建缓存
      44. taskId2BriefInfo = CacheBuilder.newBuilder().maximumSize(1024).build();
      45. // 构建分段锁
      46. //SegmentLock是自己实现的ReentrantLock数组(使用数组是为了提高并发度)
      47. segmentLock = new SegmentLock(UPDATE_CONCURRENCY);
      48. // 子类自定义初始化操作
      49. initTaskTracker(req);
      50. log.info("[TaskTracker-{}] create TaskTracker successfully.", instanceId);
      51. }
      52. /**
      53. * CommonTaskTracker:
      54. * 第53行代码处:
      55. *
      56. * @param req 服务器调度任务实例运行请求
      57. */
      58. @Override
      59. protected void initTaskTracker(ServerScheduleJobReq req) {
      60. // CommonTaskTrackerTimingPool 缩写
      61. String poolName = String.format("ctttp-%d", req.getInstanceId()) + "-%d";
      62. ThreadFactory factory = new ThreadFactoryBuilder().setNameFormat(poolName).build();
      63. this.scheduledPool = Executors.newScheduledThreadPool(2, factory);
      64. // 持久化根任务
      65. persistenceRootTask();
      66. // 开启定时状态检查
      67. int delay = Integer.parseInt(System.getProperty(PowerJobDKey.WORKER_STATUS_CHECK_PERIOD, "13"));
      68. scheduledPool.scheduleWithFixedDelay(new StatusCheckRunnable(), 3, delay, TimeUnit.SECONDS);
      69. // 如果是 MR 任务,则需要启动执行器动态检测装置
      70. ExecuteType executeType = ExecuteType.valueOf(req.getExecuteType());
      71. if (executeType == ExecuteType.MAP || executeType == ExecuteType.MAP_REDUCE) {
      72. scheduledPool.scheduleAtFixedRate(new WorkerDetector(), 1, 1, TimeUnit.MINUTES);
      73. }
      74. // 最后启动任务派发器,否则会出现 TaskTracker 还未创建完毕 ProcessorTracker 已开始汇报状态的情况
      75. scheduledPool.scheduleWithFixedDelay(new Dispatcher(), 10, 5000, TimeUnit.MILLISECONDS);
      76. }
      77. /**
      78. * 第73行代码处:
      79. * 持久化根任务,只有完成持久化才能视为任务开始running(先持久化,再报告server)
      80. */
      81. private void persistenceRootTask() {
      82. TaskDO rootTask = new TaskDO();
      83. rootTask.setStatus(TaskStatus.WAITING_DISPATCH.getValue());
      84. rootTask.setInstanceId(instanceInfo.getInstanceId());
      85. rootTask.setTaskId(ROOT_TASK_ID);
      86. rootTask.setFailedCnt(0);
      87. rootTask.setAddress(workerRuntime.getWorkerAddress());
      88. //这里需要留意下,根任务的名称为OMS_ROOT_TASK(ROOT_TASK_NAME)
      89. rootTask.setTaskName(TaskConstant.ROOT_TASK_NAME);
      90. rootTask.setCreatedTime(System.currentTimeMillis());
      91. rootTask.setLastModifiedTime(System.currentTimeMillis());
      92. rootTask.setLastReportTime(-1L);
      93. rootTask.setSubInstanceId(instanceId);
      94. if (taskPersistenceService.save(rootTask)) {
      95. log.info("[TaskTracker-{}] create root task successfully.", instanceId);
      96. } else {
      97. log.error("[TaskTracker-{}] create root task failed.", instanceId);
      98. throw new PowerJobException("create root task failed for instance: " + instanceId);
      99. }
      100. }

              在上面第77行代码处、第82行代码处和第86行代码处的initTaskTracker方法中,分别开启了三个定时任务,这里我们主要分析下StatusCheckRunnable和Dispatcher(WorkerDetector放在《较真儿学源码系列-PowerJob MapReduce源码分析》中分析):

      2.3.1.1 StatusCheckRunnable
      1. /**
      2. * StatusCheckRunnable:
      3. */
      4. @Override
      5. public void run() {
      6. try {
      7. innerRun();
      8. } catch (Exception e) {
      9. log.warn("[TaskTracker-{}] status checker execute failed, please fix the bug (@tjq)!", instanceId, e);
      10. }
      11. }
      12. /**
      13. * 第7行代码处:
      14. */
      15. @SuppressWarnings("squid:S3776")
      16. private void innerRun() {
      17. //获取任务实例产生的各个Task状态,用于分析任务实例执行情况
      18. InstanceStatisticsHolder holder = getInstanceStatisticsHolder(instanceId);
      19. long finishedNum = holder.succeedNum + holder.failedNum;
      20. long unfinishedNum = holder.waitingDispatchNum + holder.workerUnreceivedNum + holder.receivedNum + holder.runningNum;
      21. log.debug("[TaskTracker-{}] status check result: {}", instanceId, holder);
      22. //组装上报参数
      23. TaskTrackerReportInstanceStatusReq req = new TaskTrackerReportInstanceStatusReq();
      24. req.setAppId(workerRuntime.getAppId());
      25. req.setJobId(instanceInfo.getJobId());
      26. req.setInstanceId(instanceId);
      27. req.setWfInstanceId(instanceInfo.getWfInstanceId());
      28. req.setTotalTaskNum(finishedNum + unfinishedNum);
      29. req.setSucceedTaskNum(holder.succeedNum);
      30. req.setFailedTaskNum(holder.failedNum);
      31. req.setReportTime(System.currentTimeMillis());
      32. req.setStartTime(createTime);
      33. req.setSourceAddress(workerRuntime.getWorkerAddress());
      34. boolean success = false;
      35. String result = null;
      36. // 2. 如果未完成任务数为0,判断是否真正结束,并获取真正结束任务的执行结果
      37. if (unfinishedNum == 0) {
      38. // 数据库中一个任务都没有,说明根任务创建失败,该任务实例失败
      39. if (finishedNum == 0) {
      40. finished.set(true);
      41. result = SystemInstanceResult.TASK_INIT_FAILED;
      42. } else {
      43. ExecuteType executeType = ExecuteType.valueOf(instanceInfo.getExecuteType());
      44. switch (executeType) {
      45. // STANDALONE 只有一个任务,完成即结束
      46. case STANDALONE:
      47. finished.set(true);
      48. List allTask = taskPersistenceService.getAllTask(instanceId, instanceId);
      49. if (CollectionUtils.isEmpty(allTask) || allTask.size() > 1) {
      50. result = SystemInstanceResult.UNKNOWN_BUG;
      51. log.warn("[TaskTracker-{}] there must have some bug in TaskTracker.", instanceId);
      52. } else {
      53. result = allTask.get(0).getResult();
      54. success = allTask.get(0).getStatus() == TaskStatus.WORKER_PROCESS_SUCCESS.getValue();
      55. }
      56. break;
      57. case MAP:
      58. //...
      59. break;
      60. default:
      61. //...
      62. }
      63. }
      64. }
      65. // 3. 检查任务实例整体是否超时
      66. if (isTimeout()) {
      67. finished.set(true);
      68. success = false;
      69. result = SystemInstanceResult.INSTANCE_EXECUTE_TIMEOUT;
      70. }
      71. // 4. 执行完毕,报告服务器
      72. if (finished.get()) {
      73. req.setResult(result);
      74. // 上报追加的工作流上下文信息
      75. req.setAppendedWfContext(appendedWfContext);
      76. req.setInstanceStatus(success ? InstanceStatus.SUCCEED.getV() : InstanceStatus.FAILED.getV());
      77. reportFinalStatusThenDestroy(workerRuntime, req);
      78. return;
      79. }
      80. // 5. 未完成,上报状态
      81. req.setInstanceStatus(InstanceStatus.RUNNING.getV());
      82. TransportUtils.ttReportInstanceStatus(req, workerRuntime.getServerDiscoveryService().getCurrentServerAddress(), workerRuntime.getTransporter());
      83. // 6.1 定期检查 -> 重试派发后未确认的任务
      84. long currentMS = System.currentTimeMillis();
      85. if (holder.workerUnreceivedNum != 0) {
      86. taskPersistenceService.getTaskByStatus(instanceId, TaskStatus.DISPATCH_SUCCESS_WORKER_UNCHECK, 100).forEach(uncheckTask -> {
      87. long elapsedTime = currentMS - uncheckTask.getLastModifiedTime();
      88. if (elapsedTime > DISPATCH_TIME_OUT_MS) {
      89. TaskDO updateEntity = new TaskDO();
      90. updateEntity.setStatus(TaskStatus.WAITING_DISPATCH.getValue());
      91. // 特殊任务只能本机执行
      92. if (!TaskConstant.LAST_TASK_NAME.equals(uncheckTask.getTaskName())) {
      93. updateEntity.setAddress(RemoteConstant.EMPTY_ADDRESS);
      94. }
      95. // 失败次数 + 1
      96. updateEntity.setFailedCnt(uncheckTask.getFailedCnt() + 1);
      97. taskPersistenceService.updateTask(instanceId, uncheckTask.getTaskId(), updateEntity);
      98. log.warn("[TaskTracker-{}] task(id={},name={}) try to dispatch again due to unreceived the response from ProcessorTracker.",
      99. instanceId, uncheckTask.getTaskId(), uncheckTask.getTaskName());
      100. }
      101. });
      102. }
      103. // 6.2 定期检查 -> 重新执行被派发到宕机ProcessorTracker上的任务
      104. List disconnectedPTs = ptStatusHolder.getAllDisconnectedProcessorTrackers();
      105. if (!disconnectedPTs.isEmpty()) {
      106. log.warn("[TaskTracker-{}] some ProcessorTracker disconnected from TaskTracker,their address is {}.", instanceId, disconnectedPTs);
      107. if (taskPersistenceService.updateLostTasks(instanceId, disconnectedPTs, true)) {
      108. ptStatusHolder.remove(disconnectedPTs);
      109. log.warn("[TaskTracker-{}] removed these ProcessorTracker from StatusHolder: {}", instanceId, disconnectedPTs);
      110. }
      111. }
      112. }

              上报到服务端的逻辑这里就不再看了(更新instance_info实例表数据)。整体上来说,StatusCheckRunnable的作用就是检查任务的执行情况,并上报到服务端。

      2.3.1.2 Dispatcher
      1. /**
      2. * Dispatcher:
      3. */
      4. @Override
      5. public void run() {
      6. if (finished.get()) {
      7. return;
      8. }
      9. Stopwatch stopwatch = Stopwatch.createStarted();
      10. // 1. 获取可以派发任务的 ProcessorTracker
      11. List availablePtIps = ptStatusHolder.getAvailableProcessorTrackers();
      12. // 2. 没有可用 ProcessorTracker,本次不派发
      13. if (availablePtIps.isEmpty()) {
      14. log.debug("[TaskTracker-{}] no available ProcessorTracker now.", instanceId);
      15. return;
      16. }
      17. // 3. 避免大查询,分批派发任务
      18. long currentDispatchNum = 0;
      19. //这里需要留意下,最大分发次数=可用ProcessorTracker数量*实例并发度*2
      20. long maxDispatchNum = availablePtIps.size() * instanceInfo.getThreadConcurrency() * 2L;
      21. AtomicInteger index = new AtomicInteger(0);
      22. // 4. 循环查询数据库,获取需要派发的任务
      23. while (maxDispatchNum > currentDispatchNum) {
      24. //每次最多查询100个任务
      25. int dbQueryLimit = Math.min(DB_QUERY_LIMIT, (int) maxDispatchNum);
      26. //获取等待调度器调度的任务
      27. List needDispatchTasks = taskPersistenceService.getTaskByStatus(instanceId, TaskStatus.WAITING_DISPATCH, dbQueryLimit);
      28. currentDispatchNum += needDispatchTasks.size();
      29. needDispatchTasks.forEach(task -> {
      30. // 获取 ProcessorTracker 地址,如果 Task 中自带了 Address,则使用该 Address
      31. String ptAddress = task.getAddress();
      32. if (StringUtils.isEmpty(ptAddress) || RemoteConstant.EMPTY_ADDRESS.equals(ptAddress)) {
      33. //否则,从可用ProcessorTracker中取余获取一个
      34. ptAddress = availablePtIps.get(index.getAndIncrement() % availablePtIps.size());
      35. }
      36. //分发任务
      37. dispatchTask(task, ptAddress);
      38. });
      39. // 数量不足 或 查询失败,则终止循环
      40. if (needDispatchTasks.size() < dbQueryLimit) {
      41. break;
      42. }
      43. }
      44. log.debug("[TaskTracker-{}] dispatched {} tasks,using time {}.", instanceId, currentDispatchNum, stopwatch.stop());
      45. }
      46. /**
      47. * ProcessorTrackerStatusHolder:
      48. * 第14行代码处:
      49. * 获取可用 ProcessorTracker 的IP地址
      50. */
      51. public List getAvailableProcessorTrackers() {
      52. List result = Lists.newLinkedList();
      53. address2Status.forEach((address, ptStatus) -> {
      54. if (ptStatus.available()) {
      55. result.add(address);
      56. }
      57. });
      58. return result;
      59. }
      60. /**
      61. * ProcessorTrackerStatus:
      62. * 第66行代码处:
      63. * 是否可用
      64. */
      65. public boolean available() {
      66. // 未曾派发过,默认可用
      67. if (!dispatched) {
      68. return true;
      69. }
      70. // 已派发但未收到响应,则不可用
      71. if (!connected) {
      72. return false;
      73. }
      74. // 长时间未收到心跳消息,则不可用
      75. if (isTimeout()) {
      76. return false;
      77. }
      78. // 留有过多待处理任务,则不可用
      79. if (remainTaskNum >= DISPATCH_THRESHOLD) {
      80. return false;
      81. }
      82. // TODO:后续考虑加上机器健康度等信息
      83. return true;
      84. }
      85. /**
      86. * 第91行代码处:
      87. * 是否超时(超过一定时间没有收到心跳)
      88. */
      89. public boolean isTimeout() {
      90. if (dispatched) {
      91. //系统当前时间-上次活跃时间>心跳超时时间
      92. return System.currentTimeMillis() - lastActiveTime > HEARTBEAT_TIMEOUT_MS;
      93. }
      94. // 未曾派发过任务的机器,不用处理
      95. return false;
      96. }
      97. /**
      98. * HeavyTaskTracker:
      99. * 第45行代码处:
      100. * 派发任务到 ProcessorTracker
      101. *
      102. * @param task 需要被执行的任务
      103. * @param processorTrackerAddress ProcessorTracker的地址(IP:Port)
      104. */
      105. protected void dispatchTask(TaskDO task, String processorTrackerAddress) {
      106. // 1. 持久化,更新数据库(如果更新数据库失败,可能导致重复执行,先不处理)
      107. TaskDO updateEntity = new TaskDO();
      108. updateEntity.setStatus(TaskStatus.DISPATCH_SUCCESS_WORKER_UNCHECK.getValue());
      109. // 写入处理该任务的 ProcessorTracker
      110. updateEntity.setAddress(processorTrackerAddress);
      111. boolean success = taskPersistenceService.updateTask(instanceId, task.getTaskId(), updateEntity);
      112. if (!success) {
      113. log.warn("[TaskTracker-{}] dispatch task(taskId={},taskName={}) failed due to update task status failed.", instanceId, task.getTaskId(), task.getTaskName());
      114. return;
      115. }
      116. // 2. 更新 ProcessorTrackerStatus 状态
      117. ptStatusHolder.getProcessorTrackerStatus(processorTrackerAddress).setDispatched(true);
      118. // 3. 初始化缓存
      119. taskId2BriefInfo.put(task.getTaskId(), new TaskBriefInfo(task.getTaskId(), TaskStatus.DISPATCH_SUCCESS_WORKER_UNCHECK, -1L));
      120. // 4. 任务派发
      121. TaskTrackerStartTaskReq startTaskReq = new TaskTrackerStartTaskReq(instanceInfo, task, workerRuntime.getWorkerAddress());
      122. TransportUtils.ttStartPtTask(startTaskReq, processorTrackerAddress, workerRuntime.getTransporter());
      123. log.debug("[TaskTracker-{}] dispatch task(taskId={},taskName={}) successfully.", instanceId, task.getTaskId(), task.getTaskName());
      124. }
      125. /**
      126. * TransportUtils:
      127. * 第146行代码处:
      128. */
      129. public static void ttStartPtTask(TaskTrackerStartTaskReq req, String address, Transporter transporter) {
      130. final URL url = easyBuildUrl(ServerType.WORKER, WPT_PATH, WPT_HANDLER_START_TASK, address);
      131. transporter.tell(url, req);
      132. }
      133. /**
      134. * ProcessorTrackerActor:
      135. * 第157行代码处:
      136. * 处理来自TaskTracker的task执行请求
      137. *
      138. * @param req 请求
      139. */
      140. @Handler(path = RemoteConstant.WPT_HANDLER_START_TASK, processType = ProcessType.NO_BLOCKING)
      141. public void onReceiveTaskTrackerStartTaskReq(TaskTrackerStartTaskReq req) {
      142. Long instanceId = req.getInstanceInfo().getInstanceId();
      143. // 创建 ProcessorTracker 一定能成功
      144. ProcessorTracker processorTracker = ProcessorTrackerManager.getProcessorTracker(
      145. instanceId,
      146. req.getTaskTrackerAddress(),
      147. () -> new ProcessorTracker(req, workerRuntime));
      148. TaskDO task = new TaskDO();
      149. task.setTaskId(req.getTaskId());
      150. task.setTaskName(req.getTaskName());
      151. task.setTaskContent(req.getTaskContent());
      152. task.setFailedCnt(req.getTaskCurrentRetryNums());
      153. task.setSubInstanceId(req.getSubInstanceId());
      154. processorTracker.submitTask(task);
      155. }
      156. /**
      157. * ProcessorTracker:
      158. * 提交任务到线程池执行
      159. * 1.0版本:TaskTracker有任务就dispatch,导致 ProcessorTracker 本地可能堆积过多的任务,造成内存压力。为此 ProcessorTracker 在线程
      160. * 池队列堆积到一定程度时,会将数据持久化到DB,然后通过异步线程定时从数据库中取回任务,重新提交执行。
      161. * 联动:数据库的SPID设计、TaskStatus段落设计等,全部取消...
      162. * last commitId: 341953aceceafec0fbe7c3d9a3e26451656b945e
      163. * 2.0版本:ProcessorTracker定时向TaskTracker发送心跳消息,心跳消息中包含了当前线程池队列任务个数,TaskTracker根据ProcessorTracker
      164. * 的状态判断能否继续派发任务。因此,ProcessorTracker本地不会堆积过多任务,故删除 持久化机制 ╥﹏╥...!
      165. *
      166. * @param newTask 需要提交到线程池执行的任务
      167. */
      168. public void submitTask(TaskDO newTask) {
      169. // 一旦 ProcessorTracker 出现异常,所有提交到此处的任务直接返回失败,防止形成死锁
      170. // 死锁分析:TT创建PT,PT创建失败,无法定期汇报心跳,TT长时间未收到PT心跳,认为PT宕机(确实宕机了),无法选择可用的PT再次派发任务,死锁形成,GG斯密达 T_T
      171. if (lethal) {
      172. ProcessorReportTaskStatusReq report = new ProcessorReportTaskStatusReq()
      173. .setInstanceId(instanceId)
      174. .setSubInstanceId(newTask.getSubInstanceId())
      175. .setTaskId(newTask.getTaskId())
      176. .setStatus(TaskStatus.WORKER_PROCESS_FAILED.getValue())
      177. .setResult(lethalReason)
      178. .setReportTime(System.currentTimeMillis());
      179. TransportUtils.ptReportTask(report, taskTrackerAddress, workerRuntime);
      180. return;
      181. }
      182. boolean success = false;
      183. // 1. 设置值并提交执行
      184. newTask.setInstanceId(instanceInfo.getInstanceId());
      185. newTask.setAddress(taskTrackerAddress);
      186. HeavyProcessorRunnable heavyProcessorRunnable = new HeavyProcessorRunnable(instanceInfo, taskTrackerAddress, newTask, processorBean, omsLogger, statusReportRetryQueue, workerRuntime);
      187. try {
      188. threadPool.submit(heavyProcessorRunnable);
      189. success = true;
      190. } catch (RejectedExecutionException ignore) {
      191. log.warn("[ProcessorTracker-{}] submit task(taskId={},taskName={}) to ThreadPool failed due to ThreadPool has too much task waiting to process, this task will dispatch to other ProcessorTracker.",
      192. instanceId, newTask.getTaskId(), newTask.getTaskName());
      193. } catch (Exception e) {
      194. log.error("[ProcessorTracker-{}] submit task(taskId={},taskName={}) to ThreadPool failed.", instanceId, newTask.getTaskId(), newTask.getTaskName(), e);
      195. }
      196. // 2. 回复接收成功
      197. if (success) {
      198. ProcessorReportTaskStatusReq reportReq = new ProcessorReportTaskStatusReq();
      199. reportReq.setInstanceId(instanceId);
      200. reportReq.setSubInstanceId(newTask.getSubInstanceId());
      201. reportReq.setTaskId(newTask.getTaskId());
      202. reportReq.setStatus(TaskStatus.WORKER_RECEIVED.getValue());
      203. reportReq.setReportTime(System.currentTimeMillis());
      204. TransportUtils.ptReportTask(reportReq, taskTrackerAddress, workerRuntime);
      205. log.debug("[ProcessorTracker-{}] submit task(taskId={}, taskName={}) success, current queue size: {}.",
      206. instanceId, newTask.getTaskId(), newTask.getTaskName(), threadPool.getQueue().size());
      207. }
      208. }

              可以看到,在最后,Dispatcher会创建出一个HeavyProcessorRunnable的线程来执行,里面存放着需要执行的任务实例、执行地址、执行处理器等信息。接下来看下其实现:

      1. /**
      2. * HeavyProcessorRunnable:
      3. */
      4. @Override
      5. @SuppressWarnings("squid:S2142")
      6. public void run() {
      7. // 切换线程上下文类加载器(否则用的是 Worker 类加载器,不存在容器类,在序列化/反序列化时会报 ClassNotFoundException)
      8. Thread.currentThread().setContextClassLoader(processorBean.getClassLoader());
      9. try {
      10. innerRun();
      11. } catch (InterruptedException ignore) {
      12. // ignore
      13. } catch (Throwable e) {
      14. reportStatus(TaskStatus.WORKER_PROCESS_FAILED, e.toString(), null, null);
      15. log.error("[ProcessorRunnable-{}] execute failed, please contact the author(@KFCFans) to fix the bug!", task.getInstanceId(), e);
      16. } finally {
      17. ThreadLocalStore.clear();
      18. }
      19. }
      20. /**
      21. * 第10行代码处:
      22. */
      23. public void innerRun() throws InterruptedException {
      24. //获取执行处理器
      25. final BasicProcessor processor = processorBean.getProcessor();
      26. String taskId = task.getTaskId();
      27. Long instanceId = task.getInstanceId();
      28. log.debug("[ProcessorRunnable-{}] start to run task(taskId={}&taskName={})", instanceId, taskId, task.getTaskName());
      29. //缓存
      30. ThreadLocalStore.setTask(task);
      31. ThreadLocalStore.setRuntimeMeta(workerRuntime);
      32. // 0. 构造任务上下文
      33. WorkflowContext workflowContext = constructWorkflowContext();
      34. TaskContext taskContext = constructTaskContext();
      35. taskContext.setWorkflowContext(workflowContext);
      36. // 1. 上报执行信息
      37. reportStatus(TaskStatus.WORKER_PROCESSING, null, null, null);
      38. ProcessResult processResult;
      39. ExecuteType executeType = ExecuteType.valueOf(instanceInfo.getExecuteType());
      40. // 2. 根任务 & 广播执行 特殊处理
      41. if (TaskConstant.ROOT_TASK_NAME.equals(task.getTaskName()) && executeType == ExecuteType.BROADCAST) {
      42. // 广播执行:先选本机执行 preProcess,完成后 TaskTracker 再为所有 Worker 生成子 Task
      43. handleBroadcastRootTask(instanceId, taskContext);
      44. return;
      45. }
      46. // 3. 最终任务特殊处理(一定和 TaskTracker 处于相同的机器)
      47. if (TaskConstant.LAST_TASK_NAME.equals(task.getTaskName())) {
      48. handleLastTask(taskId, instanceId, taskContext, executeType);
      49. return;
      50. }
      51. // 4. 正式提交运行
      52. try {
      53. //这里的process方法也就是我们自己写的业务方法
      54. processResult = processor.process(taskContext);
      55. if (processResult == null) {
      56. processResult = new ProcessResult(false, "ProcessResult can't be null");
      57. }
      58. } catch (Throwable e) {
      59. log.warn("[ProcessorRunnable-{}] task(id={},name={}) process failed.", instanceId, taskContext.getTaskId(), taskContext.getTaskName(), e);
      60. processResult = new ProcessResult(false, e.toString());
      61. }
      62. //上报执行结果
      63. reportStatus(processResult.isSuccess() ? TaskStatus.WORKER_PROCESS_SUCCESS : TaskStatus.WORKER_PROCESS_FAILED, suit(processResult.getMsg()), null, workflowContext.getAppendedContextData());
      64. }

              Dispatcher的作用就是分发任务,并执行任务。也就是会有一个单独的线程(HeavyProcessorRunnable),会定时从任务表中拿取任务,然后执行我们自己实现的process方法,并会在执行前和执行后上报执行信息。

      2.3.2 CheckRunningInstance

      1. /**
      2. * InstanceStatusCheckService:
      3. * 检查运行中的实例
      4. * RUNNING 超时:TaskTracker down,断开与 server 的心跳连接
      5. */
      6. public void checkRunningInstance() {
      7. Stopwatch stopwatch = Stopwatch.createStarted();
      8. // 查询 DB 获取该 Server 需要负责的 AppGroup
      9. //获取appId
      10. List allAppIds = appInfoRepository.listAppIdByCurrentServer(transportService.defaultProtocol().getAddress());
      11. if (CollectionUtils.isEmpty(allAppIds)) {
      12. log.info("[InstanceStatusChecker] current server has no app's job to check");
      13. return;
      14. }
      15. try {
      16. // 检查 RUNNING 状态的任务(一定时间没收到 TaskTracker 的状态报告,视为失败)
      17. Lists.partition(allAppIds, MAX_BATCH_NUM_APP).forEach(this::handleRunningInstance);
      18. } catch (Exception e) {
      19. log.error("[InstanceStatusChecker] RunningInstance status check failed.", e);
      20. }
      21. log.info("[InstanceStatusChecker] RunningInstance status check using {}.", stopwatch.stop());
      22. }
      23. /**
      24. * 第17行代码处:
      25. */
      26. private void handleRunningInstance(List partAppIds) {
      27. // 3. 检查 RUNNING 状态的任务(一定时间没收到 TaskTracker 的状态报告,视为失败)
      28. long threshold = System.currentTimeMillis() - RUNNING_TIMEOUT_MS;
      29. //查找修改时间距离现在超过1分钟的实例
      30. List failedInstances = instanceInfoRepository.selectBriefInfoByAppIdInAndStatusAndGmtModifiedBefore(partAppIds, InstanceStatus.RUNNING.getV(), new Date(threshold), PageRequest.of(0, MAX_BATCH_NUM_INSTANCE));
      31. while (!failedInstances.isEmpty()) {
      32. // collect job id
      33. Set jobIds = failedInstances.stream().map(BriefInstanceInfo::getJobId).collect(Collectors.toSet());
      34. // query job info and map
      35. Map jobInfoMap = jobInfoRepository.findByIdIn(jobIds).stream().collect(Collectors.toMap(JobInfoDO::getId, e -> e));
      36. log.warn("[InstanceStatusCheckService] find some instances have not received status report for a long time : {}", failedInstances.stream().map(BriefInstanceInfo::getInstanceId).collect(Collectors.toList()));
      37. failedInstances.forEach(instance -> {
      38. Optional jobInfoOpt = Optional.ofNullable(jobInfoMap.get(instance.getJobId()));
      39. if (!jobInfoOpt.isPresent()) {
      40. final Optional opt = instanceInfoRepository.findById(instance.getId());
      41. opt.ifPresent(e -> updateFailedInstance(e, SystemInstanceResult.REPORT_TIMEOUT));
      42. return;
      43. }
      44. TimeExpressionType timeExpressionType = TimeExpressionType.of(jobInfoOpt.get().getTimeExpressionType());
      45. SwitchableStatus switchableStatus = SwitchableStatus.of(jobInfoOpt.get().getStatus());
      46. // 如果任务已关闭,则不进行重试,将任务置为失败即可;秒级任务也直接置为失败,由派发器重新调度
      47. if (switchableStatus != SwitchableStatus.ENABLE || TimeExpressionType.FREQUENT_TYPES.contains(timeExpressionType.getV())) {
      48. final Optional opt = instanceInfoRepository.findById(instance.getId());
      49. opt.ifPresent(e -> updateFailedInstance(e, SystemInstanceResult.REPORT_TIMEOUT));
      50. return;
      51. }
      52. // CRON 和 API一样,失败次数 + 1,根据重试配置进行重试
      53. if (instance.getRunningTimes() < jobInfoOpt.get().getInstanceRetryNum()) {
      54. //更新实例表状态为等待派发
      55. dispatchService.redispatchAsync(instance.getInstanceId(), InstanceStatus.RUNNING.getV());
      56. } else {
      57. final Optional opt = instanceInfoRepository.findById(instance.getId());
      58. opt.ifPresent(e -> updateFailedInstance(e, SystemInstanceResult.REPORT_TIMEOUT));
      59. }
      60. });
      61. threshold = System.currentTimeMillis() - RUNNING_TIMEOUT_MS;
      62. //重新查找修改时间距离现在超过1分钟的实例,并循环
      63. failedInstances = instanceInfoRepository.selectBriefInfoByAppIdInAndStatusAndGmtModifiedBefore(partAppIds, InstanceStatus.RUNNING.getV(), new Date(threshold), PageRequest.of(0, MAX_BATCH_NUM_INSTANCE));
      64. }
      65. }

      3 客户端

      3.1 PowerJobSpringWorker

              客户端的启动类同样没有什么逻辑,查看下初始化类PowerJobSpringWorker:

      1. /**
      2. * PowerJobSpringWorker:
      3. */
      4. @Override
      5. public void afterPropertiesSet() throws Exception {
      6. powerJobWorker = new PowerJobWorker(config);
      7. powerJobWorker.init();
      8. }
      9. /**
      10. * PowerJobWorker:
      11. */
      12. public void init() throws Exception {
      13. if (!initialized.compareAndSet(false, true)) {
      14. log.warn("[PowerJobWorker] please do not repeat the initialization");
      15. return;
      16. }
      17. Stopwatch stopwatch = Stopwatch.createStarted();
      18. log.info("[PowerJobWorker] start to initialize PowerJobWorker...");
      19. PowerJobWorkerConfig config = workerRuntime.getWorkerConfig();
      20. CommonUtils.requireNonNull(config, "can't find PowerJobWorkerConfig, please set PowerJobWorkerConfig first");
      21. try {
      22. //打印banner
      23. PowerBannerPrinter.print();
      24. // 校验 appName
      25. if (!config.isEnableTestMode()) {
      26. assertAppName();
      27. } else {
      28. log.warn("[PowerJobWorker] using TestMode now, it's dangerous if this is production env.");
      29. }
      30. // 初始化元数据
      31. String workerAddress = NetUtils.getLocalHost() + ":" + config.getPort();
      32. workerRuntime.setWorkerAddress(workerAddress);
      33. // 初始化 线程池
      34. final ExecutorManager executorManager = new ExecutorManager(workerRuntime.getWorkerConfig());
      35. workerRuntime.setExecutorManager(executorManager);
      36. // 初始化 ProcessorLoader
      37. //动态加载类的工厂,可以不依赖于Spring的实现
      38. ProcessorLoader processorLoader = buildProcessorLoader(workerRuntime);
      39. workerRuntime.setProcessorLoader(processorLoader);
      40. // 初始化 actor
      41. TaskTrackerActor taskTrackerActor = new TaskTrackerActor(workerRuntime);
      42. ProcessorTrackerActor processorTrackerActor = new ProcessorTrackerActor(workerRuntime);
      43. WorkerActor workerActor = new WorkerActor(workerRuntime, taskTrackerActor);
      44. // 初始化通讯引擎
      45. EngineConfig engineConfig = new EngineConfig()
      46. .setType(config.getProtocol().name())
      47. .setServerType(ServerType.WORKER)
      48. .setBindAddress(new Address().setHost(NetUtils.getLocalHost()).setPort(config.getPort()))
      49. .setActorList(Lists.newArrayList(taskTrackerActor, processorTrackerActor, workerActor));
      50. //之前服务端启动的时候也会执行PowerJobRemoteEngine.start方法
      51. EngineOutput engineOutput = remoteEngine.start(engineConfig);
      52. workerRuntime.setTransporter(engineOutput.getTransporter());
      53. // 连接 server
      54. ServerDiscoveryService serverDiscoveryService = new ServerDiscoveryService(workerRuntime.getAppId(), workerRuntime.getWorkerConfig());
      55. serverDiscoveryService.start(workerRuntime.getExecutorManager().getCoreExecutor());
      56. workerRuntime.setServerDiscoveryService(serverDiscoveryService);
      57. log.info("[PowerJobWorker] PowerJobRemoteEngine initialized successfully.");
      58. // 初始化日志系统
      59. OmsLogHandler omsLogHandler = new OmsLogHandler(workerAddress, workerRuntime.getTransporter(), serverDiscoveryService);
      60. workerRuntime.setOmsLogHandler(omsLogHandler);
      61. // 初始化存储
      62. TaskPersistenceService taskPersistenceService = new TaskPersistenceService(workerRuntime.getWorkerConfig().getStoreStrategy());
      63. taskPersistenceService.init();
      64. workerRuntime.setTaskPersistenceService(taskPersistenceService);
      65. log.info("[PowerJobWorker] local storage initialized successfully.");
      66. // 初始化定时任务
      67. workerRuntime.getExecutorManager().getCoreExecutor().scheduleAtFixedRate(new WorkerHealthReporter(workerRuntime), 0, config.getHealthReportInterval(), TimeUnit.SECONDS);
      68. workerRuntime.getExecutorManager().getCoreExecutor().scheduleWithFixedDelay(omsLogHandler.logSubmitter, 0, 5, TimeUnit.SECONDS);
      69. log.info("[PowerJobWorker] PowerJobWorker initialized successfully, using time: {}, congratulations!", stopwatch);
      70. } catch (Exception e) {
      71. log.error("[PowerJobWorker] initialize PowerJobWorker failed, using {}.", stopwatch, e);
      72. throw e;
      73. }
      74. }
      75. /**
      76. * 第31行代码处:
      77. */
      78. @SuppressWarnings("rawtypes")
      79. private void assertAppName() {
      80. PowerJobWorkerConfig config = workerRuntime.getWorkerConfig();
      81. String appName = config.getAppName();
      82. Objects.requireNonNull(appName, "appName can't be empty!");
      83. String url = "http://%s/server/assert?appName=%s";
      84. for (String server : config.getServerAddress()) {
      85. String realUrl = String.format(url, server, appName);
      86. try {
      87. //客户端启动的时候连接一下服务端,看是否能连通
      88. String resultDTOStr = CommonUtils.executeWithRetry0(() -> HttpUtils.get(realUrl));
      89. ResultDTO resultDTO = JsonUtils.parseObject(resultDTOStr, ResultDTO.class);
      90. if (resultDTO.isSuccess()) {
      91. Long appId = Long.valueOf(resultDTO.getData().toString());
      92. log.info("[PowerJobWorker] assert appName({}) succeed, the appId for this application is {}.", appName, appId);
      93. workerRuntime.setAppId(appId);
      94. return;
      95. } else {
      96. log.error("[PowerJobWorker] assert appName failed, this appName is invalid, please register the appName {} first.", appName);
      97. throw new PowerJobException(resultDTO.getMessage());
      98. }
      99. } catch (PowerJobException oe) {
      100. throw oe;
      101. } catch (Exception ignore) {
      102. log.warn("[PowerJobWorker] assert appName by url({}) failed, please check the server address.", realUrl);
      103. }
      104. }
      105. log.error("[PowerJobWorker] no available server in {}.", config.getServerAddress());
      106. throw new PowerJobException("no server available!");
      107. }
      108. /**
      109. * ExecutorManager:
      110. * 第41行代码处:
      111. */
      112. public ExecutorManager(PowerJobWorkerConfig workerConfig) {
      113. //可用cpu数
      114. final int availableProcessors = Runtime.getRuntime().availableProcessors();
      115. // 初始化定时线程池
      116. ThreadFactory coreThreadFactory = new ThreadFactoryBuilder().setNameFormat("powerjob-worker-core-%d").build();
      117. coreExecutor = new ScheduledThreadPoolExecutor(3, coreThreadFactory);
      118. ThreadFactory lightTaskReportFactory = new ThreadFactoryBuilder().setNameFormat("powerjob-worker-light-task-status-check-%d").build();
      119. // 都是 io 密集型任务
      120. lightweightTaskStatusCheckExecutor = new ScheduledThreadPoolExecutor(availableProcessors * 10, lightTaskReportFactory);
      121. ThreadFactory lightTaskExecuteFactory = new ThreadFactoryBuilder().setNameFormat("powerjob-worker-light-task-execute-%d").build();
      122. // 大部分任务都是 io 密集型
      123. lightweightTaskExecutorService = new ThreadPoolExecutor(availableProcessors * 10, availableProcessors * 10, 120L, TimeUnit.SECONDS,
      124. new ArrayBlockingQueue<>((workerConfig.getMaxLightweightTaskNum() * 2), true), lightTaskExecuteFactory, new ThreadPoolExecutor.AbortPolicy());
      125. }
      126. /**
      127. * ServerDiscoveryService:
      128. * 第68行代码处:
      129. */
      130. public void start(ScheduledExecutorService timingPool) {
      131. this.currentServerAddress = discovery();
      132. if (StringUtils.isEmpty(this.currentServerAddress) && !config.isEnableTestMode()) {
      133. throw new PowerJobException("can't find any available server, this worker has been quarantined.");
      134. }
      135. // 这里必须保证成功
      136. timingPool.scheduleAtFixedRate(() -> {
      137. try {
      138. this.currentServerAddress = discovery();
      139. } catch (Exception e) {
      140. log.error("[PowerDiscovery] fail to discovery server!", e);
      141. }
      142. }
      143. , 10, 10, TimeUnit.SECONDS);
      144. }
      145. /**
      146. * 第159行代码处和第166行代码处:
      147. */
      148. private String discovery() {
      149. if (ip2Address.isEmpty()) {
      150. config.getServerAddress().forEach(x -> ip2Address.put(x.split(":")[0], x));
      151. }
      152. String result = null;
      153. // 先对当前机器发起请求
      154. String currentServer = currentServerAddress;
      155. if (!StringUtils.isEmpty(currentServer)) {
      156. String ip = currentServer.split(":")[0];
      157. // 直接请求当前Server的HTTP服务,可以少一次网络开销,减轻Server负担
      158. String firstServerAddress = ip2Address.get(ip);
      159. if (firstServerAddress != null) {
      160. result = acquire(firstServerAddress);
      161. }
      162. }
      163. for (String httpServerAddress : config.getServerAddress()) {
      164. if (StringUtils.isEmpty(result)) {
      165. result = acquire(httpServerAddress);
      166. } else {
      167. break;
      168. }
      169. }
      170. if (StringUtils.isEmpty(result)) {
      171. log.warn("[PowerDiscovery] can't find any available server, this worker has been quarantined.");
      172. // 在 Server 高可用的前提下,连续失败多次,说明该节点与外界失联,Server已经将秒级任务转移到其他Worker,需要杀死本地的任务
      173. if (FAILED_COUNT++ > MAX_FAILED_COUNT) {
      174. log.warn("[PowerDiscovery] can't find any available server for 3 consecutive times, It's time to kill all frequent job in this worker.");
      175. List frequentInstanceIds = HeavyTaskTrackerManager.getAllFrequentTaskTrackerKeys();
      176. if (!CollectionUtils.isEmpty(frequentInstanceIds)) {
      177. frequentInstanceIds.forEach(instanceId -> {
      178. HeavyTaskTracker taskTracker = HeavyTaskTrackerManager.removeTaskTracker(instanceId);
      179. taskTracker.destroy();
      180. log.warn("[PowerDiscovery] kill frequent instance(instanceId={}) due to can't find any available server.", instanceId);
      181. });
      182. }
      183. FAILED_COUNT = 0;
      184. }
      185. return null;
      186. } else {
      187. // 重置失败次数
      188. FAILED_COUNT = 0;
      189. log.debug("[PowerDiscovery] current server is {}.", result);
      190. return result;
      191. }
      192. }
      193. /**
      194. * 第192行代码处和第198行代码处
      195. */
      196. @SuppressWarnings("rawtypes")
      197. private String acquire(String httpServerAddress) {
      198. String result = null;
      199. //构建url参数
      200. String url = buildServerDiscoveryUrl(httpServerAddress);
      201. try {
      202. //请求服务端ServerController.acquireServer方法
      203. result = CommonUtils.executeWithRetry0(() -> HttpUtils.get(url));
      204. } catch (Exception ignore) {
      205. }
      206. if (!StringUtils.isEmpty(result)) {
      207. try {
      208. ResultDTO resultDTO = JsonUtils.parseObject(result, ResultDTO.class);
      209. if (resultDTO.isSuccess()) {
      210. return resultDTO.getData().toString();
      211. }
      212. } catch (Exception ignore) {
      213. }
      214. }
      215. return null;
      216. }
      217. /**
      218. * ServerController:
      219. * 第241行代码处:
      220. */
      221. @GetMapping("/acquire")
      222. public ResultDTO acquireServer(ServerDiscoveryRequest request) {
      223. return ResultDTO.success(serverElectionService.elect(request));
      224. }
      225. /**
      226. * ServerElectionService:
      227. */
      228. public String elect(ServerDiscoveryRequest request) {
      229. if (!accurate()) {
      230. final String currentServer = request.getCurrentServer();
      231. // 如果是本机,就不需要查数据库那么复杂的操作了,直接返回成功
      232. Optional localProtocolInfoOpt = Optional.ofNullable(transportService.allProtocols().get(request.getProtocol()));
      233. //如果不是精确地请求,并且请求的参数中直接带有当前服务端的地址,则直接返回其地址即可
      234. if (localProtocolInfoOpt.isPresent() && localProtocolInfoOpt.get().getAddress().equals(currentServer)) {
      235. log.debug("[ServerElectionService] this server[{}] is worker's current server, skip check", currentServer);
      236. return currentServer;
      237. }
      238. }
      239. //上面的条件不满足,则走常规的选举流程
      240. return getServer0(request);
      241. }
      242. /**
      243. * 第269行代码处:
      244. */
      245. private boolean accurate() {
      246. //accurateSelectServerPercentage默认为50,这里是在随机判断是否是精确的请求
      247. return ThreadLocalRandom.current().nextInt(100) < accurateSelectServerPercentage;
      248. }
      249. /**
      250. * 第280行代码处:
      251. */
      252. private String getServer0(ServerDiscoveryRequest discoveryRequest) {
      253. final Long appId = discoveryRequest.getAppId();
      254. final String protocol = discoveryRequest.getProtocol();
      255. Set downServerCache = Sets.newHashSet();
      256. for (int i = 0; i < RETRY_TIMES; i++) {
      257. // 无锁获取当前数据库中的Server
      258. Optional appInfoOpt = appInfoRepository.findById(appId);
      259. if (!appInfoOpt.isPresent()) {
      260. throw new PowerJobException(appId + " is not registered!");
      261. }
      262. String appName = appInfoOpt.get().getAppName();
      263. String originServer = appInfoOpt.get().getCurrentServer();
      264. String activeAddress = activeAddress(originServer, downServerCache, protocol);
      265. if (StringUtils.isNotEmpty(activeAddress)) {
      266. return activeAddress;
      267. }
      268. // 无可用Server,重新进行Server选举,需要加锁
      269. String lockName = String.format(SERVER_ELECT_LOCK, appId);
      270. //文件锁
      271. boolean lockStatus = lockService.tryLock(lockName, 30000);
      272. if (!lockStatus) {
      273. try {
      274. Thread.sleep(500);
      275. } catch (Exception ignore) {
      276. }
      277. continue;
      278. }
      279. try {
      280. // 可能上一台机器已经完成了Server选举,需要再次判断
      281. //双重检查加锁
      282. AppInfoDO appInfo = appInfoRepository.findById(appId).orElseThrow(() -> new RuntimeException("impossible, unless we just lost our database."));
      283. String address = activeAddress(appInfo.getCurrentServer(), downServerCache, protocol);
      284. if (StringUtils.isNotEmpty(address)) {
      285. return address;
      286. }
      287. // 篡位,如果本机存在协议,则作为Server调度该 worker
      288. //优先本机调度
      289. final ProtocolInfo targetProtocolInfo = transportService.allProtocols().get(protocol);
      290. if (targetProtocolInfo != null) {
      291. // 注意,写入 AppInfoDO#currentServer 的永远是 default 的地址,仅在返回的时候特殊处理为协议地址
      292. appInfo.setCurrentServer(transportService.defaultProtocol().getAddress());
      293. appInfo.setGmtModified(new Date());
      294. appInfoRepository.saveAndFlush(appInfo);
      295. log.info("[ServerElection] this server({}) become the new server for app(appId={}).", appInfo.getCurrentServer(), appId);
      296. return targetProtocolInfo.getAddress();
      297. }
      298. } catch (Exception e) {
      299. log.error("[ServerElection] write new server to db failed for app {}.", appName, e);
      300. } finally {
      301. lockService.unlock(lockName);
      302. }
      303. }
      304. throw new PowerJobException("server elect failed for app " + appId);
      305. }
      306. /**
      307. * 第309行代码处:
      308. * 判断指定server是否存活
      309. *
      310. * @param serverAddress 需要检测的server地址
      311. * @param downServerCache 缓存,防止多次发送PING(这个QPS其实还蛮爆表的...)
      312. * @param protocol 协议,用于返回指定的地址
      313. * @return null or address
      314. */
      315. private String activeAddress(String serverAddress, Set downServerCache, String protocol) {
      316. if (downServerCache.contains(serverAddress)) {
      317. return null;
      318. }
      319. if (StringUtils.isEmpty(serverAddress)) {
      320. return null;
      321. }
      322. Ping ping = new Ping();
      323. ping.setCurrentTime(System.currentTimeMillis());
      324. //构建服务端url参数
      325. URL targetUrl = ServerURLFactory.ping2Friend(serverAddress);
      326. try {
      327. //这里会跳转到FriendActor.onReceivePing方法
      328. AskResponse response = transportService.ask(Protocol.HTTP.name(), targetUrl, ping, AskResponse.class)
      329. .toCompletableFuture()
      330. .get(PING_TIMEOUT_MS, TimeUnit.MILLISECONDS);
      331. if (response.isSuccess()) {
      332. // 检测通过的是远程 server 的暴露地址,需要返回 worker 需要的协议地址
      333. final JSONObject protocolInfo = JsonUtils.parseObject(response.getData(), JSONObject.class).getJSONObject(protocol);
      334. if (protocolInfo != null) {
      335. downServerCache.remove(serverAddress);
      336. final String protocolAddress = protocolInfo.toJavaObject(ProtocolInfo.class).getAddress();
      337. log.info("[ServerElection] server[{}] is active, it will be the master, final protocol address={}", serverAddress, protocolAddress);
      338. return protocolAddress;
      339. } else {
      340. log.warn("[ServerElection] server[{}] is active but don't have target protocol", serverAddress);
      341. }
      342. }
      343. } catch (TimeoutException te) {
      344. log.warn("[ServerElection] server[{}] was down due to ping timeout!", serverAddress);
      345. } catch (Exception e) {
      346. log.warn("[ServerElection] server[{}] was down with unknown case!", serverAddress, e);
      347. }
      348. downServerCache.add(serverAddress);
      349. return null;
      350. }
      351. /**
      352. * FriendActor:
      353. * 第381行代码处:
      354. * 处理存活检测的请求
      355. */
      356. @Handler(path = S4S_HANDLER_PING, processType = ProcessType.NO_BLOCKING)
      357. public AskResponse onReceivePing(Ping ping) {
      358. return AskResponse.succeed(transportService.allProtocols());
      359. }
      360. /**
      361. * PowerTransportService:
      362. */
      363. @Override
      364. public Map allProtocols() {
      365. //这里就是从protocolName2Info缓存中取数据
      366. return protocolName2Info;
      367. }

              由上可知,在客户端启动的时候,会通过选举的方式来选出一台服务器作为自己的server(也可能不选举,直接选择当前的服务端),并且会定时发送心跳数据来进行保活。        

              在上面第85行代码处和第86行代码处,客户端启动的时候也会有两个定时任务,查看其实现:

      3.1.1 WorkerHealthReporter

              WorkerHealthReporter是用来给客户端健康度定时上报用的:

      1. /**
      2. * WorkerHealthReporter:
      3. */
      4. @Override
      5. public void run() {
      6. // 没有可用Server,无法上报
      7. String currentServer = workerRuntime.getServerDiscoveryService().getCurrentServerAddress();
      8. if (StringUtils.isEmpty(currentServer)) {
      9. log.warn("[WorkerHealthReporter] no available server,fail to report health info!");
      10. return;
      11. }
      12. SystemMetrics systemMetrics;
      13. if (workerRuntime.getWorkerConfig().getSystemMetricsCollector() == null) {
      14. systemMetrics = SystemInfoUtils.getSystemMetrics();
      15. } else {
      16. systemMetrics = workerRuntime.getWorkerConfig().getSystemMetricsCollector().collect();
      17. }
      18. WorkerHeartbeat heartbeat = new WorkerHeartbeat();
      19. heartbeat.setSystemMetrics(systemMetrics);
      20. heartbeat.setWorkerAddress(workerRuntime.getWorkerAddress());
      21. heartbeat.setAppName(workerRuntime.getWorkerConfig().getAppName());
      22. heartbeat.setAppId(workerRuntime.getAppId());
      23. heartbeat.setHeartbeatTime(System.currentTimeMillis());
      24. heartbeat.setVersion(PowerJobWorkerVersion.getVersion());
      25. heartbeat.setProtocol(workerRuntime.getWorkerConfig().getProtocol().name());
      26. heartbeat.setClient("KingPenguin");
      27. heartbeat.setTag(workerRuntime.getWorkerConfig().getTag());
      28. // 上报 Tracker 数量
      29. heartbeat.setLightTaskTrackerNum(LightTaskTrackerManager.currentTaskTrackerSize());
      30. heartbeat.setHeavyTaskTrackerNum(HeavyTaskTrackerManager.currentTaskTrackerSize());
      31. // 是否超载
      32. if (workerRuntime.getWorkerConfig().getMaxLightweightTaskNum() <= LightTaskTrackerManager.currentTaskTrackerSize() || workerRuntime.getWorkerConfig().getMaxHeavyweightTaskNum() <= HeavyTaskTrackerManager.currentTaskTrackerSize()) {
      33. heartbeat.setOverload(true);
      34. }
      35. // 获取当前加载的容器列表
      36. heartbeat.setContainerInfos(OmsContainerFactory.getDeployedContainerInfos());
      37. // 发送请求
      38. if (StringUtils.isEmpty(currentServer)) {
      39. return;
      40. }
      41. // log
      42. log.info("[WorkerHealthReporter] report health status,appId:{},appName:{},isOverload:{},maxLightweightTaskNum:{},currentLightweightTaskNum:{},maxHeavyweightTaskNum:{},currentHeavyweightTaskNum:{}",
      43. heartbeat.getAppId(),
      44. heartbeat.getAppName(),
      45. heartbeat.isOverload(),
      46. workerRuntime.getWorkerConfig().getMaxLightweightTaskNum(),
      47. heartbeat.getLightTaskTrackerNum(),
      48. workerRuntime.getWorkerConfig().getMaxHeavyweightTaskNum(),
      49. heartbeat.getHeavyTaskTrackerNum()
      50. );
      51. TransportUtils.reportWorkerHeartbeat(heartbeat, currentServer, workerRuntime.getTransporter());
      52. }
      53. /**
      54. * SystemInfoUtils:
      55. * 第17行代码处:
      56. */
      57. public static SystemMetrics getSystemMetrics() {
      58. SystemMetrics metrics = new SystemMetrics();
      59. //赋值cpu指标
      60. fillCPUInfo(metrics);
      61. //赋值内存指标
      62. fillMemoryInfo(metrics);
      63. //赋值磁盘指标
      64. fillDiskInfo(metrics);
      65. // 在Worker完成分数计算,减小Server压力
      66. metrics.calculateScore();
      67. return metrics;
      68. }
      69. /**
      70. * TransportUtils:
      71. * 第58行代码处:
      72. */
      73. public static void reportWorkerHeartbeat(WorkerHeartbeat req, String address, Transporter transporter) {
      74. //绑定url调用信息
      75. final URL url = easyBuildUrl(ServerType.SERVER, S4W_PATH, S4W_HANDLER_WORKER_HEARTBEAT, address);
      76. transporter.tell(url, req);
      77. }

              其中SystemInfoUtils.getSystemMetrics方法是在获取系统的一些指标数据(调用的都是Java的底层Api,这里就不再详细查看了),并且计算健康度的得分(之前在第2.3.1小节中已经查看过该方法的实现了)。基于此,我们就知道了控制台首页的worker指标数据是怎么来的了:

              在上面第88行代码处向服务端发送了心跳数据,接下来就来看下服务端是如何处理的:

      1. /**
      2. * AbWorkerRequestHandler:
      3. */
      4. @Override
      5. @Handler(path = S4W_HANDLER_WORKER_HEARTBEAT, processType = ProcessType.NO_BLOCKING)
      6. public void processWorkerHeartbeat(WorkerHeartbeat heartbeat) {
      7. long startMs = System.currentTimeMillis();
      8. WorkerHeartbeatEvent event = new WorkerHeartbeatEvent()
      9. .setAppName(heartbeat.getAppName())
      10. .setAppId(heartbeat.getAppId())
      11. .setVersion(heartbeat.getVersion())
      12. .setProtocol(heartbeat.getProtocol())
      13. .setTag(heartbeat.getTag())
      14. .setWorkerAddress(heartbeat.getWorkerAddress())
      15. .setDelayMs(startMs - heartbeat.getHeartbeatTime())
      16. .setScore(heartbeat.getSystemMetrics().getScore());
      17. processWorkerHeartbeat0(heartbeat, event);
      18. //默认实现是写入日志进行监控
      19. monitorService.monitor(event);
      20. }
      21. /**
      22. * WorkerRequestHandlerImpl:
      23. * 第17行代码处:
      24. */
      25. @Override
      26. protected void processWorkerHeartbeat0(WorkerHeartbeat heartbeat, WorkerHeartbeatEvent event) {
      27. WorkerClusterManagerService.updateStatus(heartbeat);
      28. }
      29. /**
      30. * WorkerClusterManagerService:
      31. * 更新状态
      32. *
      33. * @param heartbeat Worker的心跳包
      34. */
      35. public static void updateStatus(WorkerHeartbeat heartbeat) {
      36. Long appId = heartbeat.getAppId();
      37. String appName = heartbeat.getAppName();
      38. ClusterStatusHolder clusterStatusHolder = APP_ID_2_CLUSTER_STATUS.computeIfAbsent(appId, ignore -> new ClusterStatusHolder(appName));
      39. clusterStatusHolder.updateStatus(heartbeat);
      40. }
      41. /**
      42. * ClusterStatusHolder:
      43. * 更新 worker 机器的状态
      44. *
      45. * @param heartbeat 心跳请求
      46. */
      47. public void updateStatus(WorkerHeartbeat heartbeat) {
      48. String workerAddress = heartbeat.getWorkerAddress();
      49. long heartbeatTime = heartbeat.getHeartbeatTime();
      50. WorkerInfo workerInfo = address2WorkerInfo.computeIfAbsent(workerAddress, ignore -> {
      51. WorkerInfo wf = new WorkerInfo();
      52. wf.refresh(heartbeat);
      53. return wf;
      54. });
      55. long oldTime = workerInfo.getLastActiveTime();
      56. //过期心跳数据,不处理
      57. if (heartbeatTime < oldTime) {
      58. log.warn("[ClusterStatusHolder-{}] receive the expired heartbeat from {}, serverTime: {}, heartTime: {}", appName, heartbeat.getWorkerAddress(), System.currentTimeMillis(), heartbeat.getHeartbeatTime());
      59. return;
      60. }
      61. workerInfo.refresh(heartbeat);
      62. List containerInfos = heartbeat.getContainerInfos();
      63. if (!CollectionUtils.isEmpty(containerInfos)) {
      64. containerInfos.forEach(containerInfo -> {
      65. Map infos = containerId2Infos.computeIfAbsent(containerInfo.getContainerId(), ignore -> Maps.newConcurrentMap());
      66. infos.put(workerAddress, containerInfo);
      67. });
      68. }
      69. }
      70. /**
      71. * WorkerInfo:
      72. * 第57行代码处和第67行代码处:
      73. * 刷新服务端记录的客户端的数据
      74. */
      75. public void refresh(WorkerHeartbeat workerHeartbeat) {
      76. address = workerHeartbeat.getWorkerAddress();
      77. lastActiveTime = workerHeartbeat.getHeartbeatTime();
      78. protocol = workerHeartbeat.getProtocol();
      79. client = workerHeartbeat.getClient();
      80. tag = workerHeartbeat.getTag();
      81. systemMetrics = workerHeartbeat.getSystemMetrics();
      82. containerInfos = workerHeartbeat.getContainerInfos();
      83. lightTaskTrackerNum = workerHeartbeat.getLightTaskTrackerNum();
      84. heavyTaskTrackerNum = workerHeartbeat.getHeavyTaskTrackerNum();
      85. if (workerHeartbeat.isOverload()) {
      86. overloading = true;
      87. lastOverloadTime = workerHeartbeat.getHeartbeatTime();
      88. log.warn("[WorkerInfo] worker {} is overload!", getAddress());
      89. } else {
      90. overloading = false;
      91. }
      92. }

      3.1.2 LogSubmitter

      1. /**
      2. * LogSubmitter:
      3. */
      4. @Override
      5. public void run() {
      6. boolean lockResult = reportLock.tryLock();
      7. if (!lockResult) {
      8. return;
      9. }
      10. try {
      11. final String currentServerAddress = serverDiscoveryService.getCurrentServerAddress();
      12. // 当前无可用 Server
      13. if (StringUtils.isEmpty(currentServerAddress)) {
      14. if (!logQueue.isEmpty()) {
      15. logQueue.clear();
      16. log.warn("[OmsLogHandler] because there is no available server to report logs which leads to queue accumulation, oms discarded all logs.");
      17. }
      18. return;
      19. }
      20. List logs = Lists.newLinkedList();
      21. while (!logQueue.isEmpty()) {
      22. try {
      23. //从日志队列头部取数据
      24. InstanceLogContent logContent = logQueue.poll(100, TimeUnit.MILLISECONDS);
      25. logs.add(logContent);
      26. //批处理
      27. if (logs.size() >= BATCH_SIZE) {
      28. WorkerLogReportReq req = new WorkerLogReportReq(workerAddress, Lists.newLinkedList(logs));
      29. // 不可靠请求,WEB日志不追求极致
      30. TransportUtils.reportLogs(req, currentServerAddress, transporter);
      31. logs.clear();
      32. }
      33. } catch (Exception ignore) {
      34. break;
      35. }
      36. }
      37. if (!logs.isEmpty()) {
      38. WorkerLogReportReq req = new WorkerLogReportReq(workerAddress, logs);
      39. TransportUtils.reportLogs(req, currentServerAddress, transporter);
      40. }
      41. } finally {
      42. reportLock.unlock();
      43. }
      44. }
      45. /**
      46. * TransportUtils:
      47. * 第36行代码处和第47行代码处:
      48. */
      49. public static void reportLogs(WorkerLogReportReq req, String address, Transporter transporter) {
      50. final URL url = easyBuildUrl(ServerType.SERVER, S4W_PATH, S4W_HANDLER_REPORT_LOG, address);
      51. transporter.tell(url, req);
      52. }

              上面第61行代码处会将请求发送到服务端的AbWorkerRequestHandler.processWorkerLogReport方法中:

      1. /**
      2. * AbWorkerRequestHandler:
      3. */
      4. @Override
      5. @Handler(path = S4W_HANDLER_REPORT_LOG, processType = ProcessType.NO_BLOCKING)
      6. public void processWorkerLogReport(WorkerLogReportReq req) {
      7. WorkerLogReportEvent event = new WorkerLogReportEvent()
      8. .setWorkerAddress(req.getWorkerAddress())
      9. .setLogNum(req.getInstanceLogContents().size());
      10. try {
      11. processWorkerLogReport0(req, event);
      12. event.setStatus(WorkerLogReportEvent.Status.SUCCESS);
      13. } catch (RejectedExecutionException re) {
      14. event.setStatus(WorkerLogReportEvent.Status.REJECTED);
      15. } catch (Throwable t) {
      16. event.setStatus(WorkerLogReportEvent.Status.EXCEPTION);
      17. log.warn("[WorkerRequestHandler] process worker report failed!", t);
      18. } finally {
      19. //日志监控
      20. monitorService.monitor(event);
      21. }
      22. }
      23. /**
      24. * WorkerRequestHandlerImpl:
      25. * 第12行代码处:
      26. */
      27. @Override
      28. protected void processWorkerLogReport0(WorkerLogReportReq req, WorkerLogReportEvent event) {
      29. // 这个效率应该不会拉垮吧...也就是一些判断 + Map#get 吧...
      30. instanceLogService.submitLogs(req.getWorkerAddress(), req.getInstanceLogContents());
      31. }
      32. /**
      33. * InstanceLogService:
      34. * 提交日志记录,持久化到本地数据库中
      35. *
      36. * @param workerAddress 上报机器地址
      37. * @param logs 任务实例运行时日志
      38. */
      39. @Async(value = PJThreadPool.LOCAL_DB_POOL)
      40. public void submitLogs(String workerAddress, List logs) {
      41. List logList = logs.stream().map(x -> {
      42. instanceId2LastReportTime.put(x.getInstanceId(), System.currentTimeMillis());
      43. LocalInstanceLogDO y = new LocalInstanceLogDO();
      44. BeanUtils.copyProperties(x, y);
      45. y.setWorkerAddress(workerAddress);
      46. return y;
      47. }).collect(Collectors.toList());
      48. try {
      49. //插入到local_instance_log表中
      50. CommonUtils.executeWithRetry0(() -> localInstanceLogRepository.saveAll(logList));
      51. } catch (Exception e) {
      52. log.warn("[InstanceLogService] persistent instance logs failed, these logs will be dropped: {}.", logs, e);
      53. }
      54. }

      原创不易,未得准许,请勿转载,翻版必究

    161. 相关阅读:
      Java版工程行业管理系统源码-专业的工程管理软件- 工程项目各模块及其功能点清单
      javaee SpringMVC 乱码问题解决
      docker如何查看对外暴露接口
      人工智能:人脸识别技术应用场景介绍
      20.1CubeMx配置FMC控制SDRAM【W9825G6KH-6】
      day-45 代码随想录算法训练营(19)动态规划 part 07
      SpringBoot自动配置原理解析 | 京东物流技术团队
      static关键字修饰成员变量与成员函数
      CI/CD持续集成/持续部署
      超详细Redis使用手册
    162. 原文地址:https://blog.csdn.net/weixin_30342639/article/details/132662425