|
|
|
@ -8,12 +8,10 @@ import com.xxl.job.admin.core.route.ExecutorRouteStrategyEnum;
|
|
|
|
|
import com.xxl.job.admin.core.schedule.XxlJobDynamicScheduler;
|
|
|
|
|
import com.xxl.job.admin.core.thread.JobFailMonitorHelper;
|
|
|
|
|
import com.xxl.job.admin.core.thread.JobRegistryMonitorHelper;
|
|
|
|
|
import com.xxl.job.core.biz.ExecutorBiz;
|
|
|
|
|
import com.xxl.job.core.biz.model.ReturnT;
|
|
|
|
|
import com.xxl.job.core.biz.model.TriggerParam;
|
|
|
|
|
import com.xxl.job.core.enums.ExecutorBlockStrategyEnum;
|
|
|
|
|
import com.xxl.job.core.enums.RegistryConfig;
|
|
|
|
|
import com.xxl.job.core.rpc.netcom.NetComClientProxy;
|
|
|
|
|
import org.apache.commons.collections.CollectionUtils;
|
|
|
|
|
import org.apache.commons.lang.StringUtils;
|
|
|
|
|
import org.quartz.JobExecutionContext;
|
|
|
|
@ -23,7 +21,9 @@ import org.slf4j.Logger;
|
|
|
|
|
import org.slf4j.LoggerFactory;
|
|
|
|
|
import org.springframework.scheduling.quartz.QuartzJobBean;
|
|
|
|
|
|
|
|
|
|
import java.util.*;
|
|
|
|
|
import java.util.ArrayList;
|
|
|
|
|
import java.util.Arrays;
|
|
|
|
|
import java.util.Date;
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* http job bean
|
|
|
|
@ -122,99 +122,12 @@ public class RemoteHttpJobBean extends QuartzJobBean {
|
|
|
|
|
}
|
|
|
|
|
triggerSb.append("<br>路由策略:").append(executorRouteStrategyEnum.name() + "-" + executorRouteStrategyEnum.getTitle());
|
|
|
|
|
|
|
|
|
|
// trigger remote executor
|
|
|
|
|
if (executorRouteStrategyEnum == ExecutorRouteStrategyEnum.FAILOVER) {
|
|
|
|
|
for (String address : addressList) {
|
|
|
|
|
// beat
|
|
|
|
|
ReturnT<String> beatResult = null;
|
|
|
|
|
try {
|
|
|
|
|
ExecutorBiz executorBiz = (ExecutorBiz) new NetComClientProxy(ExecutorBiz.class, address).getObject();
|
|
|
|
|
beatResult = executorBiz.beat();
|
|
|
|
|
} catch (Exception e) {
|
|
|
|
|
logger.error("", e);
|
|
|
|
|
beatResult = new ReturnT<String>(ReturnT.FAIL_CODE, ""+e );
|
|
|
|
|
}
|
|
|
|
|
triggerSb.append("<br>----------------------<br>")
|
|
|
|
|
.append("心跳检测:")
|
|
|
|
|
.append("<br>address:").append(address)
|
|
|
|
|
.append("<br>code:").append(beatResult.getCode())
|
|
|
|
|
.append("<br>msg:").append(beatResult.getMsg());
|
|
|
|
|
|
|
|
|
|
// beat success
|
|
|
|
|
if (beatResult.getCode() == ReturnT.SUCCESS_CODE) {
|
|
|
|
|
jobLog.setExecutorAddress(address);
|
|
|
|
|
|
|
|
|
|
ReturnT<String> runResult = runExecutor(triggerParam, address);
|
|
|
|
|
triggerSb.append("<br>----------------------<br>").append(runResult.getMsg());
|
|
|
|
|
|
|
|
|
|
return new ReturnT<String>(runResult.getCode(), triggerSb.toString());
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
return new ReturnT<String>(ReturnT.FAIL_CODE, triggerSb.toString());
|
|
|
|
|
} else if (executorRouteStrategyEnum == ExecutorRouteStrategyEnum.BUSYOVER) {
|
|
|
|
|
for (String address : addressList) {
|
|
|
|
|
// beat
|
|
|
|
|
ReturnT<String> idleBeatResult = null;
|
|
|
|
|
try {
|
|
|
|
|
ExecutorBiz executorBiz = (ExecutorBiz) new NetComClientProxy(ExecutorBiz.class, address).getObject();
|
|
|
|
|
idleBeatResult = executorBiz.idleBeat(triggerParam.getJobId());
|
|
|
|
|
} catch (Exception e) {
|
|
|
|
|
logger.error("", e);
|
|
|
|
|
idleBeatResult = new ReturnT<String>(ReturnT.FAIL_CODE, ""+e );
|
|
|
|
|
}
|
|
|
|
|
triggerSb.append("<br>----------------------<br>")
|
|
|
|
|
.append("空闲检测:")
|
|
|
|
|
.append("<br>address:").append(address)
|
|
|
|
|
.append("<br>code:").append(idleBeatResult.getCode())
|
|
|
|
|
.append("<br>msg:").append(idleBeatResult.getMsg());
|
|
|
|
|
|
|
|
|
|
// beat success
|
|
|
|
|
if (idleBeatResult.getCode() == ReturnT.SUCCESS_CODE) {
|
|
|
|
|
jobLog.setExecutorAddress(address);
|
|
|
|
|
|
|
|
|
|
ReturnT<String> runResult = runExecutor(triggerParam, address);
|
|
|
|
|
triggerSb.append("<br>----------------------<br>").append(runResult.getMsg());
|
|
|
|
|
|
|
|
|
|
return new ReturnT<String>(runResult.getCode(), triggerSb.toString());
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
return new ReturnT<String>(ReturnT.FAIL_CODE, triggerSb.toString());
|
|
|
|
|
} else {
|
|
|
|
|
// get address
|
|
|
|
|
String address = executorRouteStrategyEnum.getRouter().route(jobInfo.getId(), addressList);
|
|
|
|
|
jobLog.setExecutorAddress(address);
|
|
|
|
|
|
|
|
|
|
// run
|
|
|
|
|
ReturnT<String> runResult = runExecutor(triggerParam, address);
|
|
|
|
|
triggerSb.append("<br>----------------------<br>").append(runResult.getMsg());
|
|
|
|
|
|
|
|
|
|
return new ReturnT<String>(runResult.getCode(), triggerSb.toString());
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* run executor
|
|
|
|
|
* @param triggerParam
|
|
|
|
|
* @param address
|
|
|
|
|
* @return
|
|
|
|
|
*/
|
|
|
|
|
public ReturnT<String> runExecutor(TriggerParam triggerParam, String address){
|
|
|
|
|
ReturnT<String> runResult = null;
|
|
|
|
|
try {
|
|
|
|
|
ExecutorBiz executorBiz = (ExecutorBiz) new NetComClientProxy(ExecutorBiz.class, address).getObject();
|
|
|
|
|
runResult = executorBiz.run(triggerParam);
|
|
|
|
|
} catch (Exception e) {
|
|
|
|
|
logger.error("", e);
|
|
|
|
|
runResult = new ReturnT<String>(ReturnT.FAIL_CODE, ""+e );
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
StringBuffer sb = new StringBuffer("触发调度:");
|
|
|
|
|
sb.append("<br>address:").append(address);
|
|
|
|
|
sb.append("<br>code:").append(runResult.getCode());
|
|
|
|
|
sb.append("<br>msg:").append(runResult.getMsg());
|
|
|
|
|
runResult.setMsg(sb.toString());
|
|
|
|
|
// route run / trigger remote executor
|
|
|
|
|
ReturnT<String> routeRunResult = executorRouteStrategyEnum.getRouter().routeRun(triggerParam, addressList, jobLog);
|
|
|
|
|
triggerSb.append("<br>----------------------<br>").append(routeRunResult.getMsg());
|
|
|
|
|
return new ReturnT<String>(routeRunResult.getCode(), triggerSb.toString());
|
|
|
|
|
|
|
|
|
|
return runResult;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
}
|