|
@@ -9,6 +9,8 @@ import com.xxl.job.core.handler.IJobHandler;
|
9
|
9
|
import com.xxl.job.core.handler.impl.GlueJobHandler;
|
10
|
10
|
import com.xxl.job.core.log.XxlJobFileAppender;
|
11
|
11
|
import com.xxl.job.core.thread.JobThread;
|
|
12
|
+import org.slf4j.Logger;
|
|
13
|
+import org.slf4j.LoggerFactory;
|
12
|
14
|
|
13
|
15
|
import java.util.Date;
|
14
|
16
|
|
|
@@ -16,6 +18,7 @@ import java.util.Date;
|
16
|
18
|
* Created by xuxueli on 17/3/1.
|
17
|
19
|
*/
|
18
|
20
|
public class ExecutorBizImpl implements ExecutorBiz {
|
|
21
|
+ private static Logger logger = LoggerFactory.getLogger(ExecutorBizImpl.class);
|
19
|
22
|
|
20
|
23
|
@Override
|
21
|
24
|
public ReturnT<String> beat() {
|
|
@@ -55,25 +58,26 @@ public class ExecutorBizImpl implements ExecutorBiz {
|
55
|
58
|
if (!triggerParam.isGlueSwitch()) {
|
56
|
59
|
// bean model
|
57
|
60
|
|
58
|
|
- // valid handler instance
|
|
61
|
+ // valid handler
|
59
|
62
|
IJobHandler jobHandler = XxlJobExecutor.loadJobHandler(triggerParam.getExecutorHandler());
|
60
|
63
|
if (jobHandler==null) {
|
61
|
64
|
return new ReturnT(ReturnT.FAIL_CODE, "job handler for JobId=[" + triggerParam.getJobId() + "] not found.");
|
62
|
65
|
}
|
63
|
66
|
|
|
67
|
+ // valid exists job thread:change handler, need kill old thread
|
|
68
|
+ if (jobThread != null && jobThread.getHandler() != jobHandler) {
|
|
69
|
+ // kill old job thread
|
|
70
|
+ jobThread.toStop("更换任务模式或JobHandler,终止旧任务线程");
|
|
71
|
+ jobThread.interrupt();
|
|
72
|
+ XxlJobExecutor.removeJobThread(triggerParam.getJobId());
|
|
73
|
+ jobThread = null;
|
|
74
|
+ }
|
|
75
|
+
|
|
76
|
+ // make thread: new or exists invalid
|
64
|
77
|
if (jobThread == null) {
|
65
|
78
|
jobThread = XxlJobExecutor.registJobThread(triggerParam.getJobId(), jobHandler);
|
66
|
|
- } else {
|
67
|
|
- // job handler update, kill old job thread
|
68
|
|
- if (jobThread.getHandler() != jobHandler) {
|
69
|
|
- // kill old job thread
|
70
|
|
- jobThread.toStop("更换任务模式或JobHandler,终止旧任务线程");
|
71
|
|
- jobThread.interrupt();
|
72
|
|
-
|
73
|
|
- // new thread, with new job handler
|
74
|
|
- jobThread = XxlJobExecutor.registJobThread(triggerParam.getJobId(), jobHandler);
|
75
|
|
- }
|
76
|
79
|
}
|
|
80
|
+
|
77
|
81
|
} else {
|
78
|
82
|
// glue model
|
79
|
83
|
|
|
@@ -82,19 +86,29 @@ public class ExecutorBizImpl implements ExecutorBiz {
|
82
|
86
|
return new ReturnT(ReturnT.FAIL_CODE, "glueLoader for JobId=[" + triggerParam.getJobId() + "] not found.");
|
83
|
87
|
}
|
84
|
88
|
|
|
89
|
+ // valid exists job thread:change handler or glue timeout, need kill old thread
|
|
90
|
+ if (jobThread != null &&
|
|
91
|
+ !(jobThread.getHandler() instanceof GlueJobHandler
|
|
92
|
+ && ((GlueJobHandler) jobThread.getHandler()).getGlueUpdatetime()==triggerParam.getGlueUpdatetime() )) {
|
|
93
|
+ // change glue model or glue timeout, kill old job thread
|
|
94
|
+ jobThread.toStop("更换任务模式或JobHandler,终止旧任务线程");
|
|
95
|
+ jobThread.interrupt();
|
|
96
|
+ XxlJobExecutor.removeJobThread(triggerParam.getJobId());
|
|
97
|
+ jobThread = null;
|
|
98
|
+ }
|
|
99
|
+
|
|
100
|
+ // make thread: new or exists invalid
|
85
|
101
|
if (jobThread == null) {
|
86
|
|
- jobThread = XxlJobExecutor.registJobThread(triggerParam.getJobId(), new GlueJobHandler(triggerParam.getJobId()));
|
87
|
|
- } else {
|
88
|
|
- // job handler update, kill old job thread
|
89
|
|
- if (!(jobThread.getHandler() instanceof GlueJobHandler)) {
|
90
|
|
- // kill old job thread
|
91
|
|
- jobThread.toStop("更换任务模式或JobHandler,终止旧任务线程");
|
92
|
|
- jobThread.interrupt();
|
93
|
|
-
|
94
|
|
- // new thread, with new job handler
|
95
|
|
- jobThread = XxlJobExecutor.registJobThread(triggerParam.getJobId(), new GlueJobHandler(triggerParam.getJobId()));
|
|
102
|
+ IJobHandler jobHandler = null;
|
|
103
|
+ try {
|
|
104
|
+ jobHandler = GlueFactory.getInstance().loadNewInstance(triggerParam.getJobId());
|
|
105
|
+ } catch (Exception e) {
|
|
106
|
+ logger.error("", e);
|
|
107
|
+ return new ReturnT(ReturnT.FAIL_CODE, e.getMessage());
|
96
|
108
|
}
|
|
109
|
+ jobThread = XxlJobExecutor.registJobThread(triggerParam.getJobId(), new GlueJobHandler(jobHandler, triggerParam.getGlueUpdatetime()));
|
97
|
110
|
}
|
|
111
|
+
|
98
|
112
|
}
|
99
|
113
|
|
100
|
114
|
// push data to queue
|