refactor:调度组件守护线程代码重构,提升稳定性以及可维护性;

3.4.1-release
xuxueli 3 months ago
parent ab77a99f9b
commit 4d94321beb

@ -2846,8 +2846,9 @@ alter table xxl_job_log
- 3、【优化】任务参数长度调整,最长支持2048字符; - 3、【优化】任务参数长度调整,最长支持2048字符;
- 4、【升级】调度中心UI交互优化,任务及日志管理支持下拉框模糊搜索,提升交互体验; - 4、【升级】调度中心UI交互优化,任务及日志管理支持下拉框模糊搜索,提升交互体验;
- 5、【修复】XxlJobFileAppender自定义地址callbackLogPath设置无效问题修复;合并ISSUS-3963; - 5、【修复】XxlJobFileAppender自定义地址callbackLogPath设置无效问题修复;合并ISSUS-3963;
- 6、【TODO】调度中心OpenAPI完善,提供任务管理能力;封装Agent Skill并推送ClawHub; - 6、【优化】调度组件守护线程代码重构,提升稳定性以及可维护性;
- 7、【TODO】AccessToken升级:执行器维度隔离,支持线上化配置;升级双端OpenApi,适配AccessToken升级; - 7、【TODO】调度中心OpenAPI完善,提供任务管理能力;封装Agent Skill并推送ClawHub;
- 8、【TODO】AccessToken升级:执行器维度隔离,支持线上化配置;升级双端OpenApi,适配AccessToken升级;
### TODO LIST ### TODO LIST

@ -3,8 +3,9 @@ package com.xxl.job.admin.business.scheduler.thread;
import com.xxl.job.admin.business.model.XxlJobLog; import com.xxl.job.admin.business.model.XxlJobLog;
import com.xxl.job.admin.business.scheduler.config.XxlJobAdminBootstrap; import com.xxl.job.admin.business.scheduler.config.XxlJobAdminBootstrap;
import com.xxl.job.admin.framework.util.I18nUtil; import com.xxl.job.admin.framework.util.I18nUtil;
import com.xxl.job.core.openapi.model.CallbackRequest;
import com.xxl.job.core.context.XxlJobContext; import com.xxl.job.core.context.XxlJobContext;
import com.xxl.job.core.openapi.model.CallbackRequest;
import com.xxl.tool.concurrent.CyclicThread;
import com.xxl.tool.core.DateTool; import com.xxl.tool.core.DateTool;
import com.xxl.tool.response.Response; import com.xxl.tool.response.Response;
import org.slf4j.Logger; import org.slf4j.Logger;
@ -25,15 +26,14 @@ public class JobCompleteHelper {
// ---------------------- monitor ---------------------- // ---------------------- monitor ----------------------
private ThreadPoolExecutor callbackThreadPool = null; private ThreadPoolExecutor callbackThreadPool = null;
private Thread monitorThread; private CyclicThread jobMonitorThread;
private volatile boolean toStop = false;
/** /**
* start * start
*/ */
public void start(){ public void start(){
// for callback // 1、callbackThreadPool
callbackThreadPool = new ThreadPoolExecutor( callbackThreadPool = new ThreadPoolExecutor(
2, 2,
20, 20,
@ -55,88 +55,55 @@ public class JobCompleteHelper {
}); });
// for monitor // 2、jobMonitorThread
monitorThread = new Thread(new Runnable() { jobMonitorThread = new CyclicThread("JobCompleteHelper#jobMonitorThread", true, new Runnable() {
@Override @Override
public void run() { public void run() {
// 任务结果丢失处理:调度记录停留在 "运行中" 状态超过10min,且对应执行器心跳注册失败不在线,则将本地调度主动标记失败;
Date losedTime = DateTool.addMinutes(new Date(), -10);
List<Long> losedJobIds = XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().findLostJobIds(losedTime);
// wait for JobTriggerPoolHelper-init if (losedJobIds!=null && losedJobIds.size()>0) {
try { for (Long logId: losedJobIds) {
TimeUnit.MILLISECONDS.sleep(50);
} catch (Throwable e) {
if (!toStop) {
logger.error(e.getMessage(), e);
}
}
// monitor XxlJobLog jobLog = new XxlJobLog();
while (!toStop) { jobLog.setId(logId);
try {
// 任务结果丢失处理:调度记录停留在 "运行中" 状态超过10min,且对应执行器心跳注册失败不在线,则将本地调度主动标记失败;
Date losedTime = DateTool.addMinutes(new Date(), -10);
List<Long> losedJobIds = XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().findLostJobIds(losedTime);
if (losedJobIds!=null && losedJobIds.size()>0) { jobLog.setHandleTime(new Date());
for (Long logId: losedJobIds) { jobLog.setHandleCode(XxlJobContext.HANDLE_CODE_FAIL);
jobLog.setHandleMsg( I18nUtil.getString("joblog_lost_fail") );
XxlJobLog jobLog = new XxlJobLog(); XxlJobAdminBootstrap.getInstance().getJobCompleter().complete(jobLog);
jobLog.setId(logId);
jobLog.setHandleTime(new Date());
jobLog.setHandleCode(XxlJobContext.HANDLE_CODE_FAIL);
jobLog.setHandleMsg( I18nUtil.getString("joblog_lost_fail") );
XxlJobAdminBootstrap.getInstance().getJobCompleter().complete(jobLog);
}
}
} catch (Throwable e) {
if (!toStop) {
logger.error(">>>>>>>>>>> xxl-job, job fail monitor thread error:{}", e);
}
} }
try { }
TimeUnit.SECONDS.sleep(60);
} catch (Throwable e) {
if (!toStop) {
logger.error(e.getMessage(), e);
}
}
}
logger.info(">>>>>>>>>>> xxl-job, JobLosedMonitorHelper stop");
} }
}); }, 60 * 1000L, true);
monitorThread.setDaemon(true); jobMonitorThread.start();
monitorThread.setName("xxl-job, admin JobLosedMonitorHelper");
monitorThread.start();
} }
/** /**
* stop * stop
*/ */
public void stop(){ public void stop(){
toStop = true;
// stop registryOrRemoveThreadPool // 1、callbackThreadPool
callbackThreadPool.shutdownNow(); callbackThreadPool.shutdownNow();
// stop monitorThread (interrupt and wait) // 2、jobMonitorThread
monitorThread.interrupt(); jobMonitorThread.stop();
try {
monitorThread.join();
} catch (Throwable e) {
logger.error(e.getMessage(), e);
}
} }
// ---------------------- helper ---------------------- // ---------------------- helper ----------------------
/**
* callback
*
* @param callbackParamList callback param
* @return callback result
*/
public Response<String> callback(List<CallbackRequest> callbackParamList) { public Response<String> callback(List<CallbackRequest> callbackParamList) {
callbackThreadPool.execute(new Runnable() { callbackThreadPool.execute(new Runnable() {

@ -5,11 +5,11 @@ import com.xxl.job.admin.business.model.XxlJobLog;
import com.xxl.job.admin.business.scheduler.config.XxlJobAdminBootstrap; import com.xxl.job.admin.business.scheduler.config.XxlJobAdminBootstrap;
import com.xxl.job.admin.business.scheduler.trigger.TriggerTypeEnum; import com.xxl.job.admin.business.scheduler.trigger.TriggerTypeEnum;
import com.xxl.job.admin.framework.util.I18nUtil; import com.xxl.job.admin.framework.util.I18nUtil;
import com.xxl.tool.concurrent.CyclicThread;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import java.util.List; import java.util.List;
import java.util.concurrent.TimeUnit;
/** /**
* job fail-monitor helper * job fail-monitor helper
@ -22,77 +22,53 @@ public class JobFailAlarmMonitorHelper {
// ---------------------- monitor ---------------------- // ---------------------- monitor ----------------------
private Thread monitorThread; /**
private volatile boolean toStop = false; * monitor thread
*/
private CyclicThread monitorThread;
/** /**
* start * start
*/ */
public void start(){ public void start(){
monitorThread = new Thread(new Runnable() {
monitorThread = new CyclicThread("JobFailAlarmMonitorHelper#monitorThread", true, new Runnable() {
@Override @Override
public void run() { public void run() {
List<Long> failLogIds = XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().findFailJobLogIds(1000);
// monitor if (failLogIds!=null && !failLogIds.isEmpty()) {
while (!toStop) { for (long failLogId: failLogIds) {
try {
// lock log
List<Long> failLogIds = XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().findFailJobLogIds(1000); int lockRet = XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().updateAlarmStatus(failLogId, 0, -1);
if (failLogIds!=null && !failLogIds.isEmpty()) { if (lockRet < 1) {
for (long failLogId: failLogIds) { continue;
// lock log
int lockRet = XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().updateAlarmStatus(failLogId, 0, -1);
if (lockRet < 1) {
continue;
}
XxlJobLog log = XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().load(failLogId);
XxlJobInfo info = XxlJobAdminBootstrap.getInstance().getXxlJobInfoMapper().loadById(log.getJobId());
// 1、fail retry monitor
if (log.getExecutorFailRetryCount() > 0) {
XxlJobAdminBootstrap.getInstance().getJobTriggerPoolHelper().trigger(log.getJobId(), TriggerTypeEnum.RETRY, (log.getExecutorFailRetryCount()-1), log.getExecutorShardingParam(), log.getExecutorParam(), null);
String retryMsg = "<br><br><span style=\"color:#00c0ef;\" > >>>>>>>>>>>"+ I18nUtil.getString("jobconf_trigger_type_retry") +"<<<<<<<<<<< </span><br>";
log.setTriggerMsg(log.getTriggerMsg() + retryMsg);
XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().updateTriggerInfo(log);
}
// 2、fail alarm monitor
int newAlarmStatus = 0; // 告警状态:0-默认、-1=锁定状态、1-无需告警、2-告警成功、3-告警失败
if (info != null) {
boolean alarmResult = XxlJobAdminBootstrap.getInstance().getJobAlarmer().alarm(info, log);
newAlarmStatus = alarmResult?2:3;
} else {
newAlarmStatus = 1;
}
XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().updateAlarmStatus(failLogId, -1, newAlarmStatus);
}
} }
XxlJobLog log = XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().load(failLogId);
} catch (Throwable e) { XxlJobInfo info = XxlJobAdminBootstrap.getInstance().getXxlJobInfoMapper().loadById(log.getJobId());
if (!toStop) {
logger.error(">>>>>>>>>>> xxl-job, job fail monitor thread error:{}", e.getMessage(), e); // 1、fail retry monitor
if (log.getExecutorFailRetryCount() > 0) {
XxlJobAdminBootstrap.getInstance().getJobTriggerPoolHelper().trigger(log.getJobId(), TriggerTypeEnum.RETRY, (log.getExecutorFailRetryCount()-1), log.getExecutorShardingParam(), log.getExecutorParam(), null);
String retryMsg = "<br><br><span style=\"color:#00c0ef;\" > >>>>>>>>>>>"+ I18nUtil.getString("jobconf_trigger_type_retry") +"<<<<<<<<<<< </span><br>";
log.setTriggerMsg(log.getTriggerMsg() + retryMsg);
XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().updateTriggerInfo(log);
} }
}
try {
TimeUnit.SECONDS.sleep(10);
} catch (Throwable e) {
if (!toStop) {
logger.error(e.getMessage(), e);
}
}
} // 2、fail alarm monitor
int newAlarmStatus = 0; // 告警状态:0-默认、-1=锁定状态、1-无需告警、2-告警成功、3-告警失败
logger.info(">>>>>>>>>>> xxl-job, job fail monitor thread stop"); if (info != null) {
boolean alarmResult = XxlJobAdminBootstrap.getInstance().getJobAlarmer().alarm(info, log);
newAlarmStatus = alarmResult?2:3;
} else {
newAlarmStatus = 1;
}
XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().updateAlarmStatus(failLogId, -1, newAlarmStatus);
}
}
} }
}); }, 10 * 1000L, true);
monitorThread.setDaemon(true);
monitorThread.setName("xxl-job, admin JobFailMonitorHelper");
monitorThread.start(); monitorThread.start();
} }
@ -100,14 +76,7 @@ public class JobFailAlarmMonitorHelper {
* stop * stop
*/ */
public void stop(){ public void stop(){
toStop = true; monitorThread.stop();
// interrupt and wait
monitorThread.interrupt();
try {
monitorThread.join();
} catch (Throwable e) {
logger.error(e.getMessage(), e);
}
} }
} }

@ -1,7 +1,8 @@
package com.xxl.job.admin.business.scheduler.thread; package com.xxl.job.admin.business.scheduler.thread;
import com.xxl.job.admin.business.scheduler.config.XxlJobAdminBootstrap;
import com.xxl.job.admin.business.model.XxlJobLogReport; import com.xxl.job.admin.business.model.XxlJobLogReport;
import com.xxl.job.admin.business.scheduler.config.XxlJobAdminBootstrap;
import com.xxl.tool.concurrent.CyclicThread;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
@ -9,7 +10,7 @@ import java.util.Calendar;
import java.util.Date; import java.util.Date;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong;
/** /**
* job log report helper * job log report helper
@ -19,144 +20,109 @@ import java.util.concurrent.TimeUnit;
public class JobLogReportHelper { public class JobLogReportHelper {
private static final Logger logger = LoggerFactory.getLogger(JobLogReportHelper.class); private static final Logger logger = LoggerFactory.getLogger(JobLogReportHelper.class);
private CyclicThread logReportThread;
private Thread logReportThread; private AtomicLong lastCleanLogTime;
private volatile boolean toStop = false;
/** /**
* start * start
*/ */
public void start(){ public void start(){
logReportThread = new Thread(new Runnable() {
/**
* last clean log time ( Thread-safe concurrent reading and writing )
*/
lastCleanLogTime = new AtomicLong(0);
// log report thread
logReportThread = new CyclicThread("JobLogReportHelper#logReportThread", true, new Runnable() {
@Override @Override
public void run() { public void run() {
// last clean log time // 1、log-report refresh: refresh log report in 3 days
long lastCleanLogTime = 0; for (int i = 0; i < 3; i++) {
// today
while (!toStop) { Calendar itemDay = Calendar.getInstance();
itemDay.add(Calendar.DAY_OF_MONTH, -i);
// 1、log-report refresh: refresh log report in 3 days itemDay.set(Calendar.HOUR_OF_DAY, 0);
try { itemDay.set(Calendar.MINUTE, 0);
itemDay.set(Calendar.SECOND, 0);
for (int i = 0; i < 3; i++) { itemDay.set(Calendar.MILLISECOND, 0);
// today Date todayFrom = itemDay.getTime();
Calendar itemDay = Calendar.getInstance();
itemDay.add(Calendar.DAY_OF_MONTH, -i); itemDay.set(Calendar.HOUR_OF_DAY, 23);
itemDay.set(Calendar.HOUR_OF_DAY, 0); itemDay.set(Calendar.MINUTE, 59);
itemDay.set(Calendar.MINUTE, 0); itemDay.set(Calendar.SECOND, 59);
itemDay.set(Calendar.SECOND, 0); itemDay.set(Calendar.MILLISECOND, 999);
itemDay.set(Calendar.MILLISECOND, 0);
Date todayTo = itemDay.getTime();
Date todayFrom = itemDay.getTime();
// refresh log-report every minute
itemDay.set(Calendar.HOUR_OF_DAY, 23); XxlJobLogReport xxlJobLogReport = new XxlJobLogReport();
itemDay.set(Calendar.MINUTE, 59); xxlJobLogReport.setTriggerDay(todayFrom);
itemDay.set(Calendar.SECOND, 59); xxlJobLogReport.setRunningCount(0);
itemDay.set(Calendar.MILLISECOND, 999); xxlJobLogReport.setSucCount(0);
xxlJobLogReport.setFailCount(0);
Date todayTo = itemDay.getTime(); xxlJobLogReport.setUpdateTime(new Date());
// refresh log-report every minute // fill count-data
XxlJobLogReport xxlJobLogReport = new XxlJobLogReport(); Map<String, Object> triggerCountMap = XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().findLogReport(todayFrom, todayTo);
xxlJobLogReport.setTriggerDay(todayFrom); if (triggerCountMap!=null && !triggerCountMap.isEmpty()) {
xxlJobLogReport.setRunningCount(0); int triggerDayCount = triggerCountMap.containsKey("triggerDayCount")?Integer.parseInt(String.valueOf(triggerCountMap.get("triggerDayCount"))):0;
xxlJobLogReport.setSucCount(0); int triggerDayCountRunning = triggerCountMap.containsKey("triggerDayCountRunning")?Integer.parseInt(String.valueOf(triggerCountMap.get("triggerDayCountRunning"))):0;
xxlJobLogReport.setFailCount(0); int triggerDayCountSuc = triggerCountMap.containsKey("triggerDayCountSuc")?Integer.parseInt(String.valueOf(triggerCountMap.get("triggerDayCountSuc"))):0;
xxlJobLogReport.setUpdateTime(new Date()); int triggerDayCountFail = triggerDayCount - triggerDayCountRunning - triggerDayCountSuc;
// fill count-data xxlJobLogReport.setRunningCount(triggerDayCountRunning);
Map<String, Object> triggerCountMap = XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().findLogReport(todayFrom, todayTo); xxlJobLogReport.setSucCount(triggerDayCountSuc);
if (triggerCountMap!=null && !triggerCountMap.isEmpty()) { xxlJobLogReport.setFailCount(triggerDayCountFail);
int triggerDayCount = triggerCountMap.containsKey("triggerDayCount")?Integer.parseInt(String.valueOf(triggerCountMap.get("triggerDayCount"))):0;
int triggerDayCountRunning = triggerCountMap.containsKey("triggerDayCountRunning")?Integer.parseInt(String.valueOf(triggerCountMap.get("triggerDayCountRunning"))):0;
int triggerDayCountSuc = triggerCountMap.containsKey("triggerDayCountSuc")?Integer.parseInt(String.valueOf(triggerCountMap.get("triggerDayCountSuc"))):0;
int triggerDayCountFail = triggerDayCount - triggerDayCountRunning - triggerDayCountSuc;
xxlJobLogReport.setRunningCount(triggerDayCountRunning);
xxlJobLogReport.setSucCount(triggerDayCountSuc);
xxlJobLogReport.setFailCount(triggerDayCountFail);
}
// do refresh:
XxlJobAdminBootstrap.getInstance().getXxlJobLogReportMapper().saveOrUpdate(xxlJobLogReport); // 0-fail; 1-save suc; 2-update suc;
/*if (ret < 1) {
XxlJobAdminBootstrap.getInstance().getXxlJobLogReportMapper().save(xxlJobLogReport);
}*/
}
} catch (Throwable e) {
if (!toStop) {
logger.error(">>>>>>>>>>> xxl-job, JobLogReportHelper(log-report refresh) error:{}", e.getMessage(), e);
}
} }
// 2、log-clean: switch open & once each day // do refresh:
try { XxlJobAdminBootstrap.getInstance().getXxlJobLogReportMapper().saveOrUpdate(xxlJobLogReport); // 0-fail; 1-save suc; 2-update suc;
if (XxlJobAdminBootstrap.getInstance().getLogretentiondays()>0 /*if (ret < 1) {
&& System.currentTimeMillis() - lastCleanLogTime > 24*60*60*1000) { XxlJobAdminBootstrap.getInstance().getXxlJobLogReportMapper().save(xxlJobLogReport);
}*/
// expire-time }
Calendar expiredDay = Calendar.getInstance();
expiredDay.add(Calendar.DAY_OF_MONTH, -1 * XxlJobAdminBootstrap.getInstance().getLogretentiondays());
expiredDay.set(Calendar.HOUR_OF_DAY, 0);
expiredDay.set(Calendar.MINUTE, 0);
expiredDay.set(Calendar.SECOND, 0);
expiredDay.set(Calendar.MILLISECOND, 0);
Date clearBeforeTime = expiredDay.getTime();
// clean expired log
List<Long> logIds = null;
do {
logIds = XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().findClearLogIds(0, 0, clearBeforeTime, 0, 1000);
if (logIds!=null && !logIds.isEmpty()) {
XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().clearLog(logIds);
}
} while (logIds!=null && !logIds.isEmpty());
// update clean time
lastCleanLogTime = System.currentTimeMillis();
}
} catch (Throwable e) {
if (!toStop) {
logger.error(">>>>>>>>>>> xxl-job, JobLogReportHelper(log-clean) error:{}", e.getMessage(), e);
}
}
try { // 2、log-clean: switch open & once each day
TimeUnit.MINUTES.sleep(1); if (XxlJobAdminBootstrap.getInstance().getLogretentiondays()>0
} catch (Throwable e) { && System.currentTimeMillis() - lastCleanLogTime.longValue() > 24*60*60*1000) {
if (!toStop) {
logger.error(e.getMessage(), e); // expire-time
Calendar expiredDay = Calendar.getInstance();
expiredDay.add(Calendar.DAY_OF_MONTH, -1 * XxlJobAdminBootstrap.getInstance().getLogretentiondays());
expiredDay.set(Calendar.HOUR_OF_DAY, 0);
expiredDay.set(Calendar.MINUTE, 0);
expiredDay.set(Calendar.SECOND, 0);
expiredDay.set(Calendar.MILLISECOND, 0);
Date clearBeforeTime = expiredDay.getTime();
// clean expired log
List<Long> logIds = null;
do {
logIds = XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().findClearLogIds(0, 0, clearBeforeTime, 0, 1000);
if (logIds!=null && !logIds.isEmpty()) {
XxlJobAdminBootstrap.getInstance().getXxlJobLogMapper().clearLog(logIds);
} }
} } while (logIds!=null && !logIds.isEmpty());
// update clean time
lastCleanLogTime.set(System.currentTimeMillis());
} }
logger.info(">>>>>>>>>>> xxl-job, job log report thread stop");
} }
}); }, 60 * 1000L, true);
logReportThread.setDaemon(true);
logReportThread.setName("xxl-job, admin JobLogReportHelper");
logReportThread.start(); logReportThread.start();
} }
/** /**
* stop * stop
*/ */
public void stop(){ public void stop(){
toStop = true; logReportThread.stop();
// interrupt and wait
logReportThread.interrupt();
try {
logReportThread.join();
} catch (Throwable e) {
logger.error(e.getMessage(), e);
}
} }
} }

@ -3,9 +3,10 @@ package com.xxl.job.admin.business.scheduler.thread;
import com.xxl.job.admin.business.model.XxlJobGroup; import com.xxl.job.admin.business.model.XxlJobGroup;
import com.xxl.job.admin.business.model.XxlJobRegistry; import com.xxl.job.admin.business.model.XxlJobRegistry;
import com.xxl.job.admin.business.scheduler.config.XxlJobAdminBootstrap; import com.xxl.job.admin.business.scheduler.config.XxlJobAdminBootstrap;
import com.xxl.job.core.constant.Const;
import com.xxl.job.core.constant.RegistType; import com.xxl.job.core.constant.RegistType;
import com.xxl.job.core.openapi.model.RegistryRequest; import com.xxl.job.core.openapi.model.RegistryRequest;
import com.xxl.job.core.constant.Const; import com.xxl.tool.concurrent.CyclicThread;
import com.xxl.tool.core.StringTool; import com.xxl.tool.core.StringTool;
import com.xxl.tool.response.Response; import com.xxl.tool.response.Response;
import org.slf4j.Logger; import org.slf4j.Logger;
@ -20,20 +21,25 @@ import java.util.concurrent.*;
* @author xuxueli 2016-10-02 19:10:24 * @author xuxueli 2016-10-02 19:10:24
*/ */
public class JobRegistryHelper { public class JobRegistryHelper {
private static Logger logger = LoggerFactory.getLogger(JobRegistryHelper.class); private static final Logger logger = LoggerFactory.getLogger(JobRegistryHelper.class);
/**
* registry or remove thread pool
*/
private ThreadPoolExecutor registryOrRemoveThreadPool = null; private ThreadPoolExecutor registryOrRemoveThreadPool = null;
private Thread registryMonitorThread;
private volatile boolean toStop = false;
/**
* registry monitor thread
*/
private CyclicThread registryMonitorThread;
/** /**
* start * start
*/ */
public void start(){ public void start(){
// for registry or remove // 1、for registry or remove
registryOrRemoveThreadPool = new ThreadPoolExecutor( registryOrRemoveThreadPool = new ThreadPoolExecutor(
2, 2,
10, 10,
@ -54,79 +60,61 @@ public class JobRegistryHelper {
} }
}); });
// for monitor // 2、for registry monitor
registryMonitorThread = new Thread(new Runnable() { registryMonitorThread = new CyclicThread("JobRegistryHelper#registryMonitorThread", true, new Runnable() {
@Override @Override
public void run() { public void run() {
while (!toStop) { // auto registry group
try { List<XxlJobGroup> groupList = XxlJobAdminBootstrap.getInstance().getXxlJobGroupMapper().findByAddressType(0);
// auto registry group if (groupList!=null && !groupList.isEmpty()) {
List<XxlJobGroup> groupList = XxlJobAdminBootstrap.getInstance().getXxlJobGroupMapper().findByAddressType(0);
if (groupList!=null && !groupList.isEmpty()) { // remove dead address (admin/executor)
List<Integer> ids = XxlJobAdminBootstrap.getInstance().getXxlJobRegistryMapper().findDead(Const.DEAD_TIMEOUT, new Date());
// remove dead address (admin/executor) if (ids!=null && !ids.isEmpty()) {
List<Integer> ids = XxlJobAdminBootstrap.getInstance().getXxlJobRegistryMapper().findDead(Const.DEAD_TIMEOUT, new Date()); XxlJobAdminBootstrap.getInstance().getXxlJobRegistryMapper().removeDead(ids);
if (ids!=null && !ids.isEmpty()) { }
XxlJobAdminBootstrap.getInstance().getXxlJobRegistryMapper().removeDead(ids);
}
// fresh online address (admin/executor) // fresh online address (admin/executor)
HashMap<String, List<String>> appAddressMap = new HashMap<String, List<String>>(); HashMap<String, List<String>> appAddressMap = new HashMap<String, List<String>>();
List<XxlJobRegistry> list = XxlJobAdminBootstrap.getInstance().getXxlJobRegistryMapper().findAll(Const.DEAD_TIMEOUT, new Date()); List<XxlJobRegistry> list = XxlJobAdminBootstrap.getInstance().getXxlJobRegistryMapper().findAll(Const.DEAD_TIMEOUT, new Date());
if (list != null) { if (list != null) {
for (XxlJobRegistry item: list) { for (XxlJobRegistry item: list) {
if (RegistType.EXECUTOR.name().equals(item.getRegistryGroup())) { if (RegistType.EXECUTOR.name().equals(item.getRegistryGroup())) {
String appname = item.getRegistryKey(); String appname = item.getRegistryKey();
List<String> registryList = appAddressMap.get(appname); List<String> registryList = appAddressMap.get(appname);
if (registryList == null) { if (registryList == null) {
registryList = new ArrayList<String>(); registryList = new ArrayList<String>();
}
if (!registryList.contains(item.getRegistryValue())) {
registryList.add(item.getRegistryValue());
}
appAddressMap.put(appname, registryList);
}
} }
}
// fresh group address if (!registryList.contains(item.getRegistryValue())) {
for (XxlJobGroup group: groupList) { registryList.add(item.getRegistryValue());
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); appAddressMap.put(appname, registryList);
group.setUpdateTime(new Date());
XxlJobAdminBootstrap.getInstance().getXxlJobGroupMapper().update(group);
} }
} }
} catch (Throwable e) {
if (!toStop) {
logger.error(">>>>>>>>>>> xxl-job, job registry monitor thread error:{}", e);
}
} }
try {
TimeUnit.SECONDS.sleep(Const.BEAT_TIMEOUT); // fresh group address
} catch (Throwable e) { for (XxlJobGroup group: groupList) {
if (!toStop) { List<String> registryList = appAddressMap.get(group.getAppname());
logger.error(">>>>>>>>>>> xxl-job, job registry monitor thread error:{}", e); 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());
XxlJobAdminBootstrap.getInstance().getXxlJobGroupMapper().update(group);
} }
} }
logger.info(">>>>>>>>>>> xxl-job, job registry monitor thread stop");
} }
}); }, Const.BEAT_TIMEOUT * 1000L, true);
registryMonitorThread.setDaemon(true);
registryMonitorThread.setName("xxl-job, admin JobRegistryMonitorHelper-registryMonitorThread");
registryMonitorThread.start(); registryMonitorThread.start();
} }
@ -135,18 +123,12 @@ public class JobRegistryHelper {
* stop * stop
*/ */
public void stop(){ public void stop(){
toStop = true;
// stop registryOrRemoveThreadPool // 1、registryOrRemoveThreadPool
registryOrRemoveThreadPool.shutdownNow(); registryOrRemoveThreadPool.shutdownNow();
// stop monitor (interrupt and wait) // 2、registryMonitorThread
registryMonitorThread.interrupt(); registryMonitorThread.stop();
try {
registryMonitorThread.join();
} catch (Throwable e) {
logger.error(e.getMessage(), e);
}
} }

@ -44,7 +44,11 @@ public class JobScheduleHelper {
*/ */
public void start(){ public void start(){
// schedule thread // init thread flag
scheduleThreadToStop = false;
ringThreadToStop = false;
// 1、schedule thread
scheduleThread = new Thread(new Runnable() { scheduleThread = new Thread(new Runnable() {
@Override @Override
public void run() { public void run() {
@ -191,8 +195,7 @@ public class JobScheduleHelper {
scheduleThread.setName("xxl-job, admin JobScheduleHelper#scheduleThread"); scheduleThread.setName("xxl-job, admin JobScheduleHelper#scheduleThread");
scheduleThread.start(); scheduleThread.start();
// 2、ring thread
// ring thread
ringThread = new Thread(new Runnable() { ringThread = new Thread(new Runnable() {
@Override @Override
public void run() { public void run() {
@ -348,7 +351,7 @@ public class JobScheduleHelper {
} }
} }
// stop ring (wait job-in-memory stop) // 2、stop ring (wait job-in-memory stop)
ringThreadToStop = true; ringThreadToStop = true;
try { try {
TimeUnit.SECONDS.sleep(1); TimeUnit.SECONDS.sleep(1);

Loading…
Cancel
Save