xxl-job源码分析

一、数据库表梳理

image.png

1.1 源码目录介绍

  • /xxl-job-admin:调度中心(任务调度管理平台),项目源码
  • /xxl-job-core:公共Jar依赖(调度核心)
  • /xxl-job-executor-samples:执行器,Sample示例项目

1.2 数据库表介绍

  • xxl_job_lock:任务调度锁表;
  • xxl_job_registry:执行器注册表,维护在线的执行器和调度中心机器地址信息;
  • xxl_job_group:执行器信息表,维护任务执行器信息;
  • xxl_job_info:调度任务信息表:保存XXL-JOB调度任务的扩展信息,如任务分组、任务名、机器地址、执行器、执行入参和报警邮件等等;
  • xxl_job_log:调度日志表: 用于保存XXL-JOB任务调度的历史信息,如调度结果、执行结果、调度入参、调度机器和执行器等等;
  • xxl_job_log_report:调度日志报表:用户存储XXL-JOB任务调度日志的报表,调度中心报表功能页面会用到;
  • xxl_job_logglue:任务GLUE日志:用于保存GLUE更新历史,用于支持GLUE的版本回溯功能;
  • xxl_job_user:系统用户表;

二、调度中心xxl-job-admin

2.1 添加权限控制

权限注解@PermissionLimit,其实现的逻辑在PermissionInterceptor中,首先判断是否需要鉴权,如果需要则根据cookie中拿到的用户信息查库判断是否有权限登录,如果没有权限则重定向到登录页面或提示没有权限。

@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
public @interface PermissionLimit {
   
   /**
    * 登录拦截 (默认拦截)
    */
   boolean limit() default true;

   /**
    * 要求管理员权限
    *
    * @return
    */
   boolean adminuser() default false;

}
@Component
public class PermissionInterceptor implements AsyncHandlerInterceptor {

   @Resource
   private LoginService loginService;

   @Override
   public boolean preHandle(HttpServletRequest request, HttpServletResponse response, Object handler) throws Exception {
      //my-处理登录权限的逻辑
      if (!(handler instanceof HandlerMethod)) {
         return true;   // proceed with the next interceptor
      }

      // if need login
      boolean needLogin = true;
      boolean needAdminuser = false;
      HandlerMethod method = (HandlerMethod)handler;
      PermissionLimit permission = method.getMethodAnnotation(PermissionLimit.class);
      if (permission!=null) {
         needLogin = permission.limit();
         needAdminuser = permission.adminuser();
      }

      if (needLogin) {
         XxlJobUser loginUser = loginService.ifLogin(request, response);
         if (loginUser == null) {
            response.setStatus(302);
            response.setHeader("location", request.getContextPath()+"/toLogin");
            return false;
         }
         if (needAdminuser && loginUser.getRole()!=1) {
            throw new RuntimeException(I18nUtil.getString("system_permission_limit"));
         }
         request.setAttribute(LoginService.LOGIN_IDENTITY_KEY, loginUser);
      }

      return true;   // proceed with the next interceptor
   }
   
}

将PermissionInterceptor添加到web配置文件中

@Configuration
public class WebMvcConfig implements WebMvcConfigurer {

    @Resource
    private PermissionInterceptor permissionInterceptor;
    @Resource
    private CookieInterceptor cookieInterceptor;

    @Override
    public void addInterceptors(InterceptorRegistry registry) {
        registry.addInterceptor(permissionInterceptor).addPathPatterns("/**");
        registry.addInterceptor(cookieInterceptor).addPathPatterns("/**");
    }

}

2.2 配置中心初始化

xxljob的初始化和销毁动作在XxlJobAdminConfig中配置完成。

@Component
public class XxlJobAdminConfig implements InitializingBean, DisposableBean {

    private static XxlJobAdminConfig adminConfig = null;
    public static XxlJobAdminConfig getAdminConfig() {
        return adminConfig;
    }


    // ---------------------- XxlJobScheduler ----------------------

    private XxlJobScheduler xxlJobScheduler;

    @Override
    public void afterPropertiesSet() throws Exception {
        adminConfig = this;

        xxlJobScheduler = new XxlJobScheduler();
        xxlJobScheduler.init();
    }

    @Override
    public void destroy() throws Exception {
        xxlJobScheduler.destroy();
    }
}

2.3 XxlJobScheduler实例化

启动时通过XxlJobAdminConfig对XxlJobScheduler进行实例化

public class XxlJobScheduler  {
    private static final Logger logger = LoggerFactory.getLogger(XxlJobScheduler.class);

    public void init() throws Exception {
        // init i18n
        initI18n();

        // admin trigger pool start
        JobTriggerPoolHelper.toStart();

        // admin registry monitor run
        JobRegistryHelper.getInstance().start();

        // admin fail-monitor run
        JobFailMonitorHelper.getInstance().start();

        // admin lose-monitor run ( depend on JobTriggerPoolHelper )
        JobCompleteHelper.getInstance().start();

        // admin log report start
        JobLogReportHelper.getInstance().start();

        // start-schedule  ( depend on JobTriggerPoolHelper )
        JobScheduleHelper.getInstance().start();

        logger.info(">>>>>>>>> init xxl-job admin success.");
    }

    
    public void destroy() throws Exception {

        // stop-schedule
        JobScheduleHelper.getInstance().toStop();

        // admin log report stop
        JobLogReportHelper.getInstance().toStop();

        // admin lose-monitor stop
        JobCompleteHelper.getInstance().toStop();

        // admin fail-monitor stop
        JobFailMonitorHelper.getInstance().toStop();

        // admin registry stop
        JobRegistryHelper.getInstance().toStop();

        // admin trigger pool stop
        JobTriggerPoolHelper.toStop();

    }

    // ---------------------- I18n ----------------------

    private void initI18n(){
        for (ExecutorBlockStrategyEnum item:ExecutorBlockStrategyEnum.values()) {
            item.setTitle(I18nUtil.getString("jobconf_block_".concat(item.name())));
        }
    }
}

2.3.1 XxlJobScheduler中包含

  • JobTriggerPoolHelper:定时器线程池,基础线程池
  • JobRegistryHelper:注册线程池,通过拉取xxl_job_group配置进行任务执行器注册中心初始化
  • JobFailMonitorHelper:日志线程池
  • JobCompleteHelper:任务结果处理线程池,depend on JobTriggerPoolHelper
  • JobLogReportHelper:日志导出线程池
  • JobScheduleHelper:任务执行线程池,depend on JobTriggerPoolHelper
public class XxlJobScheduler  {

    public void init() throws Exception {
        // init i18n
        initI18n();

        // admin trigger pool start
        JobTriggerPoolHelper.toStart();

        // admin registry monitor run
        JobRegistryHelper.getInstance().start();

        // admin fail-monitor run
        JobFailMonitorHelper.getInstance().start();

        // admin lose-monitor run ( depend on JobTriggerPoolHelper )
        JobCompleteHelper.getInstance().start();

        // admin log report start
        JobLogReportHelper.getInstance().start();

        // start-schedule  ( depend on JobTriggerPoolHelper )
        JobScheduleHelper.getInstance().start();

        logger.info(">>>>>>>>> init xxl-job admin success.");
    }
}

2.4 JobRegistryHelper(维护和更新调度中心与执行器之间的注册信息)

2.4.1 调度中心管理注册信息

  • 1、先是初始化了注册或删除线程池(registryOrRemoveThreadPool),用来执行注册或者删除任务的
  • 2、然后创建了一个注册监控线程(registryMonitorThread),并启动

2.4.2 JobApiController

执行器调用调度中心的url来实现注册、下线、回调等操作;其主要的实现类是JobApiController

  • 调用/api/registry接口注册执行器信息
  • 调用/api/registryRemove接口下线执行器信息
  • 调用/api/callback接口执行回调操作
@Controller
@RequestMapping("/api")
public class JobApiController {

    @Resource
    private AdminBiz adminBiz;

    @RequestMapping("/{uri}")
    @ResponseBody
    @PermissionLimit(limit=false)
    public ReturnT<String> api(HttpServletRequest request, @PathVariable("uri") String uri, @RequestBody(required = false) String data) {
        // services mapping
        if ("callback".equals(uri)) {
            List<HandleCallbackParam> callbackParamList = GsonTool.fromJson(data, List.class, HandleCallbackParam.class);
            return adminBiz.callback(callbackParamList);
        } else if ("registry".equals(uri)) {
            RegistryParam registryParam = GsonTool.fromJson(data, RegistryParam.class);
            return adminBiz.registry(registryParam);
        } else if ("registryRemove".equals(uri)) {
            RegistryParam registryParam = GsonTool.fromJson(data, RegistryParam.class);
            return adminBiz.registryRemove(registryParam);
        } else {
            return new ReturnT<String>(ReturnT.FAIL_CODE, "invalid request, uri-mapping("+ uri +") not found.");
        }
    }
}

最终在JobRegistryHelper#registry方法中,registryOrRemoveThreadPool异步将 scheduler 服务注册到 xxl_job_registry表中

image.png

2.4.3 registryMonitorThread

每30秒执行一次,把 xxl_job_registry 表里的注册信息,同步刷新到 xxl_job_group 表里的 address_list列上去。


image.png
while (!toStop) {
    try {
        // auto registry group
        // 找到自动注册的信息
        List<XxlJobGroup> groupList = XxlJobAdminConfig.getAdminConfig().getXxlJobGroupDao().findByAddressType(0);
        if (groupList!=null && !groupList.isEmpty()) {

            // remove dead address (admin/executor)
            // 找到xxl_job_registry里超时的address,并删除
            List<Integer> ids = XxlJobAdminConfig.getAdminConfig().getXxlJobRegistryDao().findDead(RegistryConfig.DEAD_TIMEOUT, new Date());
            if (ids!=null && ids.size()>0) {
                XxlJobAdminConfig.getAdminConfig().getXxlJobRegistryDao().removeDead(ids);
            }

            // fresh online address (admin/executor)
            // 刷新在线的address
            HashMap<String, List<String>> appAddressMap = new HashMap<String, List<String>>();
            List<XxlJobRegistry> list = XxlJobAdminConfig.getAdminConfig().getXxlJobRegistryDao().findAll(RegistryConfig.DEAD_TIMEOUT, new Date());
            if (list != null) {
                // 把同一个appname的address归类
                for (XxlJobRegistry item: list) {
                    if (RegistryConfig.RegistType.EXECUTOR.name().equals(item.getRegistryGroup())) {
                        String appname = item.getRegistryKey();
                        List<String> registryList = appAddressMap.get(appname);
                        if (registryList == null) {
                            registryList = new ArrayList<String>();
                        }

                        if (!registryList.contains(item.getRegistryValue())) {
                            registryList.add(item.getRegistryValue());
                        }
                        appAddressMap.put(appname, registryList);
                    }
                }
            }

            // fresh group address
            for (XxlJobGroup group: groupList) {
                // 把address用逗号拼接起来
                List<String> registryList = appAddressMap.get(group.getAppname());
                String addressListStr = null;
                if (registryList!=null && !registryList.isEmpty()) {
                    Collections.sort(registryList);
                    StringBuilder addressListSB = new StringBuilder();
                    for (String item:registryList) {
                        addressListSB.append(item).append(",");
                    }
                    addressListStr = addressListSB.toString();
                    addressListStr = addressListStr.substring(0, addressListStr.length()-1);
                }
                group.setAddressList(addressListStr);
                group.setUpdateTime(new Date());
                
                //更新addressList
                XxlJobAdminConfig.getAdminConfig().getXxlJobGroupDao().update(group);
            }
        }
    } catch (Exception e) {
        if (!toStop) {
            logger.error(">>>>>>>>>>> xxl-job, job registry monitor thread error:{}", e);
        }
    }
    try {
        // 睡30秒
        TimeUnit.SECONDS.sleep(RegistryConfig.BEAT_TIMEOUT);
    } catch (InterruptedException e) {
        if (!toStop) {
            logger.error(">>>>>>>>>>> xxl-job, job registry monitor thread error:{}", e);
        }
    }
}
logger.info(">>>>>>>>>>> xxl-job, job registry monitor thread stop");

2.5 JobScheduleHelper(时间调度逻辑处理)

JobScheduleHelper是xxl-job框架的核心代码,是调度中心执行任务调度的核心代码,主要功能就是进行任务的调度,计算任务下一次执行的有效时间。

2.5.1 JobScheduleHelper有两个线程scheduleThread和ringThread

  • 1、scheduleThread负责筛选任务,挑选还没有触发或者快要触发的任务,挑选的方法是通过数据库查询job_info这个表,判断触发任务的时间是否小于当前时间+5s
  • 2、ringThread线程负责scheduleThread筛选出来的任务进行一个个执行。

2.5.2 JobScheduleHelper 任务调度器工作流程

  • 1、通过数据库的行锁和事务一致性,通过for update 来保证多个调度中心集群在同一时间内只有一个调度中心在调度任务
  • 2、周期性的遍历所有的jobInfo这个表,读取触发时间小于nowtime+5s这个时间之前的所有任务,然后进行引入以下触发机制判断
  • 3、三种触发任务机制:
    • 1)nowtime-TriggerNextTime() > PRE_READ_MS(5s) 既超过有效误差内,则查看当前任务的失效调度策略,若为立即重试一次,则立即触发调度任务,且触发类型为misfire
      1. nowtime-TriggerNextTime() < PRE_READ_MS(5s) 既没有超过有效误差,则立即调度调度任务
      1. nowtime < TriggerNextTime() 则说明这个任务马上就要触发了,放到一个时间轮上
  • 4、随后将快要触发的任务放到时间轮上,时间轮由key(将要触发的时间s),value(在当前触发s的所有任务id集合),然后更新这个任务的下一次触发时间
  • 5、这个时间轮的任务遍历交由第二个线程处理ringThread,周期在1s之内周期的扫描这个时间轮,然后执行调度任务。
public class JobScheduleHelper {
    private static Logger logger = LoggerFactory.getLogger(JobScheduleHelper.class);

    public static final long PRE_READ_MS = 5000;    // pre read

    private Thread scheduleThread;
    private Thread ringThread;

    // 时间轮的数据结构
    private volatile static Map<Integer, List<Integer>> ringData = new ConcurrentHashMap<>();

    public void start(){
        scheduleThread = new Thread(new Runnable() {
            @Override
            public void run() {

                // pre-read count: treadpool-size * trigger-qps (each trigger cost 50ms, qps = 1000/50 = 20)
                int preReadCount = (XxlJobAdminConfig.getAdminConfig().getTriggerPoolFastMax() + XxlJobAdminConfig.getAdminConfig().getTriggerPoolSlowMax()) * 20;

                while (!scheduleThreadToStop) {
                    long start = System.currentTimeMillis();

                    Connection conn = null;
                    Boolean connAutoCommit = null;
                    PreparedStatement preparedStatement = null;

                    boolean preReadSuc = true;
                    try {

                        conn = XxlJobAdminConfig.getAdminConfig().getDataSource().getConnection();
                        connAutoCommit = conn.getAutoCommit();
                        //关闭自动提交
                        conn.setAutoCommit(false);

                        preparedStatement = conn.prepareStatement(  "select * from xxl_job_lock where lock_name = 'schedule_lock' for update" );
                        preparedStatement.execute();

                        // tx start

                        // 1、pre read
                        long nowTime = System.currentTimeMillis();
                        //关键部分,预读出 当前时间+5s 的所有触发任务的时间小于这个时间预读出时间的任务
                        List<XxlJobInfo> scheduleList = XxlJobAdminConfig.getAdminConfig().getXxlJobInfoDao().scheduleJobQuery(nowTime + PRE_READ_MS, preReadCount);
                        if (scheduleList!=null && scheduleList.size()>0) {
                            // 2、push time-ring
                            for (XxlJobInfo jobInfo: scheduleList) {

                                // time-ring jump
                //遍历所有任务,过滤出所有过期任务
                    //过期任务的判断方式 当前时间>任务下一次触发时间+空闲间隔周期5s (在这段时间内任务还没触发的话说明任务超时过期了,既在5s的误差内还没出发任务)
                                if (nowTime > jobInfo.getTriggerNextTime() + PRE_READ_MS) {
                                    // 2.1、trigger-expire > 5s:pass && make next-trigger-time
                                    logger.warn(">>>>>>>>>>> xxl-job, schedule misfire, jobId = " + jobInfo.getId());

                                    // 1、misfire match 查询当前任务的过期策略
                                    MisfireStrategyEnum misfireStrategyEnum = MisfireStrategyEnum.match(jobInfo.getMisfireStrategy(), MisfireStrategyEnum.DO_NOTHING);
                                    if (MisfireStrategyEnum.FIRE_ONCE_NOW == misfireStrategyEnum) {
                                        // FIRE_ONCE_NOW 》 trigger 执行过期策略,过期立即执行一次
                                        JobTriggerPoolHelper.trigger(jobInfo.getId(), TriggerTypeEnum.MISFIRE, -1, null, null, null);
                                        logger.debug(">>>>>>>>>>> xxl-job, schedule push trigger : jobId = " + jobInfo.getId() );
                                    }

                                    // 2、fresh next
                    // 刷新接下来要执行时间
                                    refreshNextValidTime(jobInfo, new Date());
                //否则不是过期任务  任务下一次触发时间+空闲间隔周期5s>当前时间>下一次任务触发时间  (在允许的误差之内都不算过期任务)
                                } else if (nowTime > jobInfo.getTriggerNextTime()) {
                                    // 2.2、trigger-expire < 5s:direct-trigger && make next-trigger-time

                                    // 1、trigger
                   // 如果当前时间大于接下来要执行到时间则立即触发执行
                                    JobTriggerPoolHelper.trigger(jobInfo.getId(), TriggerTypeEnum.CRON, -1, null, null, null);
                                    logger.debug(">>>>>>>>>>> xxl-job, schedule push trigger : jobId = " + jobInfo.getId() );

                                    // 2、fresh next
                                    // 刷新下次执行时间
                                    refreshNextValidTime(jobInfo, new Date());

                                    // next-trigger-time in 5s, pre-read again
                                    if (jobInfo.getTriggerStatus()==1 && nowTime + PRE_READ_MS > jobInfo.getTriggerNextTime()) {

                                        // 1、make ring second
                    // 如果接下来 5 秒内还执行
                                        int ringSecond = (int)((jobInfo.getTriggerNextTime()/1000)%60);

                                        // 2、push time ring
                    // 则直接放到时间轮中
                                        pushTimeRing(ringSecond, jobInfo.getId());

                                        // 3、fresh next
                    // 刷新下次执行时间
                                        refreshNextValidTime(jobInfo, new Date(jobInfo.getTriggerNextTime()));

                                    }

                                } else {
                    //未来5秒内的,任务马上就要触发了,加入时间轮上,等待任务执行
                                    // 2.3、trigger-pre-read:time-ring trigger && make next-trigger-time

                                    // 1、make ring second
                    // 任务还没有到执行时间则直接放到时间轮中
                                    int ringSecond = (int)((jobInfo.getTriggerNextTime()/1000)%60);

                                    // 2、push time ring
                    // 则直接放到时间轮中
                                    pushTimeRing(ringSecond, jobInfo.getId());

                                    // 3、fresh next
                    // 刷新下次执行时间
                                    refreshNextValidTime(jobInfo, new Date(jobInfo.getTriggerNextTime()));

                                }

                            }

                            // 3、update trigger info
                            for (XxlJobInfo jobInfo: scheduleList) {
                // 更新任务信息
                                XxlJobAdminConfig.getAdminConfig().getXxlJobInfoDao().scheduleUpdate(jobInfo);
                            }

                        } else {
                            preReadSuc = false;
                        }
                    } catch (Exception e) {
                    } finally {
                    }
                }
                logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#scheduleThread stop");
            }
        });
        scheduleThread.setDaemon(true);
        scheduleThread.setName("xxl-job, admin JobScheduleHelper#scheduleThread");
        scheduleThread.start();
        
        ringThread = new Thread(new Runnable() {
            @Override
            public void run() {

                while (!ringThreadToStop) {

                    // align second
                    try {
                        TimeUnit.MILLISECONDS.sleep(1000 - System.currentTimeMillis() % 1000);
                    } catch (InterruptedException e) {
                        if (!ringThreadToStop) {
                            logger.error(e.getMessage(), e);
                        }
                    }

                    try {
                        // second data
                        List<Integer> ringItemData = new ArrayList<>();
                        int nowSecond = Calendar.getInstance().get(Calendar.SECOND);   // 避免处理耗时太长,跨过刻度,向前校验一个刻度;
                        for (int i = 0; i < 2; i++) {
                            List<Integer> tmpData = ringData.remove( (nowSecond+60-i)%60 );
                            if (tmpData != null) {
                                ringItemData.addAll(tmpData);
                            }
                        }

                        // ring trigger
                        logger.debug(">>>>>>>>>>> xxl-job, time-ring beat : " + nowSecond + " = " + Arrays.asList(ringItemData) );
                        if (ringItemData.size() > 0) {
                            // do trigger
                            for (int jobId: ringItemData) {
                                // do trigger
                                JobTriggerPoolHelper.trigger(jobId, TriggerTypeEnum.CRON, -1, null, null, null);
                            }
                            // clear
                            ringItemData.clear();
                        }
                    } catch (Exception e) {
                        if (!ringThreadToStop) {
                            logger.error(">>>>>>>>>>> xxl-job, JobScheduleHelper#ringThread error:{}", e);
                        }
                    }
                }
                logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#ringThread stop");
            }
        });
        ringThread.setDaemon(true);
        ringThread.setName("xxl-job, admin JobScheduleHelper#ringThread");
        ringThread.start();
    }
}

2.5.3 scheduleThread 总结:

  • 1、整秒处运行,两次间隔最多5秒,5秒内最多执行一次;
  • 2、过期超过5s的按照job配置的过期策略执行,
  • 3、过期5s内的立即执行,并且看下次执行时间是不是在5s内,在的话,放到ringData中等待ring线程处理,
  • 4、如果是在未来5s内的这种正常情况,就直接放到ringData中等待ring线程处理

2.5.4 ringThread 总结:

  • 1、整秒处运行,两次间隔最多1秒,1秒内最多执行一次;
  • 2、首先ringData,时间轮算法数据结构,其实就是个 Map<Integer, List<Integer>> ringData = new ConcurrentHashMap<>()
  • 3、时间轮由key(将要触发的时间s),value(在当前触发s的所有任务id集合)
  • 4、根据秒数刻度分类存储,然后再通过ringThread处理,处理的时候,为了避免处理时间长,往前多取了一个刻度

2.5.5 JobScheduleHelper总结:

  • 采用DB表行锁,在调度平台多节点间,避免并发调度;
  • 两个线程,scheduleThread线程查表,获取将到期任务,放入时间轮中,交由ringThread线程来处理;这种设计在定时任务的实现中常被采用;
  • 任务执行可能出现一个周期的延迟:
      1. 如调度平台节点A,将5秒内到期的任务放入ringData,然后宕机了,且TriggerNextTime已被修改;
    • 2)节点B紧接着执行时,得等到下一个到期时间,才能查询到这个任务;
    • 举例:某任务(周期1小时)20:55:00到期,20:54:56时被节点A查到,放入ringData,更新triggerLastTime为21:55:00,但20:54:58时节点A宕机了。该任务只能等到21:55:00才被执行。

2.6 JobTriggerPoolHelper(任务触发器)

image.png

2.6.1 从上面的方法中可以总结出 JobTriggerPoolHelper 大致提供了一下三个功能:

  • 启动处理线程
  • 终止处理线程
  • 进行任务触发

2.6.2 触发方式

通过查看 JobTriggerPoolHelper 中 trigger() 方法的使用者,我们可以看到有一下五种触发任务的场景:

  • 1、在调度中心页面中触发一次任务;
  • 2、由调度中心根据时间调度进行任务触发;
  • 3、任务失败重新进行任务触发;
  • 4、父任务完成进行子任务触发;
  • 5、通过API调用进行任务触发;

也可以通过触发类型枚举类 TriggerTypeEnum 来查看

/**
 * 人工触发
 */
MANUAL(I18nUtil.getString("jobconf_trigger_type_manual")),
/**
 * 根据cron表达式触发
 */
CRON(I18nUtil.getString("jobconf_trigger_type_cron")),
/**
 * 失败重试触发
 */
RETRY(I18nUtil.getString("jobconf_trigger_type_retry")),
/**
 * 父任务触发
 */
PARENT(I18nUtil.getString("jobconf_trigger_type_parent")),
/**
 * API调用触发
 */
API(I18nUtil.getString("jobconf_trigger_type_api"));
/**
 * 任务触发失败(misfire)时,jobconf_trigger_type_misfire策略决定后续处理方式:立即重试、等待下次触发、终止任务执行
 */
MISFIRE(I18nUtil.getString("jobconf_trigger_type_misfire"));

2.6.3 快任务线程池和慢任务线程池

  • xxl-job 使用了线程池来进行任务调度,一旦出现某个任务调度时间过长致使线程阻塞就会导致调度中心调度效率的下降。
  • 为了解决这一问题,xxl-job 创建了 「快任务线程池 fastTriggerPool」 和 「慢任务线程池 slowTriggerPool」 。
  • 任务默认放置在快任务线程池中进行任务触发。
  • xxl-job 设置了一个任务触发时间窗口,长度为500ms。触发器在任务触发过程中每分钟检查当前任务已触发时间,如果超过时间窗口的次数超过10次,则会将该任务降级到慢任务线程池中。

具体处理逻辑如下:

// 线程池选择逻辑
ThreadPoolExecutor triggerPool_ = fastTriggerPool;
AtomicInteger jobTimeoutCount = jobTimeoutCountMap.get(jobId);
if (jobTimeoutCount!=null && jobTimeoutCount.get() > 10) {
    // job-timeout 10 times in 1 min
    triggerPool_ = slowTriggerPool;
}

...
// 时间窗口检查逻辑
// check timeout-count-map
long minTim_now = System.currentTimeMillis()/60000;
if (minTim != minTim_now) {
    minTim = minTim_now;
    jobTimeoutCountMap.clear();
}

// incr timeout-count-map
long cost = System.currentTimeMillis()-start;
if (cost > 500) {
    // ob-timeout threshold 500ms
    AtomicInteger timeoutCount = jobTimeoutCountMap.putIfAbsent(jobId, new AtomicInteger(1));
    if (timeoutCount != null) {
        timeoutCount.incrementAndGet();
    }
}

2.6.4 任务触发

任务触发是通过快任务线程池/慢任务线程池调用 XxlJobTrigger 的 trigger() 方法实现。具体代码如下:

public static void trigger(int jobId,
                           TriggerTypeEnum triggerType,
                           int failRetryCount,
                           String executorShardingParam,
                           String executorParam,
                           String addressList) {

    // 从 xxl_job_info 获取任务信息
    XxlJobInfo jobInfo = XxlJobAdminConfig.getAdminConfig().getXxlJobInfoDao().loadById(jobId);
    if (jobInfo == null) {
        logger.warn(">>>>>>>>>>>> trigger fail, jobId invalid,jobId={}", jobId);
        return;
    }
    if (executorParam != null) {
        jobInfo.setExecutorParam(executorParam);
    }
    int finalFailRetryCount = failRetryCount>=0?failRetryCount:jobInfo.getExecutorFailRetryCount();
    // 根据执行器ID,获取执行器信息,包含 address_list 地址信息
    XxlJobGroup group = XxlJobAdminConfig.getAdminConfig().getXxlJobGroupDao().load(jobInfo.getJobGroup());
    
    // TODO 可以在这里加入发布策略的查询,如果存在发布策略,则从 address_list 里面匹配相关ip进行后续路由执行
    // TODO 如果不存在发布策略,则之间走原有的 address_list 进行后续路由执行

    // cover addressList
    if (addressList!=null && addressList.trim().length()>0) {
        group.setAddressType(1);
        group.setAddressList(addressList.trim());
    }

    // sharding param
    int[] shardingParam = null;
    if (executorShardingParam!=null){
        // 判断分片参数 为空或者格式不对时,设置默认分片参数: 0/1 (表示总执行器一台,第一台执行)
        String[] shardingArr = executorShardingParam.split("/");
        if (shardingArr.length==2 && isNumeric(shardingArr[0]) && isNumeric(shardingArr[1])) {
            shardingParam = new int[2];
            shardingParam[0] = Integer.valueOf(shardingArr[0]);
            shardingParam[1] = Integer.valueOf(shardingArr[1]);
        }
    }
    // 判断路由策略,是否分片广播执行
    if (ExecutorRouteStrategyEnum.SHARDING_BROADCAST==ExecutorRouteStrategyEnum.match(jobInfo.getExecutorRouteStrategy(), null)
            && group.getRegistryList()!=null && !group.getRegistryList().isEmpty()
            && shardingParam==null) {
        for (int i = 0; i < group.getRegistryList().size(); i++) {
            processTrigger(group, jobInfo, finalFailRetryCount, triggerType, i, group.getRegistryList().size());
        }
    } else {
        if (shardingParam == null) {
            shardingParam = new int[]{0, 1};
        }
        processTrigger(group, jobInfo, finalFailRetryCount, triggerType, shardingParam[0], shardingParam[1]);
    }
}

2.6.5 XxlJobTrigger#processTrigger 根据路由策略,选择调度地址,远程调用执行器触发任务执行

private static void processTrigger(XxlJobGroup group, XxlJobInfo jobInfo, int finalFailRetryCount, TriggerTypeEnum triggerType, int index, int total){

    // param
    ExecutorBlockStrategyEnum blockStrategy = ExecutorBlockStrategyEnum.match(jobInfo.getExecutorBlockStrategy(), ExecutorBlockStrategyEnum.SERIAL_EXECUTION);  // block strategy
    ExecutorRouteStrategyEnum executorRouteStrategyEnum = ExecutorRouteStrategyEnum.match(jobInfo.getExecutorRouteStrategy(), null);    // route strategy
    String shardingParam = (ExecutorRouteStrategyEnum.SHARDING_BROADCAST==executorRouteStrategyEnum)?String.valueOf(index).concat("/").concat(String.valueOf(total)):null;

    // 1、save log-id
    XxlJobLog jobLog = new XxlJobLog();
    jobLog.setJobGroup(jobInfo.getJobGroup());
    jobLog.setJobId(jobInfo.getId());
    jobLog.setTriggerTime(new Date());
    XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().save(jobLog);
    logger.debug(">>>>>>>>>>> xxl-job trigger start, jobId:{}", jobLog.getId());

    // 2、init trigger-param 任务参数复制
    TriggerParam triggerParam = new TriggerParam();
    triggerParam.setJobId(jobInfo.getId());
    triggerParam.setExecutorHandler(jobInfo.getExecutorHandler());
    triggerParam.setExecutorParams(jobInfo.getExecutorParam());
    triggerParam.setExecutorBlockStrategy(jobInfo.getExecutorBlockStrategy());
    triggerParam.setExecutorTimeout(jobInfo.getExecutorTimeout());
    triggerParam.setLogId(jobLog.getId());
    triggerParam.setLogDateTime(jobLog.getTriggerTime().getTime());
    triggerParam.setGlueType(jobInfo.getGlueType());
    triggerParam.setGlueSource(jobInfo.getGlueSource());
    triggerParam.setGlueUpdatetime(jobInfo.getGlueUpdatetime().getTime());
    triggerParam.setBroadcastIndex(index);
    triggerParam.setBroadcastTotal(total);

    // 3、init address 根据路由策略选择调度地址
    String address = null;
    ReturnT<String> routeAddressResult = null;
    if (group.getRegistryList()!=null && !group.getRegistryList().isEmpty()) {
        if (ExecutorRouteStrategyEnum.SHARDING_BROADCAST == executorRouteStrategyEnum) {
            if (index < group.getRegistryList().size()) {
                address = group.getRegistryList().get(index);
            } else {
                address = group.getRegistryList().get(0);
            }
        } else {
            routeAddressResult = executorRouteStrategyEnum.getRouter().route(triggerParam, group.getRegistryList());
            if (routeAddressResult.getCode() == ReturnT.SUCCESS_CODE) {
                address = routeAddressResult.getContent();
            }
        }
    } else {
        routeAddressResult = new ReturnT<String>(ReturnT.FAIL_CODE, I18nUtil.getString("jobconf_trigger_address_empty"));
    }

    // 4、trigger remote executor
    ReturnT<String> triggerResult = null;
    if (address != null) {
        //触发任务执行
        triggerResult = runExecutor(triggerParam, address);
    } else {
        triggerResult = new ReturnT<String>(ReturnT.FAIL_CODE, null);
    }

    // 5、collection trigger info
    // 日志信息拼接, 省略
    XxlJobAdminConfig.getAdminConfig().getXxlJobLogDao().updateTriggerInfo(jobLog);

    logger.debug(">>>>>>>>>>> xxl-job trigger end, jobId:{}", jobLog.getId());
}

2.6.6 路由执行策略

  • 执行器集群部署时提供丰富的路由策略,包括:第一个、最后一个、轮询、随机、一致性HASH、最不经常使用、最近最久未使用、故障转移、忙碌转移等;
  • 第一个、最后一个、轮询、随机:都是简单读address_list即可
  • 一致性HASH:TreeSet实现一致性hash算法
  • 最不经常使用、最近最久未使用:HashMap、LinkedHashMap
  • 故障转移:遍历address_list获取address时,逐个检查该address的心跳(请求返回状态);只有心跳正常的address才返回使用
  • 忙碌转移:遍历address_list获取address时,逐个检查该address是否忙碌(请求返回状态);只有状态为idle的address才返回使用
image.png

2.7 流程图

xxl-job.jpg

三、时间轮

时间轮出自Netty中的HashedWheelTimer,是一个环形结构,可以用时钟来类比,钟面上有很多bucket,每一个bucket上可以存放多个任务,使用一个List保存该时刻到期的所有任务,同时一个指针随着时间流逝一格一格转动,并执行对应bucket上所有到期的任务。任务通过取模 决定应该放入哪个bucket。和HashMap的原理类似,newTask对应put,使用List来解决 Hash 冲突。

image.png

以上图为例,假设一个bucket是1秒,则指针转动一轮表示的时间段为8s,假设当前指针指向 0,此时需要调度一个3s后执行的任务,显然应该加入到(0+3=3)的方格中,指针再走3s次就可以执行了;如果任务要在10s后执行,应该等指针走完一轮零2格再执行,因此应放入2,同时将round(1)保存到任务中。检查到期任务时只执行round为0的,bucket上其他任务的round减1。

当然,还有优化的“分层时间轮”的实现,请参考https://cnkirito.moe/timer/

3.1 XXL-JOB中的时间轮

  • XXL-JOB中的调度方式从Quartz变成了自研调度的方式,很像时间轮,可以理解为有60个bucket且每个bucket为1秒,但是没有了round的概念。
  • 具体可以看下图。
image.png

XXL-JOB中负责任务调度的有两个线程,分别为ringThread和scheduleThread,其作用如下

  • 1、scheduleThread:对任务信息进行读取,预读未来5s 即将触发的任务,放入时间轮。
  • 2、ringThread:对当前bucket和前一个bucket中的任务取出并执行。
public class JobScheduleHelper {

    // 环状结构
    private volatile static Map<Integer, List<Integer>> ringData = new ConcurrentHashMap<>();
    
    // 任务下次启动时间(单位为秒) % 60
    int ringSecond = (int)((jobInfo.getTriggerNextTime()/1000)%60);

    public void start(){
        // ring thread
        ringThread = new Thread(new Runnable() {
            @Override
            public void run() {
                while (!ringThreadToStop) {
                    try {
                        // second data
                        // 同时取两个时间刻度的任务
                        List<Integer> ringItemData = new ArrayList<>();
                        int nowSecond = Calendar.getInstance().get(Calendar.SECOND);
                        // 避免处理耗时太长,跨过刻度,向前校验一个刻度;                      
                        for (int i = 0; i < 2; i++) {
                            List<Integer> tmpData = ringData.remove( (nowSecond+60-i)%60 );
                            if (tmpData != null) {
                                ringItemData.addAll(tmpData);
                            }
                        }

                        // ring trigger
                        logger.debug(">>>>>>>>>>> xxl-job, time-ring beat : " + nowSecond + " = " + Arrays.asList(ringItemData) );
                        if (ringItemData.size() > 0) {
                            // do trigger
                            for (int jobId: ringItemData) {
                                // do trigger  运行
                                JobTriggerPoolHelper.trigger(jobId, TriggerTypeEnum.CRON, -1, null, null, null);
                            }
                            // clear
                            ringItemData.clear();
                        }
                    } catch (Exception e) {
                        if (!ringThreadToStop) {
                            logger.error(">>>>>>>>>>> xxl-job, JobScheduleHelper#ringThread error:{}", e);
                        }
                    }
                }
                logger.info(">>>>>>>>>>> xxl-job, JobScheduleHelper#ringThread stop");
            }
        });
        ringThread.setDaemon(true);
        ringThread.setName("xxl-job, admin JobScheduleHelper#ringThread");
        ringThread.start();
    }

    // 任务放进时间轮
    private void pushTimeRing(int ringSecond, int jobId){
        // push async ring
        List<Integer> ringItemData = ringData.get(ringSecond);
        if (ringItemData == null) {
            ringItemData = new ArrayList<Integer>();
            ringData.put(ringSecond, ringItemData);
        }
        ringItemData.add(jobId);

        logger.debug(">>>>>>>>>>> xxl-job, schedule push time-ring : " + ringSecond + " = " + Arrays.asList(ringItemData) );
    }
}

四、一致性Hash路由中的Hash算法

  • 大家也知道,XXL-JOB在执行任务时,任务具体在哪个执行器上运行是根据路由策略来决定的,其中有一个策略是一致性Hash策略(源码在ExecutorRouteConsistentHash.java),自然而然想到了一致性Hash算法 。
  • 一致性Hash算法 是为了解决分布式系统中负载均衡的问题时候可以使用Hash算法让固定的一部分请求落到同一台服务器上,这样每台服务器固定处理一部分请求(并维护这些请求的信息),起到负载均衡的作用。
  • 普通的余数hash(hash(比如用户id)%服务器机器数)算法伸缩性很差,当新增或者下线服务器机器时候,用户id与服务器的映射关系会大量失效。一致性hash则利用hash环对其进行了改进。
  • 一致性Hash算法 在实践中,当服务器节点比较少的时候会出现上节所说的一致性hash倾斜的问题,一个解决方法是多加机器,但是加机器是有成本的,那么就加虚拟节点 。
  • 下图为带有虚拟节点的Hash环,其中ip1-1是ip1的虚拟节点,ip2-1是ip2的虚拟节点,ip3-1是ip3的虚拟节点。


    image.png

可见 ,一致性Hash算法的关键在于Hash算法 ,保证虚拟节点 及Hash结果 的均匀性,而均匀性可以理解为减少Hash冲突 。

  • XXL-JOB中的一致性Hash的Hash函数如下。
// jobId转换为md5
// 不直接用hashCode() 是因为扩大hash取值范围,减少冲突
byte[] digest = md5.digest();
 
// 32位hashCode
long hashCode = ((long) (digest[3] & 0xFF) << 24)
 | ((long) (digest[2] & 0xFF) << 16)
 | ((long) (digest[1] & 0xFF) << 8)
 | (digest[0] & 0xFF);
 
long truncateHashCode = hashCode & 0xffffffffL;
  • 看到上图的Hash函数,让我想到了HashMap的Hash函数
f(key) = hash(key) & (table.length - 1) 
// 使用>>> 16的原因,hashCode()的高位和低位都对f(key)有了一定影响力,使得分布更加均匀,散列冲突的几率就小了。
hash(key) = (h = key.hashCode()) ^ (h >>> 16)
  • 同理,将jobId的md5编码的高低位都对Hash结果有影响,使得Hash冲突的概率减小。

五、分片任务的实现 - 维护线程上下文

  • XXL-JOB的分片任务实现了任务的分布式执行,其实是笔者调研的重点,日常开发中很多定时任务都是单机执行,对于后续数据量大的任务最好有一个分布式的解决方案。
  • 分片任务的路由策略,源代码作者提出了分片广播 的概念,刚开始还有点摸不清头脑,看了源码逐渐清晰了起来。
public enum ExecutorRouteStrategyEnum {
    FIRST(I18nUtil.getString("jobconf_route_first"), new ExecutorRouteFirst()),
    LAST(I18nUtil.getString("jobconf_route_last"), new ExecutorRouteLast()),
    ROUND(I18nUtil.getString("jobconf_route_round"), new ExecutorRouteRound()),
    RANDOM(I18nUtil.getString("jobconf_route_random"), new ExecutorRouteRandom()),
    CONSISTENT_HASH(I18nUtil.getString("jobconf_route_consistenthash"), new ExecutorRouteConsistentHash()),
    LEAST_FREQUENTLY_USED(I18nUtil.getString("jobconf_route_lfu"), new ExecutorRouteLFU()),
    LEAST_RECENTLY_USED(I18nUtil.getString("jobconf_route_lru"), new ExecutorRouteLRU()),
    FAILOVER(I18nUtil.getString("jobconf_route_failover"), new ExecutorRouteFailover()),
    BUSYOVER(I18nUtil.getString("jobconf_route_busyover"), new ExecutorRouteBusyover()),
    // 说好的实现呢???竟然是null
    SHARDING_BROADCAST(I18nUtil.getString("jobconf_route_shard"), null);
  • 再继续追查得到了结论,待我慢慢道来,首先分片任务执行参数传递的是什么?看XxlJobTrigger.trigger函数中的一段代码。
...
// 如果是分片路由,走的是这段逻辑
if (ExecutorRouteStrategyEnum.SHARDING_BROADCAST == ExecutorRouteStrategyEnum.match(jobInfo.getExecutorRouteStrategy(), null)
                && group.getRegistryList() != null && !group.getRegistryList().isEmpty()
                && shardingParam == null) {
            for (int i = 0; i < group.getRegistryList().size(); i++) {
             // 最后两个参数,i是当前机器在执行器集群当中的index,group.getRegistryList().size()为执行器总数
                processTrigger(group, jobInfo, finalFailRetryCount, triggerType, i, group.getRegistryList().size());
            }
        } 
...
  • 参数经过自研RPC传递到执行器,在执行器中具体负责任务执行的JobThread.run中,看到了如下代码。
// 分片广播的参数比set进了ShardingUtil
ShardingUtil.setShardingVo(new ShardingUtil.ShardingVO(triggerParam.getBroadcastIndex(), triggerParam.getBroadcastTotal()));
...
// 将执行参数传递给jobHandler执行
handler.execute(triggerParamTmp.getExecutorParams())
  • 接着看ShardingUtil,才发现了其中的奥秘,请看代码。
public class ShardingUtil {
 // 线程上下文
    private static InheritableThreadLocal<ShardingVO> contextHolder = new InheritableThreadLocal<ShardingVO>();
 // 分片参数对象
    public static class ShardingVO {
 
        private int index;  // sharding index
        private int total;  // sharding total
  // 次数省略 get/set
    }
 // 参数对象注入上下文
    public static void setShardingVo(ShardingVO shardingVo){
        contextHolder.set(shardingVo);
    }
 // 从上下文中取出参数对象
    public static ShardingVO getShardingVo(){
        return contextHolder.get();
    }
}
  • 显而易见,在负责分片任务的ShardingJobHandler里取出了线程上下文中的分片参数
@JobHandler(value="shardingJobHandler")
@Service
public class ShardingJobHandler extends IJobHandler {
 @Override
 public ReturnT<String> execute(String param) throws Exception {
 
  // 分片参数
  ShardingUtil.ShardingVO shardingVO = ShardingUtil.getShardingVo();
  XxlJobLogger.log("分片参数:当前分片序号 = {}, 总分片数 = {}", shardingVO.getIndex(), shardingVO.getTotal());
 
  // 业务逻辑
for (int i = 0; i < shardingVO.getTotal(); i++) {
   if (i == shardingVO.getIndex()) {
    XxlJobLogger.log("第 {} 片, 命中分片开始处理", i);
   } else {
    XxlJobLogger.log("第 {} 片, 忽略", i);
   }
  }
 
return SUCCESS;
 }
}
  • 由此得出,分布式实现是根据分片参数index及total来做的,简单来讲,就是给出了当前执行器的标识,根据这个标识将任务的数据或者逻辑进行区分,即可实现分布式运行。

题外话:至于为什么用外部注入分片参数的方式,不直接execute传递?

  • 1、可能是因为只有分片任务才用到这两个参数
  • 2、IJobHandler只有String类型参数

参考:
https://blog.csdn.net/l123lgx/article/details/136399697

https://blog.csdn.net/m0_73735578/article/details/148900849

https://blog.csdn.net/s6056826a/article/details/113446126

最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

相关阅读更多精彩内容

友情链接更多精彩内容