Commit 528460fd by huangfusuper

任务优先级开发

parent 2ca0baed
...@@ -64,6 +64,9 @@ ...@@ -64,6 +64,9 @@
<logger name="com.byit.selector" level="debug" additivity="false"> <logger name="com.byit.selector" level="debug" additivity="false">
<appender-ref ref="console"/> <appender-ref ref="console"/>
</logger> </logger>
<logger name="com.byit.task" level="debug" additivity="false">
<appender-ref ref="console"/>
</logger>
<logger name="org.springframework.web" level="debug"/> <logger name="org.springframework.web" level="debug"/>
......
...@@ -31,5 +31,27 @@ public class MythJobAutoConfigure { ...@@ -31,5 +31,27 @@ public class MythJobAutoConfigure {
TimeUnit.SECONDS, TimeUnit.SECONDS,
new LinkedBlockingQueue<Runnable>(100), new LinkedBlockingQueue<Runnable>(100),
r ->new Thread(r, "MythJob Thread of Job Run Log Callback Warehouse-" + r.hashCode())); r ->new Thread(r, "MythJob Thread of Job Run Log Callback Warehouse-" + r.hashCode()));
/**
* 高级任务的线程池
*/
public static final ThreadPoolExecutor ADVANCED_JOB_THREAD_POOL = new ThreadPoolExecutor(
10,
100,
60L,
TimeUnit.SECONDS,
new LinkedBlockingQueue<Runnable>(200),
r ->new Thread(r, "MythJob Thread of Job advanced run pool-" + r.hashCode()));
/**
* 低级任务的线程池
*/
public static final ThreadPoolExecutor LOW_LEVEL_JOB_THREAD_POOL = new ThreadPoolExecutor(
20,
200,
60L,
TimeUnit.SECONDS,
new LinkedBlockingQueue<Runnable>(300),
r ->new Thread(r, "MythJob Thread of Job low level run pool-" + r.hashCode()));
} }
...@@ -11,6 +11,8 @@ public enum NodePropertyEnum { ...@@ -11,6 +11,8 @@ public enum NodePropertyEnum {
ISNOT_VIRTUAL("1", "不是虚节点"), ISNOT_VIRTUAL("1", "不是虚节点"),
ON_FORK("0", "在工作流调度中"), ON_FORK("0", "在工作流调度中"),
OFF_FORK("1", "不在工作流调度中"), OFF_FORK("1", "不在工作流调度中"),
ADVANCED_NODE("2","高级节点"),
LOW_LEVEL_NODE("1","低级节点")
; ;
private String code; private String code;
......
...@@ -47,8 +47,8 @@ public class JobTaskRunLogAndJobTaskServiceImpl implements JobTaskRunLogAndJobTa ...@@ -47,8 +47,8 @@ public class JobTaskRunLogAndJobTaskServiceImpl implements JobTaskRunLogAndJobTa
Integer failedRemainingCount = jobTaskRunLog.getFailedRemainingCount(); Integer failedRemainingCount = jobTaskRunLog.getFailedRemainingCount();
//计算重试间隔时间 //计算重试间隔时间
Long nextReTime = (failedRetryCount-failedRemainingCount)*failedRetryInterval; Long nextReTime = (failedRetryCount-failedRemainingCount)*failedRetryInterval;
//获取本次执行完的时间 //获取当前时间
long endTime = jobTaskRunLog.getEndTime().getTime(); long endTime = System.currentTimeMillis();
//获取重试时间 //获取重试时间
Long nextTime = endTime+nextReTime; Long nextTime = endTime+nextReTime;
......
...@@ -2,6 +2,8 @@ package com.byit.task; ...@@ -2,6 +2,8 @@ package com.byit.task;
import cn.hutool.http.HttpUtil; import cn.hutool.http.HttpUtil;
import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSON;
import com.byit.conf.MythJobAutoConfigure;
import com.byit.enums.NodePropertyEnum;
import com.byit.job.dto.AdminSenPluginDto; import com.byit.job.dto.AdminSenPluginDto;
import com.byit.job.dto.DispatchResponseDto; import com.byit.job.dto.DispatchResponseDto;
import com.byit.job.enums.plugin.PluginEnum; import com.byit.job.enums.plugin.PluginEnum;
...@@ -35,6 +37,36 @@ public class JavaBeanJobTask implements TimerTask { ...@@ -35,6 +37,36 @@ public class JavaBeanJobTask implements TimerTask {
@Override @Override
public void run(Timeout timeout) { public void run(Timeout timeout) {
//获取任务级别 1最低 2最高
String priority = mythJobTaskSchedule.getPriority();
if(NodePropertyEnum.ADVANCED_NODE.getCode().equals(priority)){
log.debug("--------检测到高级节点--------");
MythJobAutoConfigure.ADVANCED_JOB_THREAD_POOL.execute(()->{
try {
runJob(mythJobTaskSchedule);
} catch (InterruptedException e) {
e.printStackTrace();
}
});
}else{
MythJobAutoConfigure.LOW_LEVEL_JOB_THREAD_POOL.execute(()->{
log.debug("--------检测到低级节点--------");
try {
runJob(mythJobTaskSchedule);
} catch (InterruptedException e) {
e.printStackTrace();
}
});
}
}
private void runJob(JobTaskSchedule mythJobTaskSchedule) throws InterruptedException{
if(!"start".equals(mythJobTaskSchedule.getNodeName())){
log.error("{}----------开始执行调度了吗?我要开始睡觉了-------------------",mythJobTaskSchedule.getNodeName());
Thread.sleep(100000);
}
//根据负责均衡方案获取对应IP //根据负责均衡方案获取对应IP
RpcLoadBalance rpcInvokerRouter = LoadBalance.match(mythJobTaskSchedule.getRoutingStrategy( ), LoadBalance.ROUND).rpcInvokerRouter; RpcLoadBalance rpcInvokerRouter = LoadBalance.match(mythJobTaskSchedule.getRoutingStrategy( ), LoadBalance.ROUND).rpcInvokerRouter;
try{ try{
......
...@@ -452,6 +452,7 @@ public class JobScheduleHelper{ ...@@ -452,6 +452,7 @@ public class JobScheduleHelper{
JobTaskRunLogWithBLOBs jobTaskRunLog = new JobTaskRunLogWithBLOBs(); JobTaskRunLogWithBLOBs jobTaskRunLog = new JobTaskRunLogWithBLOBs();
jobTaskRunLog.setRunId(mythJobTaskSchedule.getRunId()); jobTaskRunLog.setRunId(mythJobTaskSchedule.getRunId());
jobTaskRunLog.setIsVirtual(mythJobTaskSchedule.getIsVirtual());
jobTaskRunLog.setFlowId(mythJobTaskSchedule.getFlowId()); jobTaskRunLog.setFlowId(mythJobTaskSchedule.getFlowId());
jobTaskRunLog.setFlowName(mythJobTaskSchedule.getFlowName()); jobTaskRunLog.setFlowName(mythJobTaskSchedule.getFlowName());
jobTaskRunLog.setNodeId(mythJobTaskSchedule.getNodeId()); jobTaskRunLog.setNodeId(mythJobTaskSchedule.getNodeId());
......
...@@ -78,11 +78,31 @@ public class TestAddFlow { ...@@ -78,11 +78,31 @@ public class TestAddFlow {
pluginNodeConfig2.setNodeCron("0 0/7 * * * ? *"); pluginNodeConfig2.setNodeCron("0 0/7 * * * ? *");
pluginNodeConfig2.setNodeTimeout(-1L); pluginNodeConfig2.setNodeTimeout(-1L);
pluginNodeConfig2.setPluginUrls("http://127.0.0.1:8888"); pluginNodeConfig2.setPluginUrls("http://127.0.0.1:8888");
pluginNodeConfig2.setPriority("2"); pluginNodeConfig2.setPriority("1");
pluginNodeConfig2.setRoutingStrategy(LoadBalance.ROUND.name()); pluginNodeConfig2.setRoutingStrategy(LoadBalance.ROUND.name());
pluginNode2.setConfig(pluginNodeConfig2); pluginNode2.setConfig(pluginNodeConfig2);
pluginNode2.setDependNodeNameList(Collections.singletonList("start")); pluginNode2.setDependNodeNameList(Collections.singletonList("start"));
PluginNode binglie = new PluginNode();
PluginNodeConfig p2 = new PluginNodeConfig();
binglie.setName("并列节点");
binglie.setDesc("并列节点");
binglie.setType("node");
binglie.setAuthor("皇甫");
binglie.setJobType("JAVA");
binglie.setHandlerName("addJob");
binglie.setRunParam("并列节点");
p2.setFailedRetryCount(2);
p2.setFailedRetryInterval(TimeUnit.MINUTES.toSeconds(2));
p2.setNodeCron("0 0/7 * * * ? *");
p2.setNodeTimeout(-1L);
p2.setPluginUrls("http://127.0.0.1:8888");
p2.setPriority("1");
p2.setRoutingStrategy(LoadBalance.ROUND.name());
binglie.setConfig(p2);
binglie.setDependNodeNameList(Collections.singletonList("start"));
//报错节点 //报错节点
PluginNode pluginNode3 = new PluginNode(); PluginNode pluginNode3 = new PluginNode();
PluginNodeConfig pluginNodeConfig3 = new PluginNodeConfig(); PluginNodeConfig pluginNodeConfig3 = new PluginNodeConfig();
...@@ -93,7 +113,7 @@ public class TestAddFlow { ...@@ -93,7 +113,7 @@ public class TestAddFlow {
pluginNode3.setSuperSuccessRun("0"); pluginNode3.setSuperSuccessRun("0");
pluginNode3.setJobType("JAVA"); pluginNode3.setJobType("JAVA");
pluginNode3.setHandlerName("addJob"); pluginNode3.setHandlerName("addJob");
pluginNode3.setRunParam("1"); pluginNode3.setRunParam("1111");
pluginNodeConfig3.setFailedRetryCount(2); pluginNodeConfig3.setFailedRetryCount(2);
pluginNodeConfig3.setFailedRetryInterval(TimeUnit.MINUTES.toSeconds(2)); pluginNodeConfig3.setFailedRetryInterval(TimeUnit.MINUTES.toSeconds(2));
pluginNodeConfig3.setNodeCron("0 0/7 * * * ? *"); pluginNodeConfig3.setNodeCron("0 0/7 * * * ? *");
...@@ -120,7 +140,7 @@ public class TestAddFlow { ...@@ -120,7 +140,7 @@ public class TestAddFlow {
pluginNodeConfig4.setNodeCron("0 0/7 * * * ? *"); pluginNodeConfig4.setNodeCron("0 0/7 * * * ? *");
pluginNodeConfig4.setNodeTimeout(-1L); pluginNodeConfig4.setNodeTimeout(-1L);
pluginNodeConfig4.setPluginUrls("http://127.0.0.1:8888"); pluginNodeConfig4.setPluginUrls("http://127.0.0.1:8888");
pluginNodeConfig4.setPriority("2"); pluginNodeConfig4.setPriority("1");
pluginNodeConfig4.setRoutingStrategy(LoadBalance.ROUND.name()); pluginNodeConfig4.setRoutingStrategy(LoadBalance.ROUND.name());
pluginNode4.setConfig(pluginNodeConfig4); pluginNode4.setConfig(pluginNodeConfig4);
pluginNode4.setDependNodeNameList(Collections.singletonList("中间节点2")); pluginNode4.setDependNodeNameList(Collections.singletonList("中间节点2"));
...@@ -142,7 +162,7 @@ public class TestAddFlow { ...@@ -142,7 +162,7 @@ public class TestAddFlow {
pluginNodeConfig5.setPriority("2"); pluginNodeConfig5.setPriority("2");
pluginNodeConfig5.setRoutingStrategy(LoadBalance.ROUND.name()); pluginNodeConfig5.setRoutingStrategy(LoadBalance.ROUND.name());
pluginNode5.setConfig(pluginNodeConfig5); pluginNode5.setConfig(pluginNodeConfig5);
pluginNode5.setDependNodeNameList(Collections.singletonList("中间节点3")); pluginNode5.setDependNodeNameList(Arrays.asList("中间节点3","并列节点"));
return Arrays.asList(pluginNode5, pluginNode4, pluginNode3, pluginNode2, pluginNode1); return Arrays.asList(pluginNode5, pluginNode4, pluginNode3, pluginNode2, pluginNode1,binglie);
} }
} }
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment