Commit 764fe3ea by huangfusuper

修正任务不遵循cron表达式 重复调用问题

parent da51635b
package com.byit.job; package com.byit.job;
import com.byit.task.JavaTaskJobTask;
import io.netty.util.HashedWheelTimer; import io.netty.util.HashedWheelTimer;
import io.netty.util.TimerTask; import io.netty.util.TimerTask;
import lombok.extern.slf4j.Slf4j;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
...@@ -11,6 +13,7 @@ import java.util.concurrent.TimeUnit; ...@@ -11,6 +13,7 @@ import java.util.concurrent.TimeUnit;
* @author: huangfu * @author: huangfu
* @date: 2019/11/15 14:59 * @date: 2019/11/15 14:59
**/ **/
@Slf4j
public class WorkRoulette { public class WorkRoulette {
/** /**
...@@ -24,6 +27,7 @@ public class WorkRoulette { ...@@ -24,6 +27,7 @@ public class WorkRoulette {
public static void addJob(TimerTask timerTask,long triggerNextTime) { public static void addJob(TimerTask timerTask,long triggerNextTime) {
log.info("-----添加任务{}-------",((JavaTaskJobTask)timerTask).getJavaTask().getTaskName());
HASHED_WHEEL_TIMER.newTimeout(timerTask, TimeUnit.MILLISECONDS.toNanos(triggerNextTime-System.currentTimeMillis()), TimeUnit.NANOSECONDS); HASHED_WHEEL_TIMER.newTimeout(timerTask, TimeUnit.MILLISECONDS.toNanos(triggerNextTime-System.currentTimeMillis()), TimeUnit.NANOSECONDS);
} }
......
...@@ -89,7 +89,7 @@ public class RunNodeServiceImpl implements RunNodeServer, ApplicationEventPublis ...@@ -89,7 +89,7 @@ public class RunNodeServiceImpl implements RunNodeServer, ApplicationEventPublis
jobTaskService.saveJobTasks(jobTasks); jobTaskService.saveJobTasks(jobTasks);
log.info("-------------【开始修改工作流{}的下次运行时间,以及各种状态】---------------",flow); log.info("-------------【开始修改工作流{}的下次运行时间,以及各种状态】---------------",flow);
try { try {
flow.setTriggerNextTime(new CronExpression(flow.getFlowCron()).getNextValidTimeAfter(new Date()).getTime()); flow.setTriggerNextTime(new CronExpression(flow.getFlowCron()).getNextValidTimeAfter(new Date(flow.getTriggerNextTime())).getTime());
} catch (ParseException e) { } catch (ParseException e) {
flow.setTriggerNextTime(999999999999999999L); flow.setTriggerNextTime(999999999999999999L);
......
package com.byit.service.mapservice.impl; package com.byit.service.mapservice.impl;
import cn.hutool.core.date.DateUtil;
import com.byit.dto.plugin.JavaTask; import com.byit.dto.plugin.JavaTask;
import com.byit.job.WorkRoulette; import com.byit.job.WorkRoulette;
import com.byit.job.utils.CronExpression; import com.byit.job.utils.CronExpression;
...@@ -29,7 +30,10 @@ public class JavaTaskAndLogServiceServiceImpl implements JavaTaskAndLogService { ...@@ -29,7 +30,10 @@ public class JavaTaskAndLogServiceServiceImpl implements JavaTaskAndLogService {
JavaTask updateJavaTask = new JavaTask(); JavaTask updateJavaTask = new JavaTask();
try { try {
BeanUtils.copyProperties(javaTask,updateJavaTask); BeanUtils.copyProperties(javaTask,updateJavaTask);
updateJavaTask.setTriggerTime(new CronExpression(javaTask.getCron()).getNextValidTimeAfter(new Date()).getTime()); Date nextValidTimeAfter = new CronExpression(javaTask.getCron()).getNextValidTimeAfter(new Date());
System.out.println("------下次执行时间为,{}-----"+ DateUtil.format(nextValidTimeAfter,"yyyy-MM-dd HH:mm:ss"));
System.out.println("------当前线程为:-----"+ Thread.currentThread().getName()+"----hash---"+Thread.currentThread().hashCode());
updateJavaTask.setTriggerTime(new CronExpression(javaTask.getCron()).getNextValidTimeAfter(new Date(javaTask.getTriggerTime())).getTime());
} catch (ParseException e) { } catch (ParseException e) {
javaTask.setTriggerTime(999999999999999999L); javaTask.setTriggerTime(999999999999999999L);
e.printStackTrace(); e.printStackTrace();
...@@ -39,5 +43,6 @@ public class JavaTaskAndLogServiceServiceImpl implements JavaTaskAndLogService { ...@@ -39,5 +43,6 @@ public class JavaTaskAndLogServiceServiceImpl implements JavaTaskAndLogService {
//保存到调度轮 //保存到调度轮
JavaTaskJobTask javaTaskJobTask = new JavaTaskJobTask(javaTask); JavaTaskJobTask javaTaskJobTask = new JavaTaskJobTask(javaTask);
WorkRoulette.addJob(javaTaskJobTask,javaTask.getTriggerTime()); WorkRoulette.addJob(javaTaskJobTask,javaTask.getTriggerTime());
System.out.println("-----------我执行结束了------------");
} }
} }
package com.byit.task; package com.byit.task;
import com.byit.conf.MythJobAutoConfigure; import com.byit.conf.MythJobAutoConfigure;
import com.byit.conf.RpcResultHttpCallback;
import com.byit.dto.plugin.JavaTask; import com.byit.dto.plugin.JavaTask;
import com.byit.dto.plugin.JobTaskRunLog;
import com.byit.enums.FlowPropertyEnum;
import com.byit.enums.NodePropertyEnum;
import com.byit.enums.NodeRunStatusPropertyEnum; import com.byit.enums.NodeRunStatusPropertyEnum;
import com.byit.enums.ScheduleTypeEnum; import com.byit.enums.ScheduleTypeEnum;
import com.byit.model.JobTaskRunLogWithBLOBs; import com.byit.model.JobTaskRunLogWithBLOBs;
import com.byit.packet.request.PluginRpcRequestPacket; import com.byit.packet.request.PluginRpcRequestPacket;
import com.byit.packet.response.PluginRpcResponsePacket; import com.byit.packet.response.PluginRpcResponsePacket;
import com.byit.param.CommunicationParam; import com.byit.param.CommunicationParam;
import com.byit.param.defaultparam.DefaultResultCallback;
import com.byit.service.impl.JobTaskRunLogServiceImpl; import com.byit.service.impl.JobTaskRunLogServiceImpl;
import com.byit.service.impl.RunJavaServiceImpl; import com.byit.service.impl.RunJavaServiceImpl;
import com.byit.util.ServiceInfoUtil; import com.byit.util.ServiceInfoUtil;
......
...@@ -21,7 +21,7 @@ public class JavaTaskThreadRunHelper extends BaseThreadRunHelper { ...@@ -21,7 +21,7 @@ public class JavaTaskThreadRunHelper extends BaseThreadRunHelper {
private final JavaTaskService javaTaskService; private final JavaTaskService javaTaskService;
private final JavaTaskAndLogService javaTaskAndLogService; private final JavaTaskAndLogService javaTaskAndLogService;
private static final long PRE_READ_MS = 7000; private static final long PRE_READ_MS = 7000;
private static final long JAVA_TASK_WAIT_TIME = 15000; private static final long JAVA_TASK_WAIT_TIME = 60000;
public JavaTaskThreadRunHelper(DataSource dataSource, JavaTaskService javaTaskService, JavaTaskAndLogService javaTaskAndLogService) { public JavaTaskThreadRunHelper(DataSource dataSource, JavaTaskService javaTaskService, JavaTaskAndLogService javaTaskAndLogService) {
this.dataSource = dataSource; this.dataSource = dataSource;
...@@ -32,7 +32,7 @@ public class JavaTaskThreadRunHelper extends BaseThreadRunHelper { ...@@ -32,7 +32,7 @@ public class JavaTaskThreadRunHelper extends BaseThreadRunHelper {
@Override @Override
public Long start() { public Long start() {
long nowTime = System.currentTimeMillis(); long nowTime = System.currentTimeMillis();
List<JavaTask> javaTasks = javaTaskService.findAllByTriggerTimeLessThanEqual(nowTime + PRE_READ_MS); List<JavaTask> javaTasks = javaTaskService.findAllByTriggerTimeLessThanEqual(nowTime + JAVA_TASK_WAIT_TIME);
if(CollectionUtil.isNotEmpty(javaTasks)){ if(CollectionUtil.isNotEmpty(javaTasks)){
javaTasks.forEach(javaTask -> { javaTasks.forEach(javaTask -> {
//查看该任务的剩余执行次数 //查看该任务的剩余执行次数
......
...@@ -13,15 +13,11 @@ import org.springframework.stereotype.Component; ...@@ -13,15 +13,11 @@ import org.springframework.stereotype.Component;
* @author huangfu * @author huangfu
*/ */
@Component @Component
@TaskHandler(expand = "{sadasdadasdasdasd}",cron = "0 0/1 * * * ?",taskName = "sentEmailServer",autoPublish = true,publishUrl = "http://127.0.0.1:8081/myth-job-admin/api/node/autoAddJavaTask") @TaskHandler(expand = "{sadasdadasdasdasd}",cron = "0 0/2 * * * ?",taskName = "sentEmailServer",autoPublish = true,publishUrl = "http://127.0.0.1:8081/myth-job-admin/api/node/autoAddJavaTask")
@Slf4j @Slf4j
public class SentEmailServer extends BaseJobHandler { public class SentEmailServer extends BaseJobHandler {
@Autowired
private CglibServerTest cglibServerTest;
@Override @Override
public ReturnResult<String> execute(CommunicationParam param) throws Exception { public ReturnResult<String> execute(CommunicationParam param) throws Exception {
cglibServerTest.test();
System.out.println("-------------SentEmailServer-被调度执行------------"+param); System.out.println("-------------SentEmailServer-被调度执行------------"+param);
return ReturnResult.SUCCESS; return ReturnResult.SUCCESS;
} }
......
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