|
@@ -114,76 +114,82 @@ public class RemoteHttpJobBean extends QuartzJobBean {
|
114
|
114
|
return new ReturnT<String>(ReturnT.FAIL_CODE, triggerSb.toString());
|
115
|
115
|
}
|
116
|
116
|
|
117
|
|
- // trigger remote executor
|
118
|
|
- if (addressList.size() == 1) {
|
119
|
|
- String address = addressList.get(0);
|
120
|
|
- jobLog.setExecutorAddress(address);
|
121
|
|
-
|
122
|
|
- ReturnT<String> runResult = runExecutor(triggerParam, address);
|
123
|
|
- triggerSb.append("<br>----------------------<br>").append(runResult.getMsg());
|
|
117
|
+ // executor route strategy
|
|
118
|
+ ExecutorRouteStrategyEnum executorRouteStrategyEnum = ExecutorRouteStrategyEnum.match(jobInfo.getExecutorRouteStrategy(), null);
|
|
119
|
+ if (executorRouteStrategyEnum == null) {
|
|
120
|
+ triggerSb.append("<br>----------------------<br>").append("调度失败:").append("执行器路由策略为空");
|
|
121
|
+ return new ReturnT<String>(ReturnT.FAIL_CODE, triggerSb.toString());
|
|
122
|
+ }
|
|
123
|
+ triggerSb.append("<br>路由策略:").append(executorRouteStrategyEnum.name() + "-" + executorRouteStrategyEnum.getTitle());
|
124
|
124
|
|
125
|
|
- return new ReturnT<String>(runResult.getCode(), triggerSb.toString());
|
126
|
|
- } else {
|
127
|
|
- // executor route strategy
|
128
|
|
- ExecutorRouteStrategyEnum executorRouteStrategyEnum = ExecutorRouteStrategyEnum.match(jobInfo.getExecutorRouteStrategy(), null);
|
129
|
|
- triggerSb.append("<br>路由策略:").append(executorRouteStrategyEnum!=null?(executorRouteStrategyEnum.name() + "-" + executorRouteStrategyEnum.getTitle()):null);
|
130
|
|
- if (executorRouteStrategyEnum == null) {
|
131
|
|
- triggerSb.append("<br>----------------------<br>").append("调度失败:").append("执行器路由策略为空");
|
132
|
|
- return new ReturnT<String>(ReturnT.FAIL_CODE, triggerSb.toString());
|
133
|
|
- }
|
|
125
|
+ // trigger remote executor
|
|
126
|
+ if (executorRouteStrategyEnum == ExecutorRouteStrategyEnum.FAILOVER) {
|
|
127
|
+ for (String address : addressList) {
|
|
128
|
+ // beat
|
|
129
|
+ ReturnT<String> beatResult = null;
|
|
130
|
+ try {
|
|
131
|
+ ExecutorBiz executorBiz = (ExecutorBiz) new NetComClientProxy(ExecutorBiz.class, address).getObject();
|
|
132
|
+ beatResult = executorBiz.beat();
|
|
133
|
+ } catch (Exception e) {
|
|
134
|
+ logger.error("", e);
|
|
135
|
+ beatResult = new ReturnT<String>(ReturnT.FAIL_CODE, ""+e );
|
|
136
|
+ }
|
|
137
|
+ triggerSb.append("<br>----------------------<br>")
|
|
138
|
+ .append("心跳检测:")
|
|
139
|
+ .append("<br>address:").append(address)
|
|
140
|
+ .append("<br>code:").append(beatResult.getCode())
|
|
141
|
+ .append("<br>msg:").append(beatResult.getMsg());
|
134
|
142
|
|
135
|
|
- if (executorRouteStrategyEnum != ExecutorRouteStrategyEnum.FAILOVER) {
|
136
|
|
- // get address
|
137
|
|
- String address = executorRouteStrategyEnum.getRouter().route(jobInfo.getId(), addressList);
|
138
|
|
- jobLog.setExecutorAddress(address);
|
|
143
|
+ // beat success
|
|
144
|
+ if (beatResult.getCode() == ReturnT.SUCCESS_CODE) {
|
|
145
|
+ jobLog.setExecutorAddress(address);
|
139
|
146
|
|
140
|
|
- // run
|
141
|
|
- ReturnT<String> runResult = runExecutor(triggerParam, address);
|
142
|
|
- triggerSb.append("<br>----------------------<br>").append(runResult.getMsg());
|
|
147
|
+ ReturnT<String> runResult = runExecutor(triggerParam, address);
|
|
148
|
+ triggerSb.append("<br>----------------------<br>").append(runResult.getMsg());
|
143
|
149
|
|
144
|
|
- return new ReturnT<String>(runResult.getCode(), triggerSb.toString());
|
145
|
|
- } else {
|
146
|
|
- for (String address : addressList) {
|
147
|
|
- // beat
|
148
|
|
- ReturnT<String> beatResult = beatExecutor(address);
|
149
|
|
- triggerSb.append("<br>----------------------<br>").append(beatResult.getMsg());
|
|
150
|
+ return new ReturnT<String>(runResult.getCode(), triggerSb.toString());
|
|
151
|
+ }
|
|
152
|
+ }
|
|
153
|
+ return new ReturnT<String>(ReturnT.FAIL_CODE, triggerSb.toString());
|
|
154
|
+ } else if (executorRouteStrategyEnum == ExecutorRouteStrategyEnum.BUSYOVER) {
|
|
155
|
+ for (String address : addressList) {
|
|
156
|
+ // beat
|
|
157
|
+ ReturnT<String> idleBeatResult = null;
|
|
158
|
+ try {
|
|
159
|
+ ExecutorBiz executorBiz = (ExecutorBiz) new NetComClientProxy(ExecutorBiz.class, address).getObject();
|
|
160
|
+ idleBeatResult = executorBiz.idleBeat(triggerParam.getJobId());
|
|
161
|
+ } catch (Exception e) {
|
|
162
|
+ logger.error("", e);
|
|
163
|
+ idleBeatResult = new ReturnT<String>(ReturnT.FAIL_CODE, ""+e );
|
|
164
|
+ }
|
|
165
|
+ triggerSb.append("<br>----------------------<br>")
|
|
166
|
+ .append("空闲检测:")
|
|
167
|
+ .append("<br>address:").append(address)
|
|
168
|
+ .append("<br>code:").append(idleBeatResult.getCode())
|
|
169
|
+ .append("<br>msg:").append(idleBeatResult.getMsg());
|
150
|
170
|
|
151
|
|
- if (beatResult.getCode() == ReturnT.SUCCESS_CODE) {
|
152
|
|
- jobLog.setExecutorAddress(address);
|
|
171
|
+ // beat success
|
|
172
|
+ if (idleBeatResult.getCode() == ReturnT.SUCCESS_CODE) {
|
|
173
|
+ jobLog.setExecutorAddress(address);
|
153
|
174
|
|
154
|
|
- ReturnT<String> runResult = runExecutor(triggerParam, address);
|
155
|
|
- triggerSb.append("<br>----------------------<br>").append(runResult.getMsg());
|
|
175
|
+ ReturnT<String> runResult = runExecutor(triggerParam, address);
|
|
176
|
+ triggerSb.append("<br>----------------------<br>").append(runResult.getMsg());
|
156
|
177
|
|
157
|
|
- return new ReturnT<String>(runResult.getCode(), triggerSb.toString());
|
158
|
|
- }
|
|
178
|
+ return new ReturnT<String>(runResult.getCode(), triggerSb.toString());
|
159
|
179
|
}
|
160
|
|
- return new ReturnT<String>(ReturnT.FAIL_CODE, triggerSb.toString());
|
161
|
180
|
}
|
162
|
|
- }
|
163
|
|
- }
|
164
|
|
-
|
165
|
|
- /**
|
166
|
|
- * run executor
|
167
|
|
- * @param address
|
168
|
|
- * @return
|
169
|
|
- */
|
170
|
|
- public ReturnT<String> beatExecutor(String address){
|
171
|
|
- ReturnT<String> beatResult = null;
|
172
|
|
- try {
|
173
|
|
- ExecutorBiz executorBiz = (ExecutorBiz) new NetComClientProxy(ExecutorBiz.class, address).getObject();
|
174
|
|
- beatResult = executorBiz.beat();
|
175
|
|
- } catch (Exception e) {
|
176
|
|
- logger.error("", e);
|
177
|
|
- beatResult = new ReturnT<String>(ReturnT.FAIL_CODE, ""+e );
|
178
|
|
- }
|
|
181
|
+ return new ReturnT<String>(ReturnT.FAIL_CODE, triggerSb.toString());
|
|
182
|
+ } else {
|
|
183
|
+ // get address
|
|
184
|
+ String address = executorRouteStrategyEnum.getRouter().route(jobInfo.getId(), addressList);
|
|
185
|
+ jobLog.setExecutorAddress(address);
|
179
|
186
|
|
180
|
|
- StringBuffer sb = new StringBuffer("心跳检测:");
|
181
|
|
- sb.append("<br>address:").append(address);
|
182
|
|
- sb.append("<br>code:").append(beatResult.getCode());
|
183
|
|
- sb.append("<br>msg:").append(beatResult.getMsg());
|
184
|
|
- beatResult.setMsg(sb.toString());
|
|
187
|
+ // run
|
|
188
|
+ ReturnT<String> runResult = runExecutor(triggerParam, address);
|
|
189
|
+ triggerSb.append("<br>----------------------<br>").append(runResult.getMsg());
|
185
|
190
|
|
186
|
|
- return beatResult;
|
|
191
|
+ return new ReturnT<String>(runResult.getCode(), triggerSb.toString());
|
|
192
|
+ }
|
187
|
193
|
}
|
188
|
194
|
|
189
|
195
|
/**
|