d8b3675993643e6f6ea6a81b92ea092ff6119cfa.svn-base
3.7 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
package com.bjivt.job.core.handler;
import java.io.PrintWriter;
import java.io.StringWriter;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.TimeUnit;
import org.eclipse.jetty.util.ConcurrentHashSet;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.bjivt.job.core.handler.HandlerRepository.HandlerParamEnum;
import com.bjivt.job.core.handler.IJobHandler.JobHandleStatus;
import com.bjivt.job.core.log.XxlJobFileAppender;
import com.bjivt.job.core.util.HttpUtil;
/**
* handler thread
* @author xuxueli 2016-1-16 19:52:47
*/
public class HandlerThread extends Thread{
private static Logger logger = LoggerFactory.getLogger(HandlerThread.class);
private IJobHandler handler;
private LinkedBlockingQueue<Map<String, String>> handlerDataQueue;
private ConcurrentHashSet<String> logIdSet; // avoid repeat trigger for the same TRIGGER_LOG_ID
private boolean toStop = false;
public HandlerThread(IJobHandler handler) {
this.handler = handler;
handlerDataQueue = new LinkedBlockingQueue<Map<String,String>>();
logIdSet = new ConcurrentHashSet<String>();
}
public IJobHandler getHandler() {
return handler;
}
public void toStop() {
/**
* Thread.interrupt只支持终止线程的阻塞状态(wait、join、sleep),
* 在阻塞出抛出InterruptedException异常,但是并不会终止运行的线程本身;
* 所以需要注意,此处彻底销毁本线程,需要通过共享变量方式;
*/
this.toStop = true;
}
public void pushData(Map<String, String> param) {
if (param.get(HandlerParamEnum.LOG_ID.name())!=null && !logIdSet.contains(param.get(HandlerParamEnum.LOG_ID.name()))) {
handlerDataQueue.offer(param);
}
}
int i = 1;
@Override
public void run() {
while(!toStop){
try {
Map<String, String> handlerData = handlerDataQueue.poll();
if (handlerData!=null) {
i= 0;
String log_address = handlerData.get(HandlerParamEnum.LOG_ADDRESS.name());
String log_id = handlerData.get(HandlerParamEnum.LOG_ID.name());
String handler_params = handlerData.get(HandlerParamEnum.EXECUTOR_PARAMS.name());
logIdSet.remove(log_id);
// parse param
String[] handlerParams = null;
if (handler_params!=null && handler_params.trim().length()>0) {
handlerParams = handler_params.split(",");
} else {
handlerParams = new String[0];
}
// handle job
JobHandleStatus _status = JobHandleStatus.FAIL;
String _msg = null;
try {
XxlJobFileAppender.contextHolder.set(log_id);
logger.info(">>>>>>>>>>> job handle start.");
_status = handler.execute(handlerParams);
} catch (Exception e) {
logger.info("HandlerThread Exception:", e);
StringWriter out = new StringWriter();
e.printStackTrace(new PrintWriter(out));
_msg = out.toString();
}
logger.info(">>>>>>>>>>> job handle end, handlerParams:{}, _status:{}, _msg:{}",
new Object[]{handlerParams, _status, _msg});
// callback handler info
if (!toStop) {
HashMap<String, String> params = new HashMap<String, String>();
params.put("log_id", log_id);
params.put("status", _status.name());
params.put("msg", _msg);
HandlerRepository.pushCallBack(HttpUtil.addressToUrl(log_address), params);
}
} else {
i++;
logIdSet.clear();
try {
TimeUnit.MILLISECONDS.sleep(i * 100);
} catch (InterruptedException e) {
e.printStackTrace();
}
if (i>5) {
i= 0;
}
}
} catch (Exception e) {
logger.info("HandlerThread Exception:", e);
}
}
logger.info(">>>>>>>>>>>> job handlerThrad stoped, hashCode:{}", Thread.currentThread());
}
}