|
@@ -22,15 +22,34 @@ public class TriggerCallbackThread {
|
22
|
22
|
return instance;
|
23
|
23
|
}
|
24
|
24
|
|
|
25
|
+ /**
|
|
26
|
+ * job results callback queue
|
|
27
|
+ */
|
25
|
28
|
private LinkedBlockingQueue<HandleCallbackParam> callBackQueue = new LinkedBlockingQueue<HandleCallbackParam>();
|
|
29
|
+ public static void pushCallBack(HandleCallbackParam callback){
|
|
30
|
+ getInstance().callBackQueue.add(callback);
|
|
31
|
+ logger.debug(">>>>>>>>>>> xxl-job, push callback request, logId:{}", callback.getLogId());
|
|
32
|
+ }
|
26
|
33
|
|
|
34
|
+ /**
|
|
35
|
+ * callback thread
|
|
36
|
+ */
|
27
|
37
|
private Thread triggerCallbackThread;
|
28
|
|
- private boolean toStop = false;
|
|
38
|
+ private volatile boolean toStop = false;
|
29
|
39
|
public void start() {
|
|
40
|
+
|
|
41
|
+ // valid
|
|
42
|
+ if (XxlJobExecutor.getAdminBizList() == null) {
|
|
43
|
+ logger.warn(">>>>>>>>>>>> xxl-job, executor callback config fail, adminAddresses is null.");
|
|
44
|
+ return;
|
|
45
|
+ }
|
|
46
|
+
|
30
|
47
|
triggerCallbackThread = new Thread(new Runnable() {
|
31
|
48
|
|
32
|
49
|
@Override
|
33
|
50
|
public void run() {
|
|
51
|
+
|
|
52
|
+ // normal callback
|
34
|
53
|
while(!toStop){
|
35
|
54
|
try {
|
36
|
55
|
HandleCallbackParam callback = getInstance().callBackQueue.take();
|
|
@@ -41,34 +60,27 @@ public class TriggerCallbackThread {
|
41
|
60
|
int drainToNum = getInstance().callBackQueue.drainTo(callbackParamList);
|
42
|
61
|
callbackParamList.add(callback);
|
43
|
62
|
|
44
|
|
- // valid
|
45
|
|
- if (XxlJobExecutor.getAdminBizList()==null) {
|
46
|
|
- logger.warn(">>>>>>>>>>>> xxl-job callback fail, adminAddresses is null, callbackParamList:{}", callbackParamList);
|
47
|
|
- continue;
|
48
|
|
- }
|
49
|
|
-
|
50
|
63
|
// callback, will retry if error
|
51
|
|
- for (AdminBiz adminBiz: XxlJobExecutor.getAdminBizList()) {
|
52
|
|
- try {
|
53
|
|
- ReturnT<String> callbackResult = adminBiz.callback(callbackParamList);
|
54
|
|
- if (callbackResult!=null && ReturnT.SUCCESS_CODE == callbackResult.getCode()) {
|
55
|
|
- callbackResult = ReturnT.SUCCESS;
|
56
|
|
- logger.info(">>>>>>>>>>> xxl-job callback success, callbackParamList:{}, callbackResult:{}", new Object[]{callbackParamList, callbackResult});
|
57
|
|
- break;
|
58
|
|
- } else {
|
59
|
|
- logger.info(">>>>>>>>>>> xxl-job callback fail, callbackParamList:{}, callbackResult:{}", new Object[]{callbackParamList, callbackResult});
|
60
|
|
- }
|
61
|
|
- } catch (Exception e) {
|
62
|
|
- logger.error(">>>>>>>>>>> xxl-job callback error, callbackParamList:{}", callbackParamList, e);
|
63
|
|
- //getInstance().callBackQueue.addAll(callbackParamList);
|
64
|
|
- }
|
|
64
|
+ if (callbackParamList!=null && callbackParamList.size()>0) {
|
|
65
|
+ doCallback(callbackParamList);
|
65
|
66
|
}
|
66
|
|
-
|
67
|
67
|
}
|
68
|
68
|
} catch (Exception e) {
|
69
|
69
|
logger.error(e.getMessage(), e);
|
70
|
70
|
}
|
71
|
71
|
}
|
|
72
|
+
|
|
73
|
+ // last callback
|
|
74
|
+ try {
|
|
75
|
+ List<HandleCallbackParam> callbackParamList = new ArrayList<HandleCallbackParam>();
|
|
76
|
+ int drainToNum = getInstance().callBackQueue.drainTo(callbackParamList);
|
|
77
|
+ if (callbackParamList!=null && callbackParamList.size()>0) {
|
|
78
|
+ doCallback(callbackParamList);
|
|
79
|
+ }
|
|
80
|
+ } catch (Exception e) {
|
|
81
|
+ logger.error(e.getMessage(), e);
|
|
82
|
+ }
|
|
83
|
+
|
72
|
84
|
}
|
73
|
85
|
});
|
74
|
86
|
triggerCallbackThread.setDaemon(true);
|
|
@@ -78,9 +90,27 @@ public class TriggerCallbackThread {
|
78
|
90
|
toStop = true;
|
79
|
91
|
}
|
80
|
92
|
|
81
|
|
- public static void pushCallBack(HandleCallbackParam callback){
|
82
|
|
- getInstance().callBackQueue.add(callback);
|
83
|
|
- logger.debug(">>>>>>>>>>> xxl-job, push callback request, logId:{}", callback.getLogId());
|
|
93
|
+ /**
|
|
94
|
+ * do callback, will retry if error
|
|
95
|
+ * @param callbackParamList
|
|
96
|
+ */
|
|
97
|
+ private void doCallback(List<HandleCallbackParam> callbackParamList){
|
|
98
|
+ // callback, will retry if error
|
|
99
|
+ for (AdminBiz adminBiz: XxlJobExecutor.getAdminBizList()) {
|
|
100
|
+ try {
|
|
101
|
+ ReturnT<String> callbackResult = adminBiz.callback(callbackParamList);
|
|
102
|
+ if (callbackResult!=null && ReturnT.SUCCESS_CODE == callbackResult.getCode()) {
|
|
103
|
+ callbackResult = ReturnT.SUCCESS;
|
|
104
|
+ logger.info(">>>>>>>>>>> xxl-job callback success, callbackParamList:{}, callbackResult:{}", new Object[]{callbackParamList, callbackResult});
|
|
105
|
+ break;
|
|
106
|
+ } else {
|
|
107
|
+ logger.info(">>>>>>>>>>> xxl-job callback fail, callbackParamList:{}, callbackResult:{}", new Object[]{callbackParamList, callbackResult});
|
|
108
|
+ }
|
|
109
|
+ } catch (Exception e) {
|
|
110
|
+ logger.error(">>>>>>>>>>> xxl-job callback error, callbackParamList:{}", callbackParamList, e);
|
|
111
|
+ //getInstance().callBackQueue.addAll(callbackParamList);
|
|
112
|
+ }
|
|
113
|
+ }
|
84
|
114
|
}
|
85
|
115
|
|
86
|
116
|
}
|