• xxl-job源码—调度器/执行器工作原理


    目录

    一、架构图

    1.1 功能架构图

    2.2 任务调度工作原理

    二、ER图

    三、调度器

    3.1 启动过程时序图

    3.2 启动过程核心代码解析

    3.2.1 启动初始化

    3.2.2 执行器健康检查

    3.2.3 任务执行失败告警

    3.2.4 调度线程池

    3.2.5 调度中心

    3.3 任务执行时序图

    3.4 任务执行核心代码解析

    3.4.1 根据路由策略定位到执行器

    3.4.2 定位到任务执行Handler&线程

    3.4.3 任务执行线程

    四、执行器

    4.1 启动过程时序图

    4.2 启动过程核心代码解析

    4.2.1 启动初始化

    4.2.2 执行器初始化

    4.2.3 任务执行回调

    4.3 执行器注册到调度中心时序图

    4.4 执行器注册到调度中心核心代码解析

    4.4.1 执行注册(30秒一次)


    一、架构图

    1.1 功能架构图

    2.2 任务调度工作原理

    二、ER图

    三、调度器

    3.1 启动过程时序图

    3.2 启动过程核心代码解析

    入口类:XxlJobAdminConfig(PS:启动初始化XxlJobScheduler

    3.2.1 启动初始化

    调度器启动初始化:XxlJobScheduler

    1. package com.xxl.job.admin.core.scheduler;
    2. import com.xxl.job.admin.core.conf.XxlJobAdminConfig;
    3. import com.xxl.job.admin.core.thread.*;
    4. import com.xxl.job.admin.core.util.I18nUtil;
    5. import com.xxl.job.core.biz.ExecutorBiz;
    6. import com.xxl.job.core.enums.ExecutorBlockStrategyEnum;
    7. import com.xxl.rpc.remoting.invoker.call.CallType;
    8. import com.xxl.rpc.remoting.invoker.reference.XxlRpcReferenceBean;
    9. import com.xxl.rpc.remoting.invoker.route.LoadBalance;
    10. import com.xxl.rpc.remoting.net.impl.netty_http.client.NettyHttpClient;
    11. import com.xxl.rpc.serialize.impl.HessianSerializer;
    12. import org.slf4j.Logger;
    13. import org.slf4j.LoggerFactory;
    14. import java.util.concurrent.ConcurrentHashMap;
    15. import java.util.concurrent.ConcurrentMap;
    16. /**
    17. * @author xuxueli 2018-10-28 00:18:17
    18. * Description:调度程序
    19. */
    20. public class XxlJobScheduler {
    21. private static final Logger logger = LoggerFactory.getLogger(XxlJobScheduler.class);
    22. public void init() throws Exception {
    23. // init i18n
    24. /**
    25. * 主要作用:调度中心语言选择
    26. */
    27. initI18n();
    28. // admin registry monitor run
    29. /**
    30. * 主要作用是开通了一个守护线程,每隔30s扫描一次执行器的注册信息表
    31. * 1、剔除90s内没有进行健康检查的执行器信息
    32. * 2、将自动注册类型的执行器注册信息(XxlJobRegistry)经过处理更新执行器信息(XxlJobGroup)
    33. */
    34. JobRegistryMonitorHelper.getInstance().start();
    35. // admin monitor run
    36. /**
    37. * 主要作用:开通了一个守10s扫描一次失败日志
    38. * 1、如果任务失败可重试次数>0,那么重新触发任务
    39. * 2、如果任务执行失败,会进行告警,默认采用邮件形式进行告警
    40. */
    41. JobFailMonitorHelper.getInstance().start();
    42. // admin trigger pool start
    43. /**
    44. * 主要作用:初始化一个快速执行的线程池,一个稍后执行的线程池
    45. *
    46. */
    47. JobTriggerPoolHelper.toStart();
    48. // admin log report start
    49. /**
    50. * 主要作用:开通一个守护线程,每隔1min扫描一次最近3天的调度日志
    51. * 1、更新每天总任务数、正在执行数、执行成功数、执行失败数
    52. */
    53. JobLogReportHelper.getInstance().start();
    54. // start-schedule
    55. /**
    56. * 主要作用:
    57. * 1、开通一个守护线程,每隔5s扫表一次执行器的任务表
    58. * 1.1 执行已到执行时间的任务,并更新任务的下次执行时间
    59. *
    60. * 2、开通一个守护线程,每隔1s循环执行一次待执行的任务
    61. * 2.1 避免任务执行遗漏
    62. */
    63. JobScheduleHelper.getInstance().start();
    64. logger.info(">>>>>>>>> init xxl-job admin success.");
    65. }
    66. public void destroy() throws Exception {
    67. // stop-schedule
    68. JobScheduleHelper.getInstance().toStop();
    69. // admin log report stop
    70. JobLogReportHelper.getInstance().toStop();
    71. // admin trigger pool stop
    72. JobTriggerPoolHelper.toStop();
    73. // admin monitor stop
    74. JobFailMonitorHelper.getInstance().toStop();
    75. // admin registry stop
    76. JobRegistryMonitorHelper.getInstance().toStop();
    77. }
    78. // ---------------------- I18n ----------------------
    79. private void initI18n(){
    80. for (ExecutorBlockStrategyEnum item:ExecutorBlockStrategyEnum.values()) {
    81. // 根据语言选择,去对应的message_xxx.properties中取阻塞策略的title
    82. item.setTitle(I18nUtil.getString("jobconf_block_".concat(item.name())));
    83. }
    84. }
    85. // ---------------------- executor-client ----------------------
    86. private static ConcurrentMap executorBizRepository = new ConcurrentHashMap();
    87. public static ExecutorBiz getExecutorBiz(String address) throws Exception {
    88. // valid
    89. if (address==null || address.trim().length()==0) {
    90. return null;
    91. }
    92. // load-cache
    93. address = address.trim();
    94. ExecutorBiz executorBiz = executorBizRepository.get(address);
    95. if (executorBiz != null) {
    96. return executorBiz;
    97. }
    98. // set-cache
    99. // 创建ExecutorBiz的代理对象,重点在这个getObject方法
    100. XxlRpcReferenceBean referenceBean = new XxlRpcReferenceBean();
    101. referenceBean.setClient(NettyHttpClient.class);
    102. referenceBean.setSerializer(HessianSerializer.class);
    103. referenceBean.setCallType(CallType.SYNC);
    104. referenceBean.setLoadBalance(LoadBalance.ROUND);
    105. referenceBean.setIface(ExecutorBiz.class);
    106. referenceBean.setVersion(null);
    107. referenceBean.setTimeout(3000);
    108. referenceBean.setAddress(address);
    109. referenceBean.setAccessToken(XxlJobAdminConfig.getAdminConfig().getAccessToken());
    110. referenceBean.setInvokeCallback(null);
    111. referenceBean.setInvokerFactory(null);
    112. executorBiz = (ExecutorBiz) referenceBean.getObject();
    113. executorBizRepository.put(address, executorBiz);
    114. return executorBiz;
    115. }
    116. }

    3.2.2 执行器健康检查

    执行器注册健康检查:JobRegistryMonitorHelper

    1. package com.xxl.job.admin.core.thread;
    2. import com.xxl.job.admin.core.conf.XxlJobAdminConfig;
    3. import com.xxl.job.admin.core.model.XxlJobGroup;
    4. import com.xxl.job.admin.core.model.XxlJobRegistry;
    5. import com.xxl.job.core.enums.RegistryConfig;
    6. import org.slf4j.Logger;
    7. import org.slf4j.LoggerFactory;
    8. import java.util.*;
    9. import java.util.concurrent.TimeUnit;
    10. /**
    11. * job registry instance
    12. * @author xuxueli 2016-10-02 19:10:24
    13. * @Description:任务注册监控助手
    14. */
    15. public class JobRegistryMonitorHelper {
    16. private static Logger logger = LoggerFactory.getLogger(JobRegistryMonitorHelper.class);
    17. private static JobRegistryMonitorHelper instance = new JobRegistryMonitorHelper();
    18. public static JobRegistryMonitorHelper getInstance(){
    19. return instance;
    20. }
    21. private Thread registryThread;
    22. private volatile boolean toStop = false;
    23. public void start(){
    24. //创建一个线程
    25. registryThread = new Thread(new Runnable() {
    26. @Override
    27. public void run() {
    28. //当toStop为false时进入该循环(注意:toStop是用volatile修饰的)
    29. while (!toStop) {
    30. try {
    31. // auto registry group
    32. //获取类型为自动注册的执行器(XxlJobGroup)地址列表
    33. List groupList = XxlJobAdminConfig.getAdminConfig().getXxlJobGroupDao().findByAddressType(0);
    34. if (groupList!=null && !groupList.isEmpty()) {
    35. // remove dead address (admin/executor)
    36. //删除90秒内没有更新的注册机器信息,90秒没有心跳信息返回代表机器已经出现问题,所以移除
    37. List ids = XxlJobAdminConfig.getAdminConfig().getXxlJobRegistryDao().findDead(RegistryConfig.DEAD_TIMEOUT, new Date());
    38. if (ids!=null && ids.size()>0) {
    39. XxlJobAdminConfig.getAdminConfig().getXxlJobRegistryDao().removeDead(ids);
    40. }
    41. // fresh online address (admin/executor)
    42. HashMap> appAddressMap = new HashMap>();
    43. //查询90秒内有更新的注册机器信息列表
    44. List list = XxlJobAdminConfig.getAdminConfig().getXxlJobRegistryDao().findAll(RegistryConfig.DEAD_TIMEOUT, new Date());
    45. if (list != null) {
    46. //遍历注册信息列表,得到自动注册类型的执行器与其对应的地址信息关系Map
    47. for (XxlJobRegistry item: list) {
    48. if (RegistryConfig.RegistType.EXECUTOR.name().equals(item.getRegistryGroup())) {
    49. String appName = item.getRegistryKey();
    50. List registryList = appAddressMap.get(appName);
    51. if (registryList == null) {
    52. registryList = new ArrayList();
    53. }
    54. if (!registryList.contains(item.getRegistryValue())) {
    55. registryList.add(item.getRegistryValue());
    56. }
    57. //收集执行器信息,根据执行器appName做区分:appAddressMap的key为appName,value为此执行器的注册地址列表(集群环境下会有多个注册地址)
    58. appAddressMap.put(appName, registryList);
    59. }
    60. }
    61. }
    62. // fresh group address
    63. //遍历所有的自动注册的执行器
    64. for (XxlJobGroup group: groupList) {
    65. //通过执行器的appName从刚刚区分的Map中拿到该执行器下的集群机器注册地址
    66. List registryList = appAddressMap.get(group.getAppName());
    67. String addressListStr = null;
    68. if (registryList!=null && !registryList.isEmpty()) {
    69. Collections.sort(registryList);
    70. addressListStr = "";
    71. //集群的多个注册地址通过逗号拼接转为字符串
    72. for (String item:registryList) {
    73. addressListStr += item + ",";
    74. }
    75. addressListStr = addressListStr.substring(0, addressListStr.length()-1);
    76. }
    77. //集群的多个注册地址通过逗号拼接转成的字符串设置进执行器的addressList属性中,并更新执行器信息,保存到DB
    78. group.setAddressList(addressListStr);
    79. XxlJobAdminConfig.getAdminConfig().getXxlJobGroupDao().update(group);
    80. }
    81. }
    82. } catch (Exception e) {
    83. if (!toStop) {
    84. logger.error(">>>>>>>>>>> xxl-job, job registry monitor thread error:{}", e);
    85. }
    86. }
    87. try {
    88. //线程停顿30秒
    89. TimeUnit.SECONDS.sleep(RegistryConfig.BEAT_TIMEOUT);
    90. } catch (InterruptedException e) {
    91. if (!toStop) {
    92. logger.error(">>>>>>>>>>> xxl-job, job registry monitor thread error:{}", e);
    93. }
    94. }
    95. }
    96. logger.info(">>>>>>>>>>> xxl-job, job registry monitor thread stop");
    97. }
    98. });
    99. //将此线程设置成守护线程
    100. registryThread.setDaemon(true);
    101. registryThread.setName("xxl-job, admin JobRegistryMonitorHelper");
    102. //执行该线程
    103. registryThread.start();
    104. }
    105. public void toStop(){
    106. toStop = true;
    107. // interrupt and wait
    108. registryThread.interrupt();
    109. try {
    110. registryThread.join();
    111. } catch (InterruptedException e) {
    112. logger.error(e.getMessage(), e);
    113. }
    114. }
    115. }

    3.2.3 任务执行失败告警

    任务失败监控助手:JobFailMonitorHelper

    1. package com.xxl.job.admin.core.thread;
    2. import com.xxl.job.admin.core.conf.XxlJobAdminConfig;
    3. import com.xxl.job.admin.core.model.XxlJobInfo;
    4. import com.xxl.job.admin.core.model.XxlJobLog;
    5. import com.xxl.job.admin.core.trigger.TriggerTypeEnum;
    6. import com.xxl.job.admin.core.util.I18nUtil;
    7. import org.slf4j.Logger;
    8. import org.slf4j.LoggerFactory;
    9. import java.util.List;
    10. import java.util.concurrent.TimeUnit;
    11. /**
    12. * job monitor instance
    13. *
    14. * @author xuxueli 2015-9-1 18:05:56
    15. * @Description:任务失败监控助手
    16. */
    17. public class JobFailMonitorHelper {
    18. private static Logger logger = LoggerFactory.getLogger(JobFailMonitorHelper.class);
    19. private static JobFailMonitorHelper instance = new JobFailMonitorHelper();
    20. public static JobFailMonitorHelper getInstance(){
    21. return instance;
    22. }
    23. // ---------------------- monitor ----------------------
    24. private Thread monitorThread;
    25. private volatile boolean toStop = false;
    26. public void start(){
    27. //创建一个监控线程
    28. monitorThread = new Thread(new Runnable() {
    29. @Override
    30. public void run() {
    31. // monitor
    32. while (!toStop) {
    33. try {
    34. //从数据库中取出所有执行失败且告警状态为0(默认)的日志ID,虽然这里传了pageSize为1000,实际没有进行分页,取出来的是所有符合条件的失败日志
    35. List failLogIds = XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().findFailJobLogIds(1000);
    36. if (failLogIds!=null && !failLogIds.isEmpty()) {
    37. for (long failLogId: failLogIds) {
    38. // 锁定日志
    39. int lockRet = XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().updateAlarmStatus(failLogId, 0, -1);
    40. if (lockRet < 1) {
    41. continue;
    42. }
    43. //取出失败日志完整信息
    44. XxlJobLog log = XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().load(failLogId);
    45. //获取失败日志对应的任务信息
    46. XxlJobInfo info = XxlJobAdminConfig.getAdminConfig().getXxlJobInfoDao().loadById(log.getJobId());
    47. // 如果剩余失败可重试次数>0(注意:日志里存的失败重试次数实为剩余可重复次数)
    48. if (log.getExecutorFailRetryCount() > 0) {
    49. //触发任务执行
    50. JobTriggerPoolHelper.trigger(log.getJobId(), TriggerTypeEnum.RETRY, (log.getExecutorFailRetryCount()-1), log.getExecutorShardingParam(), log.getExecutorParam());
    51. String retryMsg = "

      >>>>>>>>>>>"+ I18nUtil.getString("jobconf_trigger_type_retry") +"<<<<<<<<<<<
      "
      ;
    52. log.setTriggerMsg(log.getTriggerMsg() + retryMsg);
    53. XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().updateTriggerInfo(log);
    54. }
    55. // 2、fail alarm monitor
    56. //失败告警
    57. int newAlarmStatus = 0; // 告警状态:0-默认、-1=锁定状态、1-无需告警、2-告警成功、3-告警失败
    58. if (info!=null && info.getAlarmEmail()!=null && info.getAlarmEmail().trim().length()>0) {
    59. //发送失败告警:默认是邮件通知
    60. boolean alarmResult = XxlJobAdminConfig.getAdminConfig().getJobAlarmer().alarm(info, log);
    61. newAlarmStatus = alarmResult?2:3;
    62. } else {
    63. newAlarmStatus = 1;
    64. }
    65. //更新告警状态
    66. XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().updateAlarmStatus(failLogId, -1, newAlarmStatus);
    67. }
    68. }
    69. } catch (Exception e) {
    70. if (!toStop) {
    71. logger.error(">>>>>>>>>>> xxl-job, job fail monitor thread error:{}", e);
    72. }
    73. }
    74. try {
    75. TimeUnit.SECONDS.sleep(10);
    76. } catch (Exception e) {
    77. if (!toStop) {
    78. logger.error(e.getMessage(), e);
    79. }
    80. }
    81. }
    82. logger.info(">>>>>>>>>>> xxl-job, job fail monitor thread stop");
    83. }
    84. });
    85. monitorThread.setDaemon(true);
    86. monitorThread.setName("xxl-job, admin JobFailMonitorHelper");
    87. monitorThread.start();
    88. }
    89. public void toStop(){
    90. toStop = true;
    91. // interrupt and wait
    92. monitorThread.interrupt();
    93. try {
    94. monitorThread.join();
    95. } catch (InterruptedException e) {
    96. logger.error(e.getMessage(), e);
    97. }
    98. }
    99. }

    3.2.4 调度线程池

    分配任务执行线程:JobTriggerPoolHelper

    1. package com.xxl.job.admin.core.thread;
    2. import com.xxl.job.admin.core.conf.XxlJobAdminConfig;
    3. import com.xxl.job.admin.core.trigger.TriggerTypeEnum;
    4. import com.xxl.job.admin.core.trigger.XxlJobTrigger;
    5. import org.slf4j.Logger;
    6. import org.slf4j.LoggerFactory;
    7. import java.util.concurrent.*;
    8. import java.util.concurrent.atomic.AtomicInteger;
    9. /**
    10. * job trigger thread pool helper
    11. *
    12. * @author xuxueli 2018-07-03 21:08:07
    13. */
    14. public class JobTriggerPoolHelper {
    15. private static Logger logger = LoggerFactory.getLogger(JobTriggerPoolHelper.class);
    16. // ---------------------- trigger pool ----------------------
    17. // fast/slow thread pool
    18. private ThreadPoolExecutor fastTriggerPool = null;
    19. private ThreadPoolExecutor slowTriggerPool = null;
    20. public void start(){
    21. fastTriggerPool = new ThreadPoolExecutor(
    22. 10,
    23. XxlJobAdminConfig.getAdminConfig().getTriggerPoolFastMax(),
    24. 60L,
    25. TimeUnit.SECONDS,
    26. new LinkedBlockingQueue(1000),
    27. new ThreadFactory() {
    28. @Override
    29. public Thread newThread(Runnable r) {
    30. return new Thread(r, "xxl-job, admin JobTriggerPoolHelper-fastTriggerPool-" + r.hashCode());
    31. }
    32. });
    33. slowTriggerPool = new ThreadPoolExecutor(
    34. 10,
    35. XxlJobAdminConfig.getAdminConfig().getTriggerPoolSlowMax(),
    36. 60L,
    37. TimeUnit.SECONDS,
    38. new LinkedBlockingQueue(2000),
    39. new ThreadFactory() {
    40. @Override
    41. public Thread newThread(Runnable r) {
    42. return new Thread(r, "xxl-job, admin JobTriggerPoolHelper-slowTriggerPool-" + r.hashCode());
    43. }
    44. });
    45. }
    46. public void stop() {
    47. //triggerPool.shutdown();
    48. fastTriggerPool.shutdownNow();
    49. slowTriggerPool.shutdownNow();
    50. logger.info(">>>>>>>>> xxl-job trigger thread pool shutdown success.");
    51. }
    52. // job timeout count
    53. private volatile long minTim = System.currentTimeMillis()/60000; // ms > min
    54. private volatile ConcurrentMap jobTimeoutCountMap = new ConcurrentHashMap<>();
    55. /**
    56. * add trigger
    57. */
    58. public void addTrigger(final int jobId, final TriggerTypeEnum triggerType, final int failRetryCount, final String executorShardingParam, final String executorParam) {
    59. // choose thread pool
    60. //选择线程池类型
    61. ThreadPoolExecutor triggerPool_ = fastTriggerPool;
    62. AtomicInteger jobTimeoutCount = jobTimeoutCountMap.get(jobId);
    63. // 1分钟窗口期内任务耗时达500ms超过10次则判定为慢任务,慢任务自动降级进入"Slow"线程池,避免耗尽调度线程,提高系统稳定性;
    64. if (jobTimeoutCount!=null && jobTimeoutCount.get() > 10) { // job-timeout 10 times in 1 min
    65. triggerPool_ = slowTriggerPool;
    66. }
    67. // trigger
    68. triggerPool_.execute(new Runnable() {
    69. @Override
    70. public void run() {
    71. long start = System.currentTimeMillis();
    72. try {
    73. // do trigger
    74. //这里是重点,触发任务的执行
    75. XxlJobTrigger.trigger(jobId, triggerType, failRetryCount, executorShardingParam, executorParam);
    76. } catch (Exception e) {
    77. logger.error(e.getMessage(), e);
    78. } finally {
    79. //在finally代码块中通过执行耗时时间是否>500ms将jobId与超时次数记录到到超时记录容器
    80. // check timeout-count-map
    81. //检查超时记录容器
    82. long minTim_now = System.currentTimeMillis()/60000;
    83. if (minTim != minTim_now) {
    84. //说明当前分钟数和上面的默认分钟数不同,即超过了1分钟
    85. //这里的逻辑主要是为了使得超时记录容器以1分钟为时间间隔,也就是说超时容器1分钟记录1次超时的任务,到了下一分钟清空容器,重新记录
    86. minTim = minTim_now;
    87. jobTimeoutCountMap.clear();
    88. }
    89. // incr timeout-count-map
    90. //如果线程执行时间>500毫秒:将此任务ID与次数记录进超时记录容器jobTimeoutCountMap
    91. long cost = System.currentTimeMillis()-start;
    92. if (cost > 500) { // ob-timeout threshold 500ms
    93. //这里利用AtomicInteger的原子性,保证多线程并发结果的准确性
    94. AtomicInteger timeoutCount = jobTimeoutCountMap.putIfAbsent(jobId, new AtomicInteger(1));
    95. if (timeoutCount != null) {
    96. timeoutCount.incrementAndGet();
    97. }
    98. }
    99. }
    100. }
    101. });
    102. }
    103. // ---------------------- helper ----------------------
    104. private static JobTriggerPoolHelper helper = new JobTriggerPoolHelper();
    105. public static void toStart() {
    106. helper.start();
    107. }
    108. public static void toStop() {
    109. helper.stop();
    110. }
    111. /**
    112. * @param jobId
    113. * @param triggerType
    114. * @param failRetryCount
    115. * >=0: use this param
    116. * <0: use param from job info config
    117. * @param executorShardingParam
    118. * @param executorParam
    119. * null: use job param
    120. * not null: cover job param
    121. */
    122. public static void trigger(int jobId, TriggerTypeEnum triggerType, int failRetryCount, String executorShardingParam, String executorParam) {
    123. helper.addTrigger(jobId, triggerType, failRetryCount, executorShardingParam, executorParam);
    124. }
    125. }

    3.2.5 调度中心

    核心调度程序:JobScheduleHelper

    1. package com.xxl.job.admin.core.thread;
    2. import com.xxl.job.admin.core.conf.XxlJobAdminConfig;
    3. import com.xxl.job.admin.core.cron.CronExpression;
    4. import com.xxl.job.admin.core.model.XxlJobInfo;
    5. import com.xxl.job.admin.core.trigger.TriggerTypeEnum;
    6. import org.slf4j.Logger;
    7. import org.slf4j.LoggerFactory;
    8. import java.sql.Connection;
    9. import java.sql.PreparedStatement;
    10. import java.sql.SQLException;
    11. import java.text.ParseException;
    12. import java.util.*;
    13. import java.util.concurrent.ConcurrentHashMap;
    14. import java.util.concurrent.TimeUnit;
    15. /**
    16. * @author xuxueli 2019-05-21
    17. * @Description:调度中心
    18. */
    19. public class JobScheduleHelper {
    20. private static Logger logger = LoggerFactory.getLogger(JobScheduleHelper.class);
    21. private static JobScheduleHelper instance = new JobScheduleHelper();
    22. public static JobScheduleHelper getInstance(){
    23. return instance;
    24. }
    25. public static final long PRE_READ_MS = 5000; // pre read
    26. private Thread scheduleThread;
    27. private Thread ringThread;
    28. private volatile boolean scheduleThreadToStop = false;
    29. private volatile boolean ringThreadToStop = false;
    30. private volatile static Map> ringData = new ConcurrentHashMap<>();
    31. public void start(){
    32. // schedule thread
    33. scheduleThread = new Thread(new Runnable() {
    34. @Override
    35. public void run() {
    36. try {
    37. TimeUnit.MILLISECONDS.sleep(5000 - System.currentTimeMillis()%1000 );
    38. } catch (InterruptedException e) {
    39. if (!scheduleThreadToStop) {
    40. logger.error(e.getMessage(), e);
    41. }
    42. }
    43. logger.info(">>>>>>>>> init xxl-job admin scheduler success.");
    44. // pre-read count: treadpool-size * trigger-qps (each trigger cost 50ms, qps = 1000/50 = 20)
    45. // 预读触发数 = 线程数 * 触发执行qps(限制qps:20/s)
    46. int preReadCount = (XxlJobAdminConfig.getAdminConfig().getTriggerPoolFastMax() + XxlJobAdminConfig.getAdminConfig().getTriggerPoolSlowMax()) * 20;
    47. while (!scheduleThreadToStop) {
    48. // Scan Job
    49. long start = System.currentTimeMillis();
    50. Connection conn = null;
    51. Boolean connAutoCommit = null;
    52. PreparedStatement preparedStatement = null;
    53. boolean preReadSuc = true;
    54. try {
    55. conn = XxlJobAdminConfig.getAdminConfig().getDataSource().getConnection();
    56. connAutoCommit = conn.getAutoCommit();
    57. conn.setAutoCommit(false);
    58. // 调度中心支持多节点部署,基于数据库行锁保证同时只有一个调度中心节点触发任务调度
    59. preparedStatement = conn.prepareStatement( "select * from xxl_job_lock where lock_name = 'schedule_lock' for update" );
    60. preparedStatement.execute();
    61. // tx start
    62. // 1、pre read
    63. long nowTime = System.currentTimeMillis();
    64. // 查询最近5s等待执行的任务(限制数量为preReadCount)
    65. List scheduleList = XxlJobAdminConfig.getAdminConfig().getXxlJobInfoDao().scheduleJobQuery(nowTime + PRE_READ_MS, preReadCount);
    66. if (scheduleList!=null && scheduleList.size()>0) {
    67. // 2、push time-ring
    68. for (XxlJobInfo jobInfo: scheduleList) {
    69. // time-ring jump
    70. if (nowTime > jobInfo.getTriggerNextTime() + PRE_READ_MS) {
    71. // 2.1、trigger-expire > 5s:pass && make next-trigger-time
    72. // 任务触发时间已过,则跳过执行
    73. logger.warn(">>>>>>>>>>> xxl-job, schedule misfire, jobId = " + jobInfo.getId());
    74. // fresh next
    75. // 更新该任务下次执行时间
    76. refreshNextValidTime(jobInfo, new Date());
    77. } else if (nowTime > jobInfo.getTriggerNextTime()) {
    78. // 2.2、trigger-expire < 5s:direct-trigger && make next-trigger-time
    79. // 任务触发时间已到,则立即执行
    80. // 1、trigger
    81. JobTriggerPoolHelper.trigger(jobInfo.getId(), TriggerTypeEnum.CRON, -1, null, null);
    82. logger.debug(">>>>>>>>>>> xxl-job, schedule push trigger : jobId = " + jobInfo.getId() );
    83. // 2、fresh next
    84. // 更新该任务下次执行时间
    85. refreshNextValidTime(jobInfo, new Date());
    86. // next-trigger-time in 5s, pre-read again
    87. // 任务下次触发时间在5s内,则需再执行一次
    88. if (jobInfo.getTriggerStatus()==1 && nowTime + PRE_READ_MS > jobInfo.getTriggerNextTime()) {
    89. // 针对待下次执行,任务会已过执行时间,因此必须当前立即执行,避免任务执行遗漏
    90. // 1、make ring second
    91. // 计算不足1s的剩余毫秒
    92. int ringSecond = (int)((jobInfo.getTriggerNextTime()/1000)%60);
    93. // 2、push time ring
    94. // 存放待执行的循环执行容器
    95. pushTimeRing(ringSecond, jobInfo.getId());
    96. // 3、fresh next
    97. // 更新该任务下次执行时间
    98. refreshNextValidTime(jobInfo, new Date(jobInfo.getTriggerNextTime()));
    99. }
    100. } else {
    101. // 2.3、trigger-pre-read:time-ring trigger && make next-trigger-time
    102. // 1、make ring second
    103. int ringSecond = (int)((jobInfo.getTriggerNextTime()/1000)%60);
    104. // 2、push time ring
    105. pushTimeRing(ringSecond, jobInfo.getId());
    106. // 3、fresh next
    107. refreshNextValidTime(jobInfo, new Date(jobInfo.getTriggerNextTime()));
    108. }
    109. }
    110. // 3、update trigger info
    111. for (XxlJobInfo jobInfo: scheduleList) {
    112. XxlJobAdminConfig.getAdminConfig().getXxlJobInfoDao().scheduleUpdate(jobInfo);
    113. }
    114. } else {
    115. preReadSuc = false;
    116. }
    117. // tx stop
    118. } catch (Exception e) {
    119. if (!scheduleThreadToStop) {
    120. logger.error(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread error:{}", e);
    121. }
    122. } finally {
    123. // commit
    124. if (conn != null) {
    125. try {
    126. conn.commit();
    127. } catch (SQLException e) {
    128. if (!scheduleThreadToStop) {
    129. logger.error(e.getMessage(), e);
    130. }
    131. }
    132. try {
    133. conn.setAutoCommit(connAutoCommit);
    134. } catch (SQLException e) {
    135. if (!scheduleThreadToStop) {
    136. logger.error(e.getMessage(), e);
    137. }
    138. }
    139. try {
    140. conn.close();
    141. } catch (SQLException e) {
    142. if (!scheduleThreadToStop) {
    143. logger.error(e.getMessage(), e);
    144. }
    145. }
    146. }
    147. // close PreparedStatement
    148. if (null != preparedStatement) {
    149. try {
    150. preparedStatement.close();
    151. } catch (SQLException e) {
    152. if (!scheduleThreadToStop) {
    153. logger.error(e.getMessage(), e);
    154. }
    155. }
    156. }
    157. }
    158. long cost = System.currentTimeMillis()-start;
    159. // Wait seconds, align second
    160. if (cost < 1000) { // scan-overtime, not wait
    161. try {
    162. // pre-read period: success > scan each second; fail > skip this period;
    163. TimeUnit.MILLISECONDS.sleep((preReadSuc?1000:PRE_READ_MS) - System.currentTimeMillis()%1000);
    164. } catch (InterruptedException e) {
    165. if (!scheduleThreadToStop) {
    166. logger.error(e.getMessage(), e);
    167. }
    168. }
    169. }
    170. }
    171. logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread stop");
    172. }
    173. });
    174. scheduleThread.setDaemon(true);
    175. scheduleThread.setName("xxl-job, admin JobScheduleHelper#scheduleThread");
    176. scheduleThread.start();
    177. // ring thread
    178. ringThread = new Thread(new Runnable() {
    179. @Override
    180. public void run() {
    181. // align second
    182. try {
    183. TimeUnit.MILLISECONDS.sleep(1000 - System.currentTimeMillis()%1000 );
    184. } catch (InterruptedException e) {
    185. if (!ringThreadToStop) {
    186. logger.error(e.getMessage(), e);
    187. }
    188. }
    189. while (!ringThreadToStop) {
    190. try {
    191. // second data
    192. List ringItemData = new ArrayList<>();
    193. int nowSecond = Calendar.getInstance().get(Calendar.SECOND); // 避免处理耗时太长,跨过刻度,向前校验一个刻度;
    194. for (int i = 0; i < 2; i++) {
    195. List tmpData = ringData.remove( (nowSecond+60-i)%60 );
    196. if (tmpData != null) {
    197. ringItemData.addAll(tmpData);
    198. }
    199. }
    200. // ring trigger
    201. logger.debug(">>>>>>>>>>> xxl-job, time-ring beat : " + nowSecond + " = " + Arrays.asList(ringItemData) );
    202. if (ringItemData.size() > 0) {
    203. // do trigger
    204. for (int jobId: ringItemData) {
    205. // do trigger
    206. JobTriggerPoolHelper.trigger(jobId, TriggerTypeEnum.CRON, -1, null, null);
    207. }
    208. // clear
    209. ringItemData.clear();
    210. }
    211. } catch (Exception e) {
    212. if (!ringThreadToStop) {
    213. logger.error(">>>>>>>>>>> xxl-job, JobScheduleHelper#ringThread error:{}", e);
    214. }
    215. }
    216. // next second, align second
    217. try {
    218. TimeUnit.MILLISECONDS.sleep(1000 - System.currentTimeMillis()%1000);
    219. } catch (InterruptedException e) {
    220. if (!ringThreadToStop) {
    221. logger.error(e.getMessage(), e);
    222. }
    223. }
    224. }
    225. logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#ringThread stop");
    226. }
    227. });
    228. ringThread.setDaemon(true);
    229. ringThread.setName("xxl-job, admin JobScheduleHelper#ringThread");
    230. ringThread.start();
    231. }
    232. private void refreshNextValidTime(XxlJobInfo jobInfo, Date fromTime) throws ParseException {
    233. Date nextValidTime = new CronExpression(jobInfo.getJobCron()).getNextValidTimeAfter(fromTime);
    234. if (nextValidTime != null) {
    235. jobInfo.setTriggerLastTime(jobInfo.getTriggerNextTime());
    236. jobInfo.setTriggerNextTime(nextValidTime.getTime());
    237. } else {
    238. jobInfo.setTriggerStatus(0);
    239. jobInfo.setTriggerLastTime(0);
    240. jobInfo.setTriggerNextTime(0);
    241. }
    242. }
    243. private void pushTimeRing(int ringSecond, int jobId){
    244. // push async ring
    245. List ringItemData = ringData.get(ringSecond);
    246. if (ringItemData == null) {
    247. ringItemData = new ArrayList();
    248. ringData.put(ringSecond, ringItemData);
    249. }
    250. ringItemData.add(jobId);
    251. logger.debug(">>>>>>>>>>> xxl-job, schedule push time-ring : " + ringSecond + " = " + Arrays.asList(ringItemData) );
    252. }
    253. public void toStop(){
    254. // 1、stop schedule
    255. scheduleThreadToStop = true;
    256. try {
    257. TimeUnit.SECONDS.sleep(1); // wait
    258. } catch (InterruptedException e) {
    259. logger.error(e.getMessage(), e);
    260. }
    261. if (scheduleThread.getState() != Thread.State.TERMINATED){
    262. // interrupt and wait
    263. scheduleThread.interrupt();
    264. try {
    265. scheduleThread.join();
    266. } catch (InterruptedException e) {
    267. logger.error(e.getMessage(), e);
    268. }
    269. }
    270. // if has ring data
    271. boolean hasRingData = false;
    272. if (!ringData.isEmpty()) {
    273. for (int second : ringData.keySet()) {
    274. List tmpData = ringData.get(second);
    275. if (tmpData!=null && tmpData.size()>0) {
    276. hasRingData = true;
    277. break;
    278. }
    279. }
    280. }
    281. if (hasRingData) {
    282. try {
    283. TimeUnit.SECONDS.sleep(8);
    284. } catch (InterruptedException e) {
    285. logger.error(e.getMessage(), e);
    286. }
    287. }
    288. // stop ring (wait job-in-memory stop)
    289. ringThreadToStop = true;
    290. try {
    291. TimeUnit.SECONDS.sleep(1);
    292. } catch (InterruptedException e) {
    293. logger.error(e.getMessage(), e);
    294. }
    295. if (ringThread.getState() != Thread.State.TERMINATED){
    296. // interrupt and wait
    297. ringThread.interrupt();
    298. try {
    299. ringThread.join();
    300. } catch (InterruptedException e) {
    301. logger.error(e.getMessage(), e);
    302. }
    303. }
    304. logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper stop");
    305. }
    306. }

    3.3 任务执行时序图

    3.4 任务执行核心代码解析

    3.4.1 根据路由策略定位到执行器

    任务执行:XxlJobTrigger

    1. package com.xxl.job.admin.core.trigger;
    2. import com.xxl.job.admin.core.conf.XxlJobAdminConfig;
    3. import com.xxl.job.admin.core.scheduler.XxlJobScheduler;
    4. import com.xxl.job.admin.core.model.XxlJobGroup;
    5. import com.xxl.job.admin.core.model.XxlJobInfo;
    6. import com.xxl.job.admin.core.model.XxlJobLog;
    7. import com.xxl.job.admin.core.route.ExecutorRouteStrategyEnum;
    8. import com.xxl.job.admin.core.util.I18nUtil;
    9. import com.xxl.job.core.biz.ExecutorBiz;
    10. import com.xxl.job.core.biz.model.ReturnT;
    11. import com.xxl.job.core.biz.model.TriggerParam;
    12. import com.xxl.job.core.enums.ExecutorBlockStrategyEnum;
    13. import com.xxl.rpc.util.IpUtil;
    14. import com.xxl.rpc.util.ThrowableUtil;
    15. import org.slf4j.Logger;
    16. import org.slf4j.LoggerFactory;
    17. import java.util.Date;
    18. /**
    19. * xxl-job trigger
    20. * Created by xuxueli on 17/7/13.
    21. */
    22. public class XxlJobTrigger {
    23. private static Logger logger = LoggerFactory.getLogger(XxlJobTrigger.class);
    24. /**
    25. * trigger job
    26. *
    27. * @param jobId
    28. * @param triggerType
    29. * @param failRetryCount
    30. * >=0: use this param
    31. * <0: use param from job info config
    32. * @param executorShardingParam
    33. * @param executorParam
    34. * null: use job param
    35. * not null: cover job param
    36. */
    37. public static void trigger(int jobId, TriggerTypeEnum triggerType, int failRetryCount, String executorShardingParam, String executorParam) {
    38. // load data
    39. //通过jobId查询出任务信息
    40. XxlJobInfo jobInfo = XxlJobAdminConfig.getAdminConfig().getXxlJobInfoDao().loadById(jobId);
    41. if (jobInfo == null) {
    42. logger.warn(">>>>>>>>>>>> trigger fail, jobId invalid,jobId={}", jobId);
    43. return;
    44. }
    45. if (executorParam != null) {
    46. jobInfo.setExecutorParam(executorParam);
    47. }
    48. //计算失败重试次数
    49. int finalFailRetryCount = failRetryCount>=0?failRetryCount:jobInfo.getExecutorFailRetryCount();
    50. XxlJobGroup group = XxlJobAdminConfig.getAdminConfig().getXxlJobGroupDao().load(jobInfo.getJobGroup());
    51. // sharding param
    52. //处理分片参数
    53. int[] shardingParam = null;
    54. if (executorShardingParam!=null){
    55. String[] shardingArr = executorShardingParam.split("/");
    56. if (shardingArr.length==2 && isNumeric(shardingArr[0]) && isNumeric(shardingArr[1])) {
    57. shardingParam = new int[2];
    58. shardingParam[0] = Integer.valueOf(shardingArr[0]);
    59. shardingParam[1] = Integer.valueOf(shardingArr[1]);
    60. }
    61. }
    62. //处理路由策略
    63. if (ExecutorRouteStrategyEnum.SHARDING_BROADCAST==ExecutorRouteStrategyEnum.match(jobInfo.getExecutorRouteStrategy(), null)
    64. && group.getRegistryList()!=null && !group.getRegistryList().isEmpty()
    65. && shardingParam==null) {
    66. //如果路由策略为分片广播,且执行器地址不为空,则遍历执行器地址进行广播触发任务调度
    67. for (int i = 0; i < group.getRegistryList().size(); i++) {
    68. processTrigger(group, jobInfo, finalFailRetryCount, triggerType, i, group.getRegistryList().size());
    69. }
    70. } else {
    71. //否则初始化分片参数
    72. if (shardingParam == null) {
    73. shardingParam = new int[]{0, 1};
    74. }
    75. processTrigger(group, jobInfo, finalFailRetryCount, triggerType, shardingParam[0], shardingParam[1]);
    76. }
    77. }
    78. private static boolean isNumeric(String str){
    79. try {
    80. int result = Integer.valueOf(str);
    81. return true;
    82. } catch (NumberFormatException e) {
    83. return false;
    84. }
    85. }
    86. /**
    87. * @param group job group, registry list may be empty
    88. * @param jobInfo
    89. * @param finalFailRetryCount
    90. * @param triggerType
    91. * @param index sharding index
    92. * @param total sharding index
    93. */
    94. private static void processTrigger(XxlJobGroup group, XxlJobInfo jobInfo, int finalFailRetryCount, TriggerTypeEnum triggerType, int index, int total){
    95. // param
    96. //执行器阻塞策略:调度过于密集执行器来不及处理时的处理策略,策略包括:单机串行(默认)、丢弃后续调度、覆盖之前调度;
    97. ExecutorBlockStrategyEnum blockStrategy = ExecutorBlockStrategyEnum.match(jobInfo.getExecutorBlockStrategy(), ExecutorBlockStrategyEnum.SERIAL_EXECUTION); // block strategy
    98. //执行器路由策略:执行器集群部署时提供丰富的路由策略,包括:第一个、最后一个、轮询、随机、一致性HASH、最不经常使用、最近最久未使用、故障转移、忙碌转移等
    99. ExecutorRouteStrategyEnum executorRouteStrategyEnum = ExecutorRouteStrategyEnum.match(jobInfo.getExecutorRouteStrategy(), null); // route strategy
    100. //处理分片参数:如果为分片广播则将index和total使用/进行拼接
    101. String shardingParam = (ExecutorRouteStrategyEnum.SHARDING_BROADCAST==executorRouteStrategyEnum)?String.valueOf(index).concat("/").concat(String.valueOf(total)):null;
    102. // 1、save log-id
    103. //保存任务执行日志
    104. XxlJobLog jobLog = new XxlJobLog();
    105. jobLog.setJobGroup(jobInfo.getJobGroup());
    106. jobLog.setJobId(jobInfo.getId());
    107. jobLog.setTriggerTime(new Date());
    108. XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().save(jobLog);
    109. logger.debug(">>>>>>>>>>> xxl-job trigger start, jobId:{}", jobLog.getId());
    110. // 2、init trigger-param
    111. //初始化触发器参数
    112. TriggerParam triggerParam = new TriggerParam();
    113. triggerParam.setJobId(jobInfo.getId());
    114. //取出执行器任务handler名称
    115. triggerParam.setExecutorHandler(jobInfo.getExecutorHandler());
    116. triggerParam.setExecutorParams(jobInfo.getExecutorParam());
    117. triggerParam.setExecutorBlockStrategy(jobInfo.getExecutorBlockStrategy());
    118. triggerParam.setExecutorTimeout(jobInfo.getExecutorTimeout());
    119. triggerParam.setLogId(jobLog.getId());
    120. triggerParam.setLogDateTime(jobLog.getTriggerTime().getTime());
    121. triggerParam.setGlueType(jobInfo.getGlueType());
    122. triggerParam.setGlueSource(jobInfo.getGlueSource());
    123. triggerParam.setGlueUpdatetime(jobInfo.getGlueUpdatetime().getTime());
    124. triggerParam.setBroadcastIndex(index);
    125. triggerParam.setBroadcastTotal(total);
    126. // 3、init address
    127. //初始化执行器地址信息
    128. String address = null;
    129. ReturnT routeAddressResult = null;
    130. if (group.getRegistryList()!=null && !group.getRegistryList().isEmpty()) {
    131. //如果执行器的注册地址不为空
    132. if (ExecutorRouteStrategyEnum.SHARDING_BROADCAST == executorRouteStrategyEnum) {
    133. //如果路由策略是分片广播,取出index对应的执行器地址信息
    134. if (index < group.getRegistryList().size()) {
    135. address = group.getRegistryList().get(index);
    136. } else {
    137. address = group.getRegistryList().get(0);
    138. }
    139. } else {
    140. //获取执行器路由地址信息
    141. //⚠️⚠️️⚠️非常重要:此处使用了策略模式, 根据不同的策略 使用不同的实现类,从而选举出本次使用执行器集群地址列表中不同的地址信息,这是xxl-job实现任务路由的地方
    142. routeAddressResult = executorRouteStrategyEnum.getRouter().route(triggerParam, group.getRegistryList());
    143. if (routeAddressResult.getCode() == ReturnT.SUCCESS_CODE) {
    144. address = routeAddressResult.getContent();
    145. }
    146. }
    147. } else {
    148. routeAddressResult = new ReturnT(ReturnT.FAIL_CODE, I18nUtil.getString("jobconf_trigger_address_empty"));
    149. }
    150. // 4、trigger remote executor
    151. //触发远程执行器执行任务(执行器中的handler)
    152. ReturnT triggerResult = null;
    153. if (address != null) {
    154. //这里是重点⚠️⚠️⚠️:启动执行器, 向执行器发送指令都是从这个方法中执行的
    155. triggerResult = runExecutor(triggerParam, address);
    156. } else {
    157. triggerResult = new ReturnT(ReturnT.FAIL_CODE, null);
    158. }
    159. // 5、collection trigger info
    160. StringBuffer triggerMsgSb = new StringBuffer();
    161. triggerMsgSb.append(I18nUtil.getString("jobconf_trigger_type")).append(":").append(triggerType.getTitle());
    162. triggerMsgSb.append("
      "
      ).append(I18nUtil.getString("jobconf_trigger_admin_adress")).append(":").append(IpUtil.getIp());
    163. triggerMsgSb.append("
      "
      ).append(I18nUtil.getString("jobconf_trigger_exe_regtype")).append(":")
    164. .append( (group.getAddressType() == 0)?I18nUtil.getString("jobgroup_field_addressType_0"):I18nUtil.getString("jobgroup_field_addressType_1") );
    165. triggerMsgSb.append("
      "
      ).append(I18nUtil.getString("jobconf_trigger_exe_regaddress")).append(":").append(group.getRegistryList());
    166. triggerMsgSb.append("
      "
      ).append(I18nUtil.getString("jobinfo_field_executorRouteStrategy")).append(":").append(executorRouteStrategyEnum.getTitle());
    167. if (shardingParam != null) {
    168. triggerMsgSb.append("("+shardingParam+")");
    169. }
    170. triggerMsgSb.append("
      "
      ).append(I18nUtil.getString("jobinfo_field_executorBlockStrategy")).append(":").append(blockStrategy.getTitle());
    171. triggerMsgSb.append("
      "
      ).append(I18nUtil.getString("jobinfo_field_timeout")).append(":").append(jobInfo.getExecutorTimeout());
    172. triggerMsgSb.append("
      "
      ).append(I18nUtil.getString("jobinfo_field_executorFailRetryCount")).append(":").append(finalFailRetryCount);
    173. triggerMsgSb.append("

      >>>>>>>>>>>"+ I18nUtil.getString("jobconf_trigger_run") +"<<<<<<<<<<<
      "
      )
    174. .append((routeAddressResult!=null&&routeAddressResult.getMsg()!=null)?routeAddressResult.getMsg()+"

      "
      :"").append(triggerResult.getMsg()!=null?triggerResult.getMsg():"");
    175. // 6、save log trigger-info
    176. //更新执行日志
    177. jobLog.setExecutorAddress(address);
    178. jobLog.setExecutorHandler(jobInfo.getExecutorHandler());
    179. jobLog.setExecutorParam(jobInfo.getExecutorParam());
    180. jobLog.setExecutorShardingParam(shardingParam);
    181. jobLog.setExecutorFailRetryCount(finalFailRetryCount);
    182. //jobLog.setTriggerTime();
    183. jobLog.setTriggerCode(triggerResult.getCode());
    184. jobLog.setTriggerMsg(triggerMsgSb.toString());
    185. XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().updateTriggerInfo(jobLog);
    186. logger.debug(">>>>>>>>>>> xxl-job trigger end, jobId:{}", jobLog.getId());
    187. }
    188. /**
    189. * run executor
    190. * @param triggerParam
    191. * @param address
    192. * @return
    193. */
    194. public static ReturnT runExecutor(TriggerParam triggerParam, String address){
    195. ReturnT runResult = null;
    196. try {
    197. //这是重点:获取ExecutorBiz代理对象,先从缓存(内存map)中取,取不到new一个XxlRpcReferenceBean放进缓存map
    198. ExecutorBiz executorBiz = XxlJobScheduler.getExecutorBiz(address);
    199. // 这个run方法不会最终执行,仅仅只是为了触发代理对象的invoke方法,同时将目标的类型传送给服务端,因为在代理对象的invoke的方法里面没有执行目标对象的方法
    200. runResult = executorBiz.run(triggerParam);
    201. } catch (Exception e) {
    202. logger.error(">>>>>>>>>>> xxl-job trigger error, please check if the executor[{}] is running.", address, e);
    203. runResult = new ReturnT(ReturnT.FAIL_CODE, ThrowableUtil.toString(e));
    204. }
    205. StringBuffer runResultSB = new StringBuffer(I18nUtil.getString("jobconf_trigger_run") + ":");
    206. runResultSB.append("
      address:"
      ).append(address);
    207. runResultSB.append("
      code:"
      ).append(runResult.getCode());
    208. runResultSB.append("
      msg:"
      ).append(runResult.getMsg());
    209. runResult.setMsg(runResultSB.toString());
    210. return runResult;
    211. }
    212. }

    3.4.2 定位到任务执行Handler&线程

    获得执行Handler&线程:ExecutorBizImpl

    1. package com.xxl.job.core.biz.impl;
    2. import com.xxl.job.core.biz.ExecutorBiz;
    3. import com.xxl.job.core.biz.model.LogResult;
    4. import com.xxl.job.core.biz.model.ReturnT;
    5. import com.xxl.job.core.biz.model.TriggerParam;
    6. import com.xxl.job.core.enums.ExecutorBlockStrategyEnum;
    7. import com.xxl.job.core.executor.XxlJobExecutor;
    8. import com.xxl.job.core.glue.GlueFactory;
    9. import com.xxl.job.core.glue.GlueTypeEnum;
    10. import com.xxl.job.core.handler.IJobHandler;
    11. import com.xxl.job.core.handler.impl.GlueJobHandler;
    12. import com.xxl.job.core.handler.impl.ScriptJobHandler;
    13. import com.xxl.job.core.log.XxlJobFileAppender;
    14. import com.xxl.job.core.thread.JobThread;
    15. import org.slf4j.Logger;
    16. import org.slf4j.LoggerFactory;
    17. import java.util.Date;
    18. /**
    19. * Created by xuxueli on 17/3/1.
    20. * @Description:阻塞处理策略,当执行器节点存在多个相同任务id的任务未执行完成,则需要基于阻塞策略对任务进行取舍: 串行策略:默认策略,任务进行排队、丢弃旧任务策略、丢弃新任务策略
    21. */
    22. public class ExecutorBizImpl implements ExecutorBiz {
    23. private static Logger logger = LoggerFactory.getLogger(ExecutorBizImpl.class);
    24. @Override
    25. public ReturnT beat() {
    26. return ReturnT.SUCCESS;
    27. }
    28. @Override
    29. public ReturnT idleBeat(int jobId) {
    30. // isRunningOrHasQueue
    31. boolean isRunningOrHasQueue = false;
    32. JobThread jobThread = XxlJobExecutor.loadJobThread(jobId);
    33. if (jobThread != null && jobThread.isRunningOrHasQueue()) {
    34. isRunningOrHasQueue = true;
    35. }
    36. if (isRunningOrHasQueue) {
    37. return new ReturnT(ReturnT.FAIL_CODE, "job thread is running or has trigger queue.");
    38. }
    39. return ReturnT.SUCCESS;
    40. }
    41. @Override
    42. public ReturnT kill(int jobId) {
    43. // kill handlerThread, and create new one
    44. JobThread jobThread = XxlJobExecutor.loadJobThread(jobId);
    45. if (jobThread != null) {
    46. XxlJobExecutor.removeJobThread(jobId, "scheduling center kill job.");
    47. return ReturnT.SUCCESS;
    48. }
    49. return new ReturnT(ReturnT.SUCCESS_CODE, "job thread already killed.");
    50. }
    51. @Override
    52. public ReturnT log(long logDateTim, long logId, int fromLineNum) {
    53. // log filename: logPath/yyyy-MM-dd/9999.log
    54. String logFileName = XxlJobFileAppender.makeLogFileName(new Date(logDateTim), logId);
    55. LogResult logResult = XxlJobFileAppender.readLog(logFileName, fromLineNum);
    56. return new ReturnT(logResult);
    57. }
    58. @Override
    59. public ReturnT run(TriggerParam triggerParam) {
    60. // load old:jobHandler + jobThread
    61. JobThread jobThread = XxlJobExecutor.loadJobThread(triggerParam.getJobId());
    62. IJobHandler jobHandler = jobThread!=null?jobThread.getHandler():null;
    63. String removeOldReason = null;
    64. // valid:jobHandler + jobThread
    65. GlueTypeEnum glueTypeEnum = GlueTypeEnum.match(triggerParam.getGlueType());
    66. if (GlueTypeEnum.BEAN == glueTypeEnum) {
    67. // new jobhandler
    68. IJobHandler newJobHandler = XxlJobExecutor.loadJobHandler(triggerParam.getExecutorHandler());
    69. // valid old jobThread
    70. if (jobThread!=null && jobHandler != newJobHandler) {
    71. // change handler, need kill old thread
    72. removeOldReason = "change jobhandler or glue type, and terminate the old job thread.";
    73. jobThread = null;
    74. jobHandler = null;
    75. }
    76. // valid handler
    77. if (jobHandler == null) {
    78. jobHandler = newJobHandler;
    79. if (jobHandler == null) {
    80. return new ReturnT(ReturnT.FAIL_CODE, "job handler [" + triggerParam.getExecutorHandler() + "] not found.");
    81. }
    82. }
    83. } else if (GlueTypeEnum.GLUE_GROOVY == glueTypeEnum) {
    84. // valid old jobThread
    85. if (jobThread != null &&
    86. !(jobThread.getHandler() instanceof GlueJobHandler
    87. && ((GlueJobHandler) jobThread.getHandler()).getGlueUpdatetime()==triggerParam.getGlueUpdatetime() )) {
    88. // change handler or gluesource updated, need kill old thread
    89. removeOldReason = "change job source or glue type, and terminate the old job thread.";
    90. jobThread = null;
    91. jobHandler = null;
    92. }
    93. // valid handler
    94. if (jobHandler == null) {
    95. try {
    96. IJobHandler originJobHandler = GlueFactory.getInstance().loadNewInstance(triggerParam.getGlueSource());
    97. jobHandler = new GlueJobHandler(originJobHandler, triggerParam.getGlueUpdatetime());
    98. } catch (Exception e) {
    99. logger.error(e.getMessage(), e);
    100. return new ReturnT(ReturnT.FAIL_CODE, e.getMessage());
    101. }
    102. }
    103. } else if (glueTypeEnum!=null && glueTypeEnum.isScript()) {
    104. // valid old jobThread
    105. if (jobThread != null &&
    106. !(jobThread.getHandler() instanceof ScriptJobHandler
    107. && ((ScriptJobHandler) jobThread.getHandler()).getGlueUpdatetime()==triggerParam.getGlueUpdatetime() )) {
    108. // change script or gluesource updated, need kill old thread
    109. removeOldReason = "change job source or glue type, and terminate the old job thread.";
    110. jobThread = null;
    111. jobHandler = null;
    112. }
    113. // valid handler
    114. if (jobHandler == null) {
    115. jobHandler = new ScriptJobHandler(triggerParam.getJobId(), triggerParam.getGlueUpdatetime(), triggerParam.getGlueSource(), GlueTypeEnum.match(triggerParam.getGlueType()));
    116. }
    117. } else {
    118. return new ReturnT(ReturnT.FAIL_CODE, "glueType[" + triggerParam.getGlueType() + "] is not valid.");
    119. }
    120. // executor block strategy
    121. if (jobThread != null) {
    122. ExecutorBlockStrategyEnum blockStrategy = ExecutorBlockStrategyEnum.match(triggerParam.getExecutorBlockStrategy(), null);
    123. if (ExecutorBlockStrategyEnum.DISCARD_LATER == blockStrategy) {
    124. // discard when running
    125. if (jobThread.isRunningOrHasQueue()) {
    126. return new ReturnT(ReturnT.FAIL_CODE, "block strategy effect:"+ExecutorBlockStrategyEnum.DISCARD_LATER.getTitle());
    127. }
    128. } else if (ExecutorBlockStrategyEnum.COVER_EARLY == blockStrategy) {
    129. // kill running jobThread
    130. if (jobThread.isRunningOrHasQueue()) {
    131. removeOldReason = "block strategy effect:" + ExecutorBlockStrategyEnum.COVER_EARLY.getTitle();
    132. jobThread = null;
    133. }
    134. } else {
    135. // just queue trigger
    136. }
    137. }
    138. // replace thread (new or exists invalid)
    139. if (jobThread == null) {
    140. // 任务执行
    141. jobThread = XxlJobExecutor.registJobThread(triggerParam.getJobId(), jobHandler, removeOldReason);
    142. }
    143. // push data to queue
    144. ReturnT pushResult = jobThread.pushTriggerQueue(triggerParam);
    145. return pushResult;
    146. }
    147. }

    3.4.3 任务执行线程

    任务线程执行&放入回调队列:JobThread

    1. package com.xxl.job.core.thread;
    2. import com.xxl.job.core.biz.model.HandleCallbackParam;
    3. import com.xxl.job.core.biz.model.ReturnT;
    4. import com.xxl.job.core.biz.model.TriggerParam;
    5. import com.xxl.job.core.executor.XxlJobExecutor;
    6. import com.xxl.job.core.handler.IJobHandler;
    7. import com.xxl.job.core.log.XxlJobFileAppender;
    8. import com.xxl.job.core.log.XxlJobLogger;
    9. import com.xxl.job.core.util.ShardingUtil;
    10. import org.slf4j.Logger;
    11. import org.slf4j.LoggerFactory;
    12. import java.io.PrintWriter;
    13. import java.io.StringWriter;
    14. import java.util.Collections;
    15. import java.util.Date;
    16. import java.util.HashSet;
    17. import java.util.Set;
    18. import java.util.concurrent.*;
    19. /**
    20. * handler thread
    21. * @author xuxueli 2016-1-16 19:52:47
    22. */
    23. public class JobThread extends Thread{
    24. private static Logger logger = LoggerFactory.getLogger(JobThread.class);
    25. private int jobId;
    26. private IJobHandler handler;
    27. private LinkedBlockingQueue triggerQueue;
    28. private Set triggerLogIdSet; // avoid repeat trigger for the same TRIGGER_LOG_ID
    29. private volatile boolean toStop = false;
    30. private String stopReason;
    31. private boolean running = false; // if running job
    32. private int idleTimes = 0; // idel times
    33. public JobThread(int jobId, IJobHandler handler) {
    34. this.jobId = jobId;
    35. this.handler = handler;
    36. this.triggerQueue = new LinkedBlockingQueue();
    37. this.triggerLogIdSet = Collections.synchronizedSet(new HashSet());
    38. }
    39. public IJobHandler getHandler() {
    40. return handler;
    41. }
    42. /**
    43. * new trigger to queue
    44. *
    45. * @param triggerParam
    46. * @return
    47. */
    48. public ReturnT pushTriggerQueue(TriggerParam triggerParam) {
    49. // avoid repeat
    50. if (triggerLogIdSet.contains(triggerParam.getLogId())) {
    51. logger.info(">>>>>>>>>>> repeate trigger job, logId:{}", triggerParam.getLogId());
    52. return new ReturnT(ReturnT.FAIL_CODE, "repeate trigger job, logId:" + triggerParam.getLogId());
    53. }
    54. triggerLogIdSet.add(triggerParam.getLogId());
    55. triggerQueue.add(triggerParam);
    56. return ReturnT.SUCCESS;
    57. }
    58. /**
    59. * kill job thread
    60. *
    61. * @param stopReason
    62. */
    63. public void toStop(String stopReason) {
    64. /**
    65. * Thread.interrupt只支持终止线程的阻塞状态(wait、join、sleep),
    66. * 在阻塞出抛出InterruptedException异常,但是并不会终止运行的线程本身;
    67. * 所以需要注意,此处彻底销毁本线程,需要通过共享变量方式;
    68. */
    69. this.toStop = true;
    70. this.stopReason = stopReason;
    71. }
    72. /**
    73. * is running job
    74. * @return
    75. */
    76. public boolean isRunningOrHasQueue() {
    77. return running || triggerQueue.size()>0;
    78. }
    79. @Override
    80. public void run() {
    81. // init
    82. try {
    83. handler.init();
    84. } catch (Throwable e) {
    85. logger.error(e.getMessage(), e);
    86. }
    87. // execute
    88. while(!toStop){
    89. running = false;
    90. idleTimes++;
    91. TriggerParam triggerParam = null;
    92. ReturnT executeResult = null;
    93. try {
    94. // to check toStop signal, we need cycle, so wo cannot use queue.take(), instand of poll(timeout)
    95. triggerParam = triggerQueue.poll(3L, TimeUnit.SECONDS);
    96. if (triggerParam!=null) {
    97. running = true;
    98. idleTimes = 0;
    99. triggerLogIdSet.remove(triggerParam.getLogId());
    100. // log filename, like "logPath/yyyy-MM-dd/9999.log"
    101. String logFileName = XxlJobFileAppender.makeLogFileName(new Date(triggerParam.getLogDateTime()), triggerParam.getLogId());
    102. XxlJobFileAppender.contextHolder.set(logFileName);
    103. ShardingUtil.setShardingVo(new ShardingUtil.ShardingVO(triggerParam.getBroadcastIndex(), triggerParam.getBroadcastTotal()));
    104. // execute
    105. XxlJobLogger.log("
      ----------- xxl-job job execute start -----------
      ----------- Param:"
      + triggerParam.getExecutorParams());
    106. if (triggerParam.getExecutorTimeout() > 0) {
    107. // limit timeout
    108. Thread futureThread = null;
    109. try {
    110. final TriggerParam triggerParamTmp = triggerParam;
    111. FutureTask> futureTask = new FutureTask>(new Callable>() {
    112. @Override
    113. public ReturnT call() throws Exception {
    114. return handler.execute(triggerParamTmp.getExecutorParams());
    115. }
    116. });
    117. futureThread = new Thread(futureTask);
    118. futureThread.start();
    119. executeResult = futureTask.get(triggerParam.getExecutorTimeout(), TimeUnit.SECONDS);
    120. } catch (TimeoutException e) {
    121. XxlJobLogger.log("
      ----------- xxl-job job execute timeout"
      );
    122. XxlJobLogger.log(e);
    123. executeResult = new ReturnT(IJobHandler.FAIL_TIMEOUT.getCode(), "job execute timeout ");
    124. } finally {
    125. futureThread.interrupt();
    126. }
    127. } else {
    128. // just execute
    129. executeResult = handler.execute(triggerParam.getExecutorParams());
    130. }
    131. if (executeResult == null) {
    132. executeResult = IJobHandler.FAIL;
    133. } else {
    134. executeResult.setMsg(
    135. (executeResult!=null&&executeResult.getMsg()!=null&&executeResult.getMsg().length()>50000)
    136. ?executeResult.getMsg().substring(0, 50000).concat("...")
    137. :executeResult.getMsg());
    138. executeResult.setContent(null); // limit obj size
    139. }
    140. XxlJobLogger.log("
      ----------- xxl-job job execute end(finish) -----------
      ----------- ReturnT:"
      + executeResult);
    141. } else {
    142. if (idleTimes > 30) {
    143. if(triggerQueue.size() == 0) { // avoid concurrent trigger causes jobId-lost
    144. XxlJobExecutor.removeJobThread(jobId, "excutor idel times over limit.");
    145. }
    146. }
    147. }
    148. } catch (Throwable e) {
    149. if (toStop) {
    150. XxlJobLogger.log("
      ----------- JobThread toStop, stopReason:"
      + stopReason);
    151. }
    152. StringWriter stringWriter = new StringWriter();
    153. e.printStackTrace(new PrintWriter(stringWriter));
    154. String errorMsg = stringWriter.toString();
    155. executeResult = new ReturnT(ReturnT.FAIL_CODE, errorMsg);
    156. XxlJobLogger.log("
      ----------- JobThread Exception:"
      + errorMsg + "
      ----------- xxl-job job execute end(error) -----------"
      );
    157. } finally {
    158. if(triggerParam != null) {
    159. // callback handler info
    160. if (!toStop) {
    161. // commonm
    162. TriggerCallbackThread.pushCallBack(new HandleCallbackParam(triggerParam.getLogId(), triggerParam.getLogDateTime(), executeResult));
    163. } else {
    164. // is killed
    165. ReturnT stopResult = new ReturnT(ReturnT.FAIL_CODE, stopReason + " [job running, killed]");
    166. TriggerCallbackThread.pushCallBack(new HandleCallbackParam(triggerParam.getLogId(), triggerParam.getLogDateTime(), stopResult));
    167. }
    168. }
    169. }
    170. }
    171. // callback trigger request in queue
    172. while(triggerQueue !=null && triggerQueue.size()>0){
    173. TriggerParam triggerParam = triggerQueue.poll();
    174. if (triggerParam!=null) {
    175. // is killed
    176. ReturnT stopResult = new ReturnT(ReturnT.FAIL_CODE, stopReason + " [job not executed, in the job queue, killed.]");
    177. TriggerCallbackThread.pushCallBack(new HandleCallbackParam(triggerParam.getLogId(), triggerParam.getLogDateTime(), stopResult));
    178. }
    179. }
    180. // destroy
    181. try {
    182. handler.destroy();
    183. } catch (Throwable e) {
    184. logger.error(e.getMessage(), e);
    185. }
    186. logger.info(">>>>>>>>>>> xxl-job JobThread stoped, hashCode:{}", Thread.currentThread());
    187. }
    188. }

    四、执行器

    4.1 启动过程时序图

    4.2 启动过程核心代码解析

    入口类:XxlJobSpringExecutor(PS:继承XxlJobExecutor,所有的任务触发最终都是通过这个类

    4.2.1 启动初始化

    执行器启动初始化:XxlJobSpringExecutor

    1. package com.xxl.job.core.executor.impl;
    2. import com.xxl.job.core.biz.model.ReturnT;
    3. import com.xxl.job.core.executor.XxlJobExecutor;
    4. import com.xxl.job.core.glue.GlueFactory;
    5. import com.xxl.job.core.handler.IJobHandler;
    6. import com.xxl.job.core.handler.annotation.JobHandler;
    7. import com.xxl.job.core.handler.annotation.XxlJob;
    8. import com.xxl.job.core.handler.impl.MethodJobHandler;
    9. import org.springframework.beans.BeansException;
    10. import org.springframework.beans.factory.DisposableBean;
    11. import org.springframework.beans.factory.InitializingBean;
    12. import org.springframework.context.ApplicationContext;
    13. import org.springframework.context.ApplicationContextAware;
    14. import org.springframework.core.annotation.AnnotationUtils;
    15. import java.lang.reflect.Method;
    16. import java.util.Map;
    17. /**
    18. * xxl-job executor (for spring)
    19. *
    20. * @author xuxueli 2018-11-01 09:24:52
    21. * @Description:执行器启动过程
    22. */
    23. public class XxlJobSpringExecutor extends XxlJobExecutor implements ApplicationContextAware, InitializingBean, DisposableBean {
    24. // start
    25. @Override
    26. public void afterPropertiesSet() throws Exception {
    27. // init JobHandler Repository
    28. /**
    29. * 老版本:获取使用了JobHandler注解的bean,即任务处理器,并注册到jobHandlerRepository(ConcurrentMap)缓存中
    30. */
    31. initJobHandlerRepository(applicationContext);
    32. // init JobHandler Repository (for method)
    33. /**
    34. * 新版本:获取使用了XxlJob注解的method,即任务处理器,并注册到jobHandlerRepository(ConcurrentMap)缓存中
    35. */
    36. initJobHandlerMethodRepository(applicationContext);
    37. // refresh GlueFactory
    38. GlueFactory.refreshInstance(1);
    39. // super start
    40. /**
    41. * 核心启动
    42. */
    43. super.start();
    44. }
    45. // destroy
    46. @Override
    47. public void destroy() {
    48. super.destroy();
    49. }
    50. private void initJobHandlerRepository(ApplicationContext applicationContext) {
    51. if (applicationContext == null) {
    52. return;
    53. }
    54. // init job handler action
    55. Map serviceBeanMap = applicationContext.getBeansWithAnnotation(JobHandler.class);
    56. if (serviceBeanMap != null && serviceBeanMap.size() > 0) {
    57. for (Object serviceBean : serviceBeanMap.values()) {
    58. if (serviceBean instanceof IJobHandler) {
    59. String name = serviceBean.getClass().getAnnotation(JobHandler.class).value();
    60. IJobHandler handler = (IJobHandler) serviceBean;
    61. if (loadJobHandler(name) != null) {
    62. throw new RuntimeException("xxl-job jobhandler[" + name + "] naming conflicts.");
    63. }
    64. registJobHandler(name, handler);
    65. }
    66. }
    67. }
    68. }
    69. private void initJobHandlerMethodRepository(ApplicationContext applicationContext) {
    70. if (applicationContext == null) {
    71. return;
    72. }
    73. // init job handler from method
    74. String[] beanDefinitionNames = applicationContext.getBeanDefinitionNames();
    75. if (beanDefinitionNames!=null && beanDefinitionNames.length>0) {
    76. for (String beanDefinitionName : beanDefinitionNames) {
    77. Object bean = applicationContext.getBean(beanDefinitionName);
    78. Method[] methods = bean.getClass().getDeclaredMethods();
    79. for (Method method: methods) {
    80. XxlJob xxlJob = AnnotationUtils.findAnnotation(method, XxlJob.class);
    81. if (xxlJob != null) {
    82. // name
    83. String name = xxlJob.value();
    84. if (name.trim().length() == 0) {
    85. throw new RuntimeException("xxl-job method-jobhandler name invalid, for[" + bean.getClass() + "#"+ method.getName() +"] .");
    86. }
    87. if (loadJobHandler(name) != null) {
    88. throw new RuntimeException("xxl-job jobhandler[" + name + "] naming conflicts.");
    89. }
    90. // execute method
    91. if (!(method.getParameterTypes()!=null && method.getParameterTypes().length==1 && method.getParameterTypes()[0].isAssignableFrom(String.class))) {
    92. throw new RuntimeException("xxl-job method-jobhandler param-classtype invalid, for[" + bean.getClass() + "#"+ method.getName() +"] , " +
    93. "The correct method format like \" public ReturnT execute(String param) \" .");
    94. }
    95. if (!method.getReturnType().isAssignableFrom(ReturnT.class)) {
    96. throw new RuntimeException("xxl-job method-jobhandler return-classtype invalid, for[" + bean.getClass() + "#"+ method.getName() +"] , " +
    97. "The correct method format like \" public ReturnT execute(String param) \" .");
    98. }
    99. method.setAccessible(true);
    100. // init and destory
    101. Method initMethod = null;
    102. Method destroyMethod = null;
    103. if(xxlJob.init().trim().length() > 0) {
    104. try {
    105. initMethod = bean.getClass().getDeclaredMethod(xxlJob.init());
    106. initMethod.setAccessible(true);
    107. } catch (NoSuchMethodException e) {
    108. throw new RuntimeException("xxl-job method-jobhandler initMethod invalid, for[" + bean.getClass() + "#"+ method.getName() +"] .");
    109. }
    110. }
    111. if(xxlJob.destroy().trim().length() > 0) {
    112. try {
    113. destroyMethod = bean.getClass().getDeclaredMethod(xxlJob.destroy());
    114. destroyMethod.setAccessible(true);
    115. } catch (NoSuchMethodException e) {
    116. throw new RuntimeException("xxl-job method-jobhandler destroyMethod invalid, for[" + bean.getClass() + "#"+ method.getName() +"] .");
    117. }
    118. }
    119. // registry jobhandler
    120. registJobHandler(name, new MethodJobHandler(bean, method, initMethod, destroyMethod));
    121. }
    122. }
    123. }
    124. }
    125. }
    126. // ---------------------- applicationContext ----------------------
    127. private static ApplicationContext applicationContext;
    128. @Override
    129. public void setApplicationContext(ApplicationContext applicationContext) throws BeansException {
    130. this.applicationContext = applicationContext;
    131. }
    132. public static ApplicationContext getApplicationContext() {
    133. return applicationContext;
    134. }
    135. }

    4.2.2 执行器初始化

    执行器相关信息准备:XxlJobExecutor

    1. package com.xxl.job.core.executor;
    2. import com.xxl.job.core.biz.AdminBiz;
    3. import com.xxl.job.core.biz.ExecutorBiz;
    4. import com.xxl.job.core.biz.client.AdminBizClient;
    5. import com.xxl.job.core.biz.impl.ExecutorBizImpl;
    6. import com.xxl.job.core.handler.IJobHandler;
    7. import com.xxl.job.core.log.XxlJobFileAppender;
    8. import com.xxl.job.core.thread.ExecutorRegistryThread;
    9. import com.xxl.job.core.thread.JobLogFileCleanThread;
    10. import com.xxl.job.core.thread.JobThread;
    11. import com.xxl.job.core.thread.TriggerCallbackThread;
    12. import com.xxl.rpc.registry.ServiceRegistry;
    13. import com.xxl.rpc.remoting.net.impl.netty_http.server.NettyHttpServer;
    14. import com.xxl.rpc.remoting.provider.XxlRpcProviderFactory;
    15. import com.xxl.rpc.serialize.Serializer;
    16. import com.xxl.rpc.serialize.impl.HessianSerializer;
    17. import com.xxl.rpc.util.IpUtil;
    18. import com.xxl.rpc.util.NetUtil;
    19. import org.slf4j.Logger;
    20. import org.slf4j.LoggerFactory;
    21. import java.util.*;
    22. import java.util.concurrent.ConcurrentHashMap;
    23. import java.util.concurrent.ConcurrentMap;
    24. /**
    25. * Created by xuxueli on 2016/3/2 21:14.
    26. * @Description:任务执行器
    27. */
    28. public class XxlJobExecutor {
    29. private static final Logger logger = LoggerFactory.getLogger(XxlJobExecutor.class);
    30. // ---------------------- param ----------------------
    31. private String adminAddresses;
    32. private String appName;
    33. private String ip;
    34. private int port;
    35. private String accessToken;
    36. private String logPath;
    37. private int logRetentionDays;
    38. public void setAdminAddresses(String adminAddresses) {
    39. this.adminAddresses = adminAddresses;
    40. }
    41. public void setAppName(String appName) {
    42. this.appName = appName;
    43. }
    44. public void setIp(String ip) {
    45. this.ip = ip;
    46. }
    47. public void setPort(int port) {
    48. this.port = port;
    49. }
    50. public void setAccessToken(String accessToken) {
    51. this.accessToken = accessToken;
    52. }
    53. public void setLogPath(String logPath) {
    54. this.logPath = logPath;
    55. }
    56. public void setLogRetentionDays(int logRetentionDays) {
    57. this.logRetentionDays = logRetentionDays;
    58. }
    59. // ---------------------- start + stop ----------------------
    60. public void start() throws Exception {
    61. // init logpath
    62. /**
    63. * 主要作用:初始化日志路径
    64. */
    65. XxlJobFileAppender.initLogPath(logPath);
    66. // init invoker, admin-client
    67. /**
    68. * 主要作用:初始化注册中心列表 (把注册地址放到 List)
    69. * 1、初始化执行器与调度中心建立通信的地址
    70. */
    71. initAdminBizList(adminAddresses, accessToken);
    72. // init JobLogFileCleanThread
    73. /**
    74. * 主要作用:
    75. * 1、启动日志文件清理线程 (一天清理一次)
    76. * 2、每天清理一次过期日志,配置参数必须大于3才有效
    77. */
    78. JobLogFileCleanThread.getInstance().start(logRetentionDays);
    79. // init TriggerCallbackThread
    80. /**
    81. * 主要作用:开启触发器回调线程
    82. */
    83. TriggerCallbackThread.getInstance().start();
    84. // init executor-server
    85. /**
    86. * 主要端口:指定端口
    87. */
    88. port = port>0?port: NetUtil.findAvailablePort(9999);
    89. /**
    90. * 主要端口:指定IP
    91. */
    92. ip = (ip!=null&&ip.trim().length()>0)?ip: IpUtil.getIp();
    93. /**
    94. * 主要作用:启动注册中心的RPC服务,使得执行器项目可以通过RPC进行注册和心跳检测
    95. * 1、将执行器注册到调度中心 30秒一次
    96. */
    97. initRpcProvider(ip, port, appName, accessToken);
    98. }
    99. public void destroy(){
    100. // destory executor-server
    101. stopRpcProvider();
    102. // destory jobThreadRepository
    103. if (jobThreadRepository.size() > 0) {
    104. for (Map.Entry item: jobThreadRepository.entrySet()) {
    105. removeJobThread(item.getKey(), "web container destroy and kill the job.");
    106. }
    107. jobThreadRepository.clear();
    108. }
    109. jobHandlerRepository.clear();
    110. // destory JobLogFileCleanThread
    111. JobLogFileCleanThread.getInstance().toStop();
    112. // destory TriggerCallbackThread
    113. TriggerCallbackThread.getInstance().toStop();
    114. }
    115. // ---------------------- admin-client (rpc invoker) ----------------------
    116. private static List adminBizList;
    117. private static Serializer serializer = new HessianSerializer();
    118. private void initAdminBizList(String adminAddresses, String accessToken) throws Exception {
    119. if (adminAddresses!=null && adminAddresses.trim().length()>0) {
    120. for (String address: adminAddresses.trim().split(",")) {
    121. if (address!=null && address.trim().length()>0) {
    122. AdminBiz adminBiz = new AdminBizClient(address.trim(), accessToken);
    123. if (adminBizList == null) {
    124. adminBizList = new ArrayList();
    125. }
    126. adminBizList.add(adminBiz);
    127. }
    128. }
    129. }
    130. }
    131. public static List getAdminBizList(){
    132. return adminBizList;
    133. }
    134. public static Serializer getSerializer() {
    135. return serializer;
    136. }
    137. // ---------------------- executor-server (rpc provider) ----------------------
    138. private XxlRpcProviderFactory xxlRpcProviderFactory = null;
    139. private void initRpcProvider(String ip, int port, String appName, String accessToken) throws Exception {
    140. // init, provider factory
    141. String address = IpUtil.getIpPort(ip, port);
    142. Map serviceRegistryParam = new HashMap();
    143. serviceRegistryParam.put("appName", appName);
    144. serviceRegistryParam.put("address", address);
    145. //初始化RPC配置
    146. xxlRpcProviderFactory = new XxlRpcProviderFactory();
    147. xxlRpcProviderFactory.setServer(NettyHttpServer.class);
    148. xxlRpcProviderFactory.setSerializer(HessianSerializer.class);
    149. xxlRpcProviderFactory.setCorePoolSize(20);
    150. xxlRpcProviderFactory.setMaxPoolSize(200);
    151. xxlRpcProviderFactory.setIp(ip);
    152. xxlRpcProviderFactory.setPort(port);
    153. xxlRpcProviderFactory.setAccessToken(accessToken);
    154. xxlRpcProviderFactory.setServiceRegistry(ExecutorServiceRegistry.class);
    155. xxlRpcProviderFactory.setServiceRegistryParam(serviceRegistryParam);
    156. // add services
    157. //给xxlRpcProviderFactory加入服务
    158. xxlRpcProviderFactory.addService(ExecutorBiz.class.getName(), null, new ExecutorBizImpl());
    159. // start
    160. /**
    161. * 开启注册
    162. */
    163. xxlRpcProviderFactory.start();
    164. }
    165. public static class ExecutorServiceRegistry extends ServiceRegistry {
    166. @Override
    167. public void start(Map param) {
    168. // start registry
    169. ExecutorRegistryThread.getInstance().start(param.get("appName"), param.get("address"));
    170. }
    171. @Override
    172. public void stop() {
    173. // stop registry
    174. ExecutorRegistryThread.getInstance().toStop();
    175. }
    176. @Override
    177. public boolean registry(Set keys, String value) {
    178. return false;
    179. }
    180. @Override
    181. public boolean remove(Set keys, String value) {
    182. return false;
    183. }
    184. @Override
    185. public Map> discovery(Set keys) {
    186. return null;
    187. }
    188. @Override
    189. public TreeSet discovery(String key) {
    190. return null;
    191. }
    192. }
    193. private void stopRpcProvider() {
    194. // stop provider factory
    195. try {
    196. xxlRpcProviderFactory.stop();
    197. } catch (Exception e) {
    198. logger.error(e.getMessage(), e);
    199. }
    200. }
    201. /**
    202. * 一个JobThread对应一个JobHandler
    203. */
    204. // ---------------------- job handler repository(key为XxlJob注解中的bean名称,value为JobHandler) ----------------------
    205. private static ConcurrentMap jobHandlerRepository = new ConcurrentHashMap();
    206. public static IJobHandler registJobHandler(String name, IJobHandler jobHandler){
    207. logger.info(">>>>>>>>>>> xxl-job register jobhandler success, name:{}, jobHandler:{}", name, jobHandler);
    208. return jobHandlerRepository.put(name, jobHandler);
    209. }
    210. public static IJobHandler loadJobHandler(String name){
    211. return jobHandlerRepository.get(name);
    212. }
    213. // ---------------------- job thread repository(key为任务id,即jobId,value为JobThread对象)----------------------
    214. private static ConcurrentMap jobThreadRepository = new ConcurrentHashMap();
    215. public static JobThread registJobThread(int jobId, IJobHandler handler, String removeOldReason){
    216. JobThread newJobThread = new JobThread(jobId, handler);
    217. newJobThread.start();
    218. logger.info(">>>>>>>>>>> xxl-job regist JobThread success, jobId:{}, handler:{}", new Object[]{jobId, handler});
    219. JobThread oldJobThread = jobThreadRepository.put(jobId, newJobThread); // putIfAbsent | oh my god, map's put method return the old value!!!
    220. if (oldJobThread != null) {
    221. oldJobThread.toStop(removeOldReason);
    222. oldJobThread.interrupt();
    223. }
    224. return newJobThread;
    225. }
    226. public static void removeJobThread(int jobId, String removeOldReason){
    227. JobThread oldJobThread = jobThreadRepository.remove(jobId);
    228. if (oldJobThread != null) {
    229. oldJobThread.toStop(removeOldReason);
    230. oldJobThread.interrupt();
    231. }
    232. }
    233. public static JobThread loadJobThread(int jobId){
    234. JobThread jobThread = jobThreadRepository.get(jobId);
    235. return jobThread;
    236. }
    237. }

    4.2.3 任务执行回调

    任务执行后结果回调:TriggerCallbackThread

    1. package com.xxl.job.core.thread;
    2. import com.xxl.job.core.biz.AdminBiz;
    3. import com.xxl.job.core.biz.model.HandleCallbackParam;
    4. import com.xxl.job.core.biz.model.ReturnT;
    5. import com.xxl.job.core.enums.RegistryConfig;
    6. import com.xxl.job.core.executor.XxlJobExecutor;
    7. import com.xxl.job.core.log.XxlJobFileAppender;
    8. import com.xxl.job.core.log.XxlJobLogger;
    9. import com.xxl.job.core.util.FileUtil;
    10. import org.slf4j.Logger;
    11. import org.slf4j.LoggerFactory;
    12. import java.io.File;
    13. import java.util.ArrayList;
    14. import java.util.Date;
    15. import java.util.List;
    16. import java.util.concurrent.LinkedBlockingQueue;
    17. import java.util.concurrent.TimeUnit;
    18. /**
    19. * Created by xuxueli on 16/7/22.
    20. */
    21. public class TriggerCallbackThread {
    22. private static Logger logger = LoggerFactory.getLogger(TriggerCallbackThread.class);
    23. private static TriggerCallbackThread instance = new TriggerCallbackThread();
    24. public static TriggerCallbackThread getInstance(){
    25. return instance;
    26. }
    27. /**
    28. * job results callback queue
    29. */
    30. private LinkedBlockingQueue callBackQueue = new LinkedBlockingQueue();
    31. public static void pushCallBack(HandleCallbackParam callback){
    32. getInstance().callBackQueue.add(callback);
    33. logger.debug(">>>>>>>>>>> xxl-job, push callback request, logId:{}", callback.getLogId());
    34. }
    35. /**
    36. * callback thread
    37. */
    38. private Thread triggerCallbackThread;
    39. private Thread triggerRetryCallbackThread;
    40. private volatile boolean toStop = false;
    41. public void start() {
    42. // valid
    43. if (XxlJobExecutor.getAdminBizList() == null) {
    44. logger.warn(">>>>>>>>>>> xxl-job, executor callback config fail, adminAddresses is null.");
    45. return;
    46. }
    47. // callback
    48. triggerCallbackThread = new Thread(new Runnable() {
    49. @Override
    50. public void run() {
    51. // normal callback
    52. while(!toStop){
    53. try {
    54. HandleCallbackParam callback = getInstance().callBackQueue.take();
    55. if (callback != null) {
    56. // callback list param
    57. List callbackParamList = new ArrayList();
    58. int drainToNum = getInstance().callBackQueue.drainTo(callbackParamList);
    59. callbackParamList.add(callback);
    60. // callback, will retry if error
    61. if (callbackParamList!=null && callbackParamList.size()>0) {
    62. doCallback(callbackParamList);
    63. }
    64. }
    65. } catch (Exception e) {
    66. if (!toStop) {
    67. logger.error(e.getMessage(), e);
    68. }
    69. }
    70. }
    71. // last callback
    72. try {
    73. List callbackParamList = new ArrayList();
    74. int drainToNum = getInstance().callBackQueue.drainTo(callbackParamList);
    75. if (callbackParamList!=null && callbackParamList.size()>0) {
    76. doCallback(callbackParamList);
    77. }
    78. } catch (Exception e) {
    79. if (!toStop) {
    80. logger.error(e.getMessage(), e);
    81. }
    82. }
    83. logger.info(">>>>>>>>>>> xxl-job, executor callback thread destory.");
    84. }
    85. });
    86. triggerCallbackThread.setDaemon(true);
    87. triggerCallbackThread.setName("xxl-job, executor TriggerCallbackThread");
    88. triggerCallbackThread.start();
    89. // retry
    90. triggerRetryCallbackThread = new Thread(new Runnable() {
    91. @Override
    92. public void run() {
    93. while(!toStop){
    94. try {
    95. retryFailCallbackFile();
    96. } catch (Exception e) {
    97. if (!toStop) {
    98. logger.error(e.getMessage(), e);
    99. }
    100. }
    101. try {
    102. TimeUnit.SECONDS.sleep(RegistryConfig.BEAT_TIMEOUT);
    103. } catch (InterruptedException e) {
    104. if (!toStop) {
    105. logger.error(e.getMessage(), e);
    106. }
    107. }
    108. }
    109. logger.info(">>>>>>>>>>> xxl-job, executor retry callback thread destory.");
    110. }
    111. });
    112. triggerRetryCallbackThread.setDaemon(true);
    113. triggerRetryCallbackThread.start();
    114. }
    115. public void toStop(){
    116. toStop = true;
    117. // stop callback, interrupt and wait
    118. if (triggerCallbackThread != null) { // support empty admin address
    119. triggerCallbackThread.interrupt();
    120. try {
    121. triggerCallbackThread.join();
    122. } catch (InterruptedException e) {
    123. logger.error(e.getMessage(), e);
    124. }
    125. }
    126. // stop retry, interrupt and wait
    127. if (triggerRetryCallbackThread != null) {
    128. triggerRetryCallbackThread.interrupt();
    129. try {
    130. triggerRetryCallbackThread.join();
    131. } catch (InterruptedException e) {
    132. logger.error(e.getMessage(), e);
    133. }
    134. }
    135. }
    136. /**
    137. * do callback, will retry if error
    138. * @param callbackParamList
    139. */
    140. private void doCallback(List callbackParamList){
    141. boolean callbackRet = false;
    142. // callback, will retry if error
    143. for (AdminBiz adminBiz: XxlJobExecutor.getAdminBizList()) {
    144. try {
    145. ReturnT callbackResult = adminBiz.callback(callbackParamList);
    146. if (callbackResult!=null && ReturnT.SUCCESS_CODE == callbackResult.getCode()) {
    147. callbackLog(callbackParamList, "
      ----------- xxl-job job callback finish."
      );
    148. callbackRet = true;
    149. break;
    150. } else {
    151. callbackLog(callbackParamList, "
      ----------- xxl-job job callback fail, callbackResult:"
      + callbackResult);
    152. }
    153. } catch (Exception e) {
    154. callbackLog(callbackParamList, "
      ----------- xxl-job job callback error, errorMsg:"
      + e.getMessage());
    155. }
    156. }
    157. if (!callbackRet) {
    158. appendFailCallbackFile(callbackParamList);
    159. }
    160. }
    161. /**
    162. * callback log
    163. */
    164. private void callbackLog(List callbackParamList, String logContent){
    165. for (HandleCallbackParam callbackParam: callbackParamList) {
    166. String logFileName = XxlJobFileAppender.makeLogFileName(new Date(callbackParam.getLogDateTim()), callbackParam.getLogId());
    167. XxlJobFileAppender.contextHolder.set(logFileName);
    168. XxlJobLogger.log(logContent);
    169. }
    170. }
    171. // ---------------------- fail-callback file ----------------------
    172. private static String failCallbackFilePath = XxlJobFileAppender.getLogPath().concat(File.separator).concat("callbacklog").concat(File.separator);
    173. private static String failCallbackFileName = failCallbackFilePath.concat("xxl-job-callback-{x}").concat(".log");
    174. private void appendFailCallbackFile(List callbackParamList){
    175. // valid
    176. if (callbackParamList==null || callbackParamList.size()==0) {
    177. return;
    178. }
    179. // append file
    180. byte[] callbackParamList_bytes = XxlJobExecutor.getSerializer().serialize(callbackParamList);
    181. File callbackLogFile = new File(failCallbackFileName.replace("{x}", String.valueOf(System.currentTimeMillis())));
    182. if (callbackLogFile.exists()) {
    183. for (int i = 0; i < 100; i++) {
    184. callbackLogFile = new File(failCallbackFileName.replace("{x}", String.valueOf(System.currentTimeMillis()).concat("-").concat(String.valueOf(i)) ));
    185. if (!callbackLogFile.exists()) {
    186. break;
    187. }
    188. }
    189. }
    190. FileUtil.writeFileContent(callbackLogFile, callbackParamList_bytes);
    191. }
    192. private void retryFailCallbackFile(){
    193. // valid
    194. File callbackLogPath = new File(failCallbackFilePath);
    195. if (!callbackLogPath.exists()) {
    196. return;
    197. }
    198. if (callbackLogPath.isFile()) {
    199. callbackLogPath.delete();
    200. }
    201. if (!(callbackLogPath.isDirectory() && callbackLogPath.list()!=null && callbackLogPath.list().length>0)) {
    202. return;
    203. }
    204. // load and clear file, retry
    205. for (File callbaclLogFile: callbackLogPath.listFiles()) {
    206. byte[] callbackParamList_bytes = FileUtil.readFileContent(callbaclLogFile);
    207. List callbackParamList = (List) XxlJobExecutor.getSerializer().deserialize(callbackParamList_bytes, HandleCallbackParam.class);
    208. callbaclLogFile.delete();
    209. doCallback(callbackParamList);
    210. }
    211. }
    212. }

    4.3 执行器注册到调度中心时序图

    4.4 执行器注册到调度中心核心代码解析

    4.4.1 执行注册(30秒一次)

    将执行器注册到调度中心:ExecutorRegistryThread

    1. package com.xxl.job.core.thread;
    2. import com.xxl.job.core.biz.AdminBiz;
    3. import com.xxl.job.core.biz.model.RegistryParam;
    4. import com.xxl.job.core.biz.model.ReturnT;
    5. import com.xxl.job.core.enums.RegistryConfig;
    6. import com.xxl.job.core.executor.XxlJobExecutor;
    7. import org.slf4j.Logger;
    8. import org.slf4j.LoggerFactory;
    9. import java.util.concurrent.TimeUnit;
    10. /**
    11. * Created by xuxueli on 17/3/2.
    12. * @Description:将执行器注册到调度中心
    13. */
    14. public class ExecutorRegistryThread {
    15. private static Logger logger = LoggerFactory.getLogger(ExecutorRegistryThread.class);
    16. private static ExecutorRegistryThread instance = new ExecutorRegistryThread();
    17. public static ExecutorRegistryThread getInstance(){
    18. return instance;
    19. }
    20. private Thread registryThread;
    21. private volatile boolean toStop = false;
    22. public void start(final String appName, final String address){
    23. // valid
    24. if (appName==null || appName.trim().length()==0) {
    25. logger.warn(">>>>>>>>>>> xxl-job, executor registry config fail, appName is null.");
    26. return;
    27. }
    28. if (XxlJobExecutor.getAdminBizList() == null) {
    29. logger.warn(">>>>>>>>>>> xxl-job, executor registry config fail, adminAddresses is null.");
    30. return;
    31. }
    32. registryThread = new Thread(new Runnable() {
    33. @Override
    34. public void run() {
    35. // registry
    36. while (!toStop) {
    37. try {
    38. RegistryParam registryParam = new RegistryParam(RegistryConfig.RegistType.EXECUTOR.name(), appName, address);
    39. for (AdminBiz adminBiz: XxlJobExecutor.getAdminBizList()) {
    40. try {
    41. ReturnT registryResult = adminBiz.registry(registryParam);
    42. if (registryResult!=null && ReturnT.SUCCESS_CODE == registryResult.getCode()) {
    43. registryResult = ReturnT.SUCCESS;
    44. logger.debug(">>>>>>>>>>> xxl-job registry success, registryParam:{}, registryResult:{}", new Object[]{registryParam, registryResult});
    45. break;
    46. } else {
    47. logger.info(">>>>>>>>>>> xxl-job registry fail, registryParam:{}, registryResult:{}", new Object[]{registryParam, registryResult});
    48. }
    49. } catch (Exception e) {
    50. logger.info(">>>>>>>>>>> xxl-job registry error, registryParam:{}", registryParam, e);
    51. }
    52. }
    53. } catch (Exception e) {
    54. if (!toStop) {
    55. logger.error(e.getMessage(), e);
    56. }
    57. }
    58. try {
    59. if (!toStop) {
    60. TimeUnit.SECONDS.sleep(RegistryConfig.BEAT_TIMEOUT);
    61. }
    62. } catch (InterruptedException e) {
    63. if (!toStop) {
    64. logger.warn(">>>>>>>>>>> xxl-job, executor registry thread interrupted, error msg:{}", e.getMessage());
    65. }
    66. }
    67. }
    68. // registry remove
    69. try {
    70. RegistryParam registryParam = new RegistryParam(RegistryConfig.RegistType.EXECUTOR.name(), appName, address);
    71. for (AdminBiz adminBiz: XxlJobExecutor.getAdminBizList()) {
    72. try {
    73. ReturnT registryResult = adminBiz.registryRemove(registryParam);
    74. if (registryResult!=null && ReturnT.SUCCESS_CODE == registryResult.getCode()) {
    75. registryResult = ReturnT.SUCCESS;
    76. logger.info(">>>>>>>>>>> xxl-job registry-remove success, registryParam:{}, registryResult:{}", new Object[]{registryParam, registryResult});
    77. break;
    78. } else {
    79. logger.info(">>>>>>>>>>> xxl-job registry-remove fail, registryParam:{}, registryResult:{}", new Object[]{registryParam, registryResult});
    80. }
    81. } catch (Exception e) {
    82. if (!toStop) {
    83. logger.info(">>>>>>>>>>> xxl-job registry-remove error, registryParam:{}", registryParam, e);
    84. }
    85. }
    86. }
    87. } catch (Exception e) {
    88. if (!toStop) {
    89. logger.error(e.getMessage(), e);
    90. }
    91. }
    92. logger.info(">>>>>>>>>>> xxl-job, executor registry thread destory.");
    93. }
    94. });
    95. registryThread.setDaemon(true);
    96. registryThread.setName("xxl-job, executor ExecutorRegistryThread");
    97. registryThread.start();
    98. }
    99. public void toStop() {
    100. toStop = true;
    101. // interrupt and wait
    102. registryThread.interrupt();
    103. try {
    104. registryThread.join();
    105. } catch (InterruptedException e) {
    106. logger.error(e.getMessage(), e);
    107. }
    108. }
    109. }

  • 相关阅读:
    安装插件时Vscode XHR Failed 报错ERR_CERT_AUTHORITY_INVALID
    结构型-代理模式
    Dubbo框架基本使用
    数据结构-二叉树的递归遍历
    Linux中Docker挂载mysql/mariadb等数据库,数据库问题汇总
    React核心原理与实际开发
    Complete Binary Tree
    [Linux 基础] 一篇带你了解linux权限问题
    js 对象循环遍历
    CV 面试指南—深度学习知识点总结(5)
  • 原文地址:https://blog.csdn.net/YYQ_QYY/article/details/126366987