|
|
|
@ -33,12 +33,8 @@ public class ExecutorBizImpl implements ExecutorBiz {
|
|
|
|
|
public ReturnT<String> kill(int jobId) {
|
|
|
|
|
// kill handlerThread, and create new one
|
|
|
|
|
JobThread jobThread = XxlJobExecutor.loadJobThread(jobId);
|
|
|
|
|
|
|
|
|
|
if (jobThread != null) {
|
|
|
|
|
IJobHandler handler = jobThread.getHandler();
|
|
|
|
|
jobThread.toStop("人工手动终止");
|
|
|
|
|
jobThread.interrupt();
|
|
|
|
|
XxlJobExecutor.removeJobThread(jobId);
|
|
|
|
|
XxlJobExecutor.removeJobThread(jobId, "人工手动终止");
|
|
|
|
|
return ReturnT.SUCCESS;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
@ -56,81 +52,103 @@ public class ExecutorBizImpl implements ExecutorBiz {
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
public ReturnT<String> run(TriggerParam triggerParam) {
|
|
|
|
|
// load old thread
|
|
|
|
|
// load old:jobHandler + jobThread
|
|
|
|
|
JobThread jobThread = XxlJobExecutor.loadJobThread(triggerParam.getJobId());
|
|
|
|
|
IJobHandler jobHandler = jobThread!=null?jobThread.getHandler():null;
|
|
|
|
|
String removeOldReason = null;
|
|
|
|
|
|
|
|
|
|
// valid:jobHandler + jobThread
|
|
|
|
|
if (GlueTypeEnum.BEAN==GlueTypeEnum.match(triggerParam.getGlueType())) {
|
|
|
|
|
|
|
|
|
|
// valid handler
|
|
|
|
|
IJobHandler jobHandler = XxlJobExecutor.loadJobHandler(triggerParam.getExecutorHandler());
|
|
|
|
|
if (jobHandler==null) {
|
|
|
|
|
return new ReturnT<String>(ReturnT.FAIL_CODE, "job handler [" + triggerParam.getExecutorHandler() + "] not found.");
|
|
|
|
|
}
|
|
|
|
|
// valid old jobThread
|
|
|
|
|
if (jobThread != null && jobHandler!=null && jobThread.getHandler() != jobHandler) {
|
|
|
|
|
// change handler, need kill old thread
|
|
|
|
|
removeOldReason = "更新JobHandler或更换任务模式,终止旧任务线程";
|
|
|
|
|
|
|
|
|
|
// valid exists job thread:change handler, need kill old thread
|
|
|
|
|
if (jobThread != null && jobThread.getHandler() != jobHandler) {
|
|
|
|
|
// kill old job thread
|
|
|
|
|
jobThread.toStop("更换任务模式或JobHandler,终止旧任务线程");
|
|
|
|
|
jobThread.interrupt();
|
|
|
|
|
XxlJobExecutor.removeJobThread(triggerParam.getJobId());
|
|
|
|
|
jobThread = null;
|
|
|
|
|
jobHandler = null;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// make thread: new or exists invalid
|
|
|
|
|
if (jobThread == null) {
|
|
|
|
|
jobThread = XxlJobExecutor.registJobThread(triggerParam.getJobId(), jobHandler);
|
|
|
|
|
// valid handler
|
|
|
|
|
if (jobHandler == null) {
|
|
|
|
|
jobHandler = XxlJobExecutor.loadJobHandler(triggerParam.getExecutorHandler());
|
|
|
|
|
if (jobHandler == null) {
|
|
|
|
|
return new ReturnT<String>(ReturnT.FAIL_CODE, "job handler [" + triggerParam.getExecutorHandler() + "] not found.");
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
} else if (GlueTypeEnum.GLUE_GROOVY==GlueTypeEnum.match(triggerParam.getGlueType())) {
|
|
|
|
|
|
|
|
|
|
// valid exists job thread:change handler or gluesource updated, need kill old thread
|
|
|
|
|
// valid old jobThread
|
|
|
|
|
if (jobThread != null &&
|
|
|
|
|
!(jobThread.getHandler() instanceof GlueJobHandler
|
|
|
|
|
&& ((GlueJobHandler) jobThread.getHandler()).getGlueUpdatetime()==triggerParam.getGlueUpdatetime() )) {
|
|
|
|
|
// change glue model or gluesource updated, kill old job thread
|
|
|
|
|
jobThread.toStop("更换任务模式或JobHandler,终止旧任务线程");
|
|
|
|
|
jobThread.interrupt();
|
|
|
|
|
XxlJobExecutor.removeJobThread(triggerParam.getJobId());
|
|
|
|
|
// change handler or gluesource updated, need kill old thread
|
|
|
|
|
removeOldReason = "更新任务逻辑或更换任务模式,终止旧任务线程";
|
|
|
|
|
|
|
|
|
|
jobThread = null;
|
|
|
|
|
jobHandler = null;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// make thread: new or exists invalid
|
|
|
|
|
if (jobThread == null) {
|
|
|
|
|
IJobHandler jobHandler = null;
|
|
|
|
|
// valid handler
|
|
|
|
|
if (jobHandler == null) {
|
|
|
|
|
try {
|
|
|
|
|
jobHandler = GlueFactory.getInstance().loadNewInstance(triggerParam.getGlueSource());
|
|
|
|
|
IJobHandler originJobHandler = GlueFactory.getInstance().loadNewInstance(triggerParam.getGlueSource());
|
|
|
|
|
jobHandler = new GlueJobHandler(originJobHandler, triggerParam.getGlueUpdatetime());
|
|
|
|
|
} catch (Exception e) {
|
|
|
|
|
logger.error("", e);
|
|
|
|
|
return new ReturnT<String>(ReturnT.FAIL_CODE, e.getMessage());
|
|
|
|
|
}
|
|
|
|
|
jobThread = XxlJobExecutor.registJobThread(triggerParam.getJobId(), new GlueJobHandler(jobHandler, triggerParam.getGlueUpdatetime()));
|
|
|
|
|
}
|
|
|
|
|
} else if (GlueTypeEnum.GLUE_SHELL==GlueTypeEnum.match(triggerParam.getGlueType())
|
|
|
|
|
|| GlueTypeEnum.GLUE_PYTHON==GlueTypeEnum.match(triggerParam.getGlueType()) ) {
|
|
|
|
|
|
|
|
|
|
// valid exists job thread:change script or gluesource updated, need kill old thread
|
|
|
|
|
// valid old jobThread
|
|
|
|
|
if (jobThread != null &&
|
|
|
|
|
!(jobThread.getHandler() instanceof ScriptJobHandler
|
|
|
|
|
&& ((ScriptJobHandler) jobThread.getHandler()).getGlueUpdatetime()==triggerParam.getGlueUpdatetime() )) {
|
|
|
|
|
// change glue model or gluesource updated, kill old job thread
|
|
|
|
|
jobThread.toStop("更换任务模式或JobHandler,终止旧任务线程");
|
|
|
|
|
jobThread.interrupt();
|
|
|
|
|
XxlJobExecutor.removeJobThread(triggerParam.getJobId());
|
|
|
|
|
// change script or gluesource updated, need kill old thread
|
|
|
|
|
removeOldReason = "更新任务逻辑或更换任务模式,终止旧任务线程";
|
|
|
|
|
|
|
|
|
|
jobThread = null;
|
|
|
|
|
jobHandler = null;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// make thread: new or exists invalid
|
|
|
|
|
if (jobThread == null) {
|
|
|
|
|
ScriptJobHandler scriptJobHandler = new ScriptJobHandler(triggerParam.getJobId(), triggerParam.getGlueUpdatetime(), triggerParam.getGlueSource(), GlueTypeEnum.match(triggerParam.getGlueType()));
|
|
|
|
|
jobThread = XxlJobExecutor.registJobThread(triggerParam.getJobId(), scriptJobHandler);
|
|
|
|
|
// valid handler
|
|
|
|
|
if (jobHandler == null) {
|
|
|
|
|
jobHandler = new ScriptJobHandler(triggerParam.getJobId(), triggerParam.getGlueUpdatetime(), triggerParam.getGlueSource(), GlueTypeEnum.match(triggerParam.getGlueType()));
|
|
|
|
|
}
|
|
|
|
|
} else {
|
|
|
|
|
return new ReturnT<String>(ReturnT.FAIL_CODE, "glueType[" + triggerParam.getGlueType() + "] is not valid.");
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// executor block strategy
|
|
|
|
|
if (jobThread != null) {
|
|
|
|
|
ExecutorBlockStrategyEnum blockStrategy = ExecutorBlockStrategyEnum.match(triggerParam.getExecutorBlockStrategy(), null);
|
|
|
|
|
if (ExecutorBlockStrategyEnum.DISCARD_LATER == blockStrategy) {
|
|
|
|
|
// discard when running
|
|
|
|
|
if (jobThread.isRunningOrHasQueue()) {
|
|
|
|
|
return new ReturnT<String>(ReturnT.FAIL_CODE, "阻塞处理策略-生效:"+ExecutorBlockStrategyEnum.DISCARD_LATER.getTitle());
|
|
|
|
|
}
|
|
|
|
|
} else if (ExecutorBlockStrategyEnum.COVER_EARLY == blockStrategy) {
|
|
|
|
|
// kill running jobThread
|
|
|
|
|
if (jobThread.isRunningOrHasQueue()) {
|
|
|
|
|
removeOldReason = "阻塞处理策略-生效:" + ExecutorBlockStrategyEnum.COVER_EARLY.getTitle();
|
|
|
|
|
|
|
|
|
|
jobThread = null;
|
|
|
|
|
}
|
|
|
|
|
} else {
|
|
|
|
|
// just queue trigger
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// replace thread (new or exists invalid)
|
|
|
|
|
if (jobThread == null) {
|
|
|
|
|
jobThread = XxlJobExecutor.registJobThread(triggerParam.getJobId(), jobHandler, removeOldReason);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// push data to queue
|
|
|
|
|
ExecutorBlockStrategyEnum blockStrategy = ExecutorBlockStrategyEnum.match(triggerParam.getExecutorBlockStrategy(), null);
|
|
|
|
|
ReturnT<String> pushResult = jobThread.pushTriggerQueue(triggerParam, blockStrategy);
|
|
|
|
|
ReturnT<String> pushResult = jobThread.pushTriggerQueue(triggerParam);
|
|
|
|
|
return pushResult;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|