xxl-job源码分析
xxl-job
系統(tǒng)說明
安裝
安裝部署參考文檔:分布式任務(wù)調(diào)度平臺xxl-job
功能
定時調(diào)度、服務(wù)解耦、靈活控制跑批時間(停止、開啟、重新設(shè)定時間、手動觸發(fā))
XXL-JOB是一個輕量級分布式任務(wù)調(diào)度平臺,其核心設(shè)計(jì)目標(biāo)是開發(fā)迅速、學(xué)習(xí)簡單、輕量級、易擴(kuò)展。現(xiàn)已開放源代碼并接入多家公司線上產(chǎn)品線,開箱即用
概念
執(zhí)行器列表:一個執(zhí)行器是一個項(xiàng)目
任務(wù):一個任務(wù)是一個項(xiàng)目中的 JobHandler
一個xxl-job服務(wù)可以有多個執(zhí)行器(項(xiàng)目),一個項(xiàng)目下可以有多個任務(wù)(JobHandler),他們是如何關(guān)聯(lián)的?
頁面操作:
代碼操作:
架構(gòu)圖
拋出疑問
系統(tǒng)分析
執(zhí)行器依賴jar包
com.xuxueli:xxl-job-core:2.1.0
com.xuxueli:xxl-registry-client:1.0.2
com.xuxueli:xxl-rpc-core:1.4.1
調(diào)度中心啟動過程
// 1. 加載 XxlJobAdminConfig,adminConfig = this XxlJobAdminConfig.java// 啟動過程代碼 @Component public class XxlJobScheduler implements InitializingBean, DisposableBean {private static final Logger logger = LoggerFactory.getLogger(XxlJobScheduler.class);@Overridepublic void afterPropertiesSet() throws Exception {// init i18ninitI18n();// admin registry monitor run// 2. 啟動注冊監(jiān)控器(將注冊到register表中的IP加載到group表)/ 30執(zhí)行一次JobRegistryMonitorHelper.getInstance().start();// admin monitor run// 3. 啟動失敗日志監(jiān)控器(失敗重試,失敗郵件發(fā)送)JobFailMonitorHelper.getInstance().start();// admin-server// 4. 初始化RPC服務(wù)initRpcProvider();// start-schedule// 5. 啟動定時任務(wù)調(diào)度器(執(zhí)行任務(wù),緩存任務(wù))JobScheduleHelper.getInstance().start();logger.info(">>>>>>>>> init xxl-job admin success.");}...... }執(zhí)行器啟動過程
@Override public void start() throws Exception {// init JobHandler Repository// 將執(zhí)行 JobHandler 注冊到緩存中 jobHandlerRepository(ConcurrentMap)initJobHandlerRepository(applicationContext);// refresh GlueFactory// 刷新GLUEGlueFactory.refreshInstance(1);// super start// 核心啟動項(xiàng)super.start(); }public void start() throws Exception {// 初始化日志路徑 // private static String logBasePath = "/data/applogs/xxl-job/jobhandler";XxlJobFileAppender.initLogPath(this.logPath);// 初始化注冊中心列表 (把注冊地址放到 List)this.initAdminBizList(this.adminAddresses, this.accessToken);// 啟動日志文件清理線程 (一天清理一次)// 每天清理一次過期日志,配置參數(shù)必須大于3才有效JobLogFileCleanThread.getInstance().start((long)this.logRetentionDays);// 開啟觸發(fā)器回調(diào)線程TriggerCallbackThread.getInstance().start();// 指定端口this.port = this.port > 0 ? this.port : NetUtil.findAvailablePort(9999);// 指定IPthis.ip = this.ip != null && this.ip.trim().length() > 0 ? this.ip : IpUtil.getIp();// 初始化RPC 將執(zhí)行器注冊到調(diào)度中心 30秒一次this.initRpcProvider(this.ip, this.port, this.appName, this.accessToken); }執(zhí)行器注冊到調(diào)度中心
執(zhí)行器
// 注冊執(zhí)行器入口 XxlJobExecutor.java->initRpcProvider()->xxlRpcProviderFactory.start();// 開啟注冊 XxlRpcProviderFactory.java->start();// 執(zhí)行注冊 ExecutorRegistryThread.java->start(); // RPC 注冊代碼 for (AdminBiz adminBiz: XxlJobExecutor.getAdminBizList()) {try {ReturnT<String> registryResult = adminBiz.registry(registryParam);if (registryResult!=null && ReturnT.SUCCESS_CODE == registryResult.getCode()) {registryResult = ReturnT.SUCCESS;logger.debug(">>>>>>>>>>> xxl-job registry success, registryParam:{}, registryResult:{}", new Object[]{registryParam, registryResult});break;} else {logger.info(">>>>>>>>>>> xxl-job registry fail, registryParam:{}, registryResult:{}", new Object[]{registryParam, registryResult});}} catch (Exception e) {logger.info(">>>>>>>>>>> xxl-job registry error, registryParam:{}", registryParam, e);}}調(diào)度中心
// RPC 注冊服務(wù) AdminBizImpl.java->registry();數(shù)據(jù)庫
調(diào)度中心調(diào)用執(zhí)行器
/* 調(diào)度中心執(zhí)行步驟 */ // 1. 調(diào)用執(zhí)行器 XxlJobTrigger.java->runExecutor();// 2. 獲取執(zhí)行器 XxlJobScheduler.java->getExecutorBiz();// 3. 調(diào)用 ExecutorBizImpl.java->run();/* 執(zhí)行器執(zhí)行步驟 */ // 1. 執(zhí)行器接口 ExecutorBiz.java->run();// 2. 執(zhí)行器實(shí)現(xiàn) ExecutorBizImpl.java->run();// 3. 把jobInfo 從 jobThreadRepository (ConcurrentMap) 中獲取一個新線程,并開啟新線程 XxlJobExecutor.java->registJobThread();// 4. 保存到當(dāng)前線程隊(duì)列 JobThread.java->pushTriggerQueue();// 5. 執(zhí)行 JobThread.java->handler.execute(triggerParam.getExecutorParams());調(diào)度中心(Admin)
實(shí)現(xiàn) org.springframework.beans.factory.InitializingBean類,重寫 afterPropertiesSet 方法,在初始化bean的時候都會執(zhí)行該方法
DisposableBean spring停止時執(zhí)行
結(jié)束加載項(xiàng)
手動執(zhí)行方式
JobInfoController.java
@RequestMapping("/trigger") @ResponseBody //@PermissionLimit(limit = false) public ReturnT<String> triggerJob(int id, String executorParam) {// force cover job paramif (executorParam == null) {executorParam = "";}JobTriggerPoolHelper.trigger(id, TriggerTypeEnum.MANUAL, -1, null, executorParam);return ReturnT.SUCCESS; }定時調(diào)度策略
調(diào)度策略執(zhí)行圖
調(diào)度策略源碼
JobScheduleHelper.java->start();路由策略
第一個
固定選擇第一個機(jī)器
ExecutorRouteFirst.java->route();最后一個
固定選擇最后一個機(jī)器
ExecutorRouteLast.java->route();輪詢
隨機(jī)選擇在線的機(jī)器
ExecutorRouteRound.java->route();private static int count(int jobId) {// cache clearif (System.currentTimeMillis() > CACHE_VALID_TIME) {routeCountEachJob.clear();CACHE_VALID_TIME = System.currentTimeMillis() + 1000*60*60*24;}// count++Integer count = routeCountEachJob.get(jobId);count = (count==null || count>1000000)?(new Random().nextInt(100)):++count; // 初始化時主動Random一次,緩解首次壓力routeCountEachJob.put(jobId, count);return count; }隨機(jī)
隨機(jī)獲取地址列表中的一個
ExecutorRouteRandom.java->route();一致性HASH
一個job通過hash算法固定使用一臺機(jī)器,且所有任務(wù)均勻散列在不同機(jī)器
ExecutorRouteConsistentHash.java->route();public String hashJob(int jobId, List<String> addressList) {// ------A1------A2-------A3------// -----------J1------------------TreeMap<Long, String> addressRing = new TreeMap<Long, String>();for (String address: addressList) {for (int i = 0; i < VIRTUAL_NODE_NUM; i++) {long addressHash = hash("SHARD-" + address + "-NODE-" + i);addressRing.put(addressHash, address);}}long jobHash = hash(String.valueOf(jobId));// 取出鍵值 >= jobHashSortedMap<Long, String> lastRing = addressRing.tailMap(jobHash);if (!lastRing.isEmpty()) {return lastRing.get(lastRing.firstKey());}return addressRing.firstEntry().getValue(); }最不經(jīng)常使用
使用頻率最低的機(jī)器優(yōu)先被選舉
把地址列表加入到內(nèi)存中,等下次執(zhí)行時剔除無效的地址,判斷地址列表中執(zhí)行次數(shù)最少的地址取出
頻率、次數(shù)
最近最久未使用
最久未使用的機(jī)器優(yōu)先被選舉
用鏈表的方式存儲地址,第一個地址使用后下次該任務(wù)過來使用第二個地址,依次類推(PS:有點(diǎn)類似輪詢策略)
與輪詢策略的區(qū)別:
次數(shù)
故障轉(zhuǎn)移
按照順序依次進(jìn)行心跳檢測,第一個心跳檢測成功的機(jī)器選定為目標(biāo)執(zhí)行器并發(fā)起調(diào)度
ExecutorRouteFailover.java->route();忙碌轉(zhuǎn)移
按照順序依次進(jìn)行空閑檢測,第一個空閑檢測成功的機(jī)器選定為目標(biāo)執(zhí)行器并發(fā)起調(diào)度
ExecutorRouteBusyover.java->route();分片廣播
廣播觸發(fā)對應(yīng)集群中所有機(jī)器執(zhí)行一次任務(wù),同時傳遞分片參數(shù);可根據(jù)分片參數(shù)開發(fā)分片任務(wù)
阻塞處理策略
為了解決執(zhí)行線程因并發(fā)問題、執(zhí)行效率慢、任務(wù)多等原因而做的一種線程處理機(jī)制,主要包括 串行、丟棄后續(xù)調(diào)度、覆蓋之前調(diào)度,一般常用策略是串行機(jī)制
ExecutorBlockStrategyEnum.javaSERIAL_EXECUTION("Serial execution"), // 串行 DISCARD_LATER("Discard Later"), // 丟棄后續(xù)調(diào)度 COVER_EARLY("Cover Early"); // 覆蓋之前調(diào)度ExecutorBizImpl.java->run();// executor block strategy if (jobThread != null) {ExecutorBlockStrategyEnum blockStrategy = ExecutorBlockStrategyEnum.match(triggerParam.getExecutorBlockStrategy(), null);if (ExecutorBlockStrategyEnum.DISCARD_LATER == blockStrategy) {// discard when runningif (jobThread.isRunningOrHasQueue()) {return new ReturnT<String>(ReturnT.FAIL_CODE, "block strategy effect:"+ExecutorBlockStrategyEnum.DISCARD_LATER.getTitle());}} else if (ExecutorBlockStrategyEnum.COVER_EARLY == blockStrategy) {// kill running jobThreadif (jobThread.isRunningOrHasQueue()) {removeOldReason = "block strategy effect:" + ExecutorBlockStrategyEnum.COVER_EARLY.getTitle();jobThread = null;}} else {// just queue trigger} }單機(jī)串行
對當(dāng)前線程不做任何處理,并在當(dāng)前線程的隊(duì)列里增加一個執(zhí)行任務(wù)
丟棄后續(xù)調(diào)度
如果當(dāng)前線程阻塞,后續(xù)任務(wù)不再執(zhí)行,直接返回失敗
覆蓋之前調(diào)度
創(chuàng)建一個移除原因,新建一個線程去執(zhí)行后續(xù)任務(wù)
運(yùn)行模式
ExecutorBizImpl.java->run();BEAN
java里的bean對象
GLUE(Java)
利用java的反射機(jī)制,通過代碼字符串生成實(shí)體類
IJobHandler originJobHandler = GlueFactory.getInstance().loadNewInstance(triggerParam.getGlueSource());GroovyClassLoaderGLUE(Shell Python PHP Nodejs PowerShell)
按照文件命名規(guī)則創(chuàng)建一個執(zhí)行腳本文件和一個日志輸出文件,通過腳本執(zhí)行器執(zhí)行
失敗重試次數(shù)
任務(wù)失敗后記錄到 xxl_job_log 中,由失敗監(jiān)控線程查詢處理失敗的任務(wù)且失敗次數(shù)大于0,繼續(xù)執(zhí)行
任務(wù)超時時間
把超時時間給 triggerParam 觸發(fā)參數(shù),在調(diào)用執(zhí)行器的任務(wù)時超時時間,有點(diǎn)類似HttpClient的超時時間
執(zhí)行器(Exector)
注冊自己的機(jī)器地址
注冊項(xiàng)目中的 JobHandler
提供被調(diào)度中心調(diào)用的接口
public interface ExecutorBiz {/*** 供調(diào)度中心檢測機(jī)器是否存活** beat* @return*/public ReturnT<String> beat();/*** 供調(diào)度中心檢測機(jī)器是否空閑** @param jobId* @return*/public ReturnT<String> idleBeat(int jobId);/*** kill* @param jobId* @return*/public ReturnT<String> kill(int jobId);/*** log* @param logDateTim* @param logId* @param fromLineNum* @return*/public ReturnT<LogResult> log(long logDateTim, long logId, int fromLineNum);/*** 執(zhí)行觸發(fā)器* * @param triggerParam* @return*/public ReturnT<String> run(TriggerParam triggerParam);}總結(jié)
學(xué)到了什么
轉(zhuǎn)載于:https://www.cnblogs.com/guoyinli/p/11555035.html
總結(jié)
以上是生活随笔為你收集整理的xxl-job源码分析的全部內(nèi)容,希望文章能夠幫你解決所遇到的問題。
- 上一篇: Python实现图形学DDA算法
- 下一篇: oracle between and m