Commit 7b4ef077 by guo_minglei@163.com

Merge remote-tracking branch 'origin/developer' into developer

parents c5c64824 764fe3ea
package com.byit.job;
import com.byit.task.JavaTaskJobTask;
import io.netty.util.HashedWheelTimer;
import io.netty.util.TimerTask;
import lombok.extern.slf4j.Slf4j;
import java.util.concurrent.TimeUnit;
......@@ -11,6 +13,7 @@ import java.util.concurrent.TimeUnit;
* @author: huangfu
* @date: 2019/11/15 14:59
**/
@Slf4j
public class WorkRoulette {
/**
......@@ -24,6 +27,7 @@ public class WorkRoulette {
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);
}
......
......@@ -89,7 +89,7 @@ public class RunNodeServiceImpl implements RunNodeServer, ApplicationEventPublis
jobTaskService.saveJobTasks(jobTasks);
log.info("-------------【开始修改工作流{}的下次运行时间,以及各种状态】---------------",flow);
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) {
flow.setTriggerNextTime(999999999999999999L);
......
package com.byit.service.mapservice.impl;
import cn.hutool.core.date.DateUtil;
import com.byit.dto.plugin.JavaTask;
import com.byit.job.WorkRoulette;
import com.byit.job.utils.CronExpression;
......@@ -29,7 +30,10 @@ public class JavaTaskAndLogServiceServiceImpl implements JavaTaskAndLogService {
JavaTask updateJavaTask = new JavaTask();
try {
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) {
javaTask.setTriggerTime(999999999999999999L);
e.printStackTrace();
......@@ -39,5 +43,6 @@ public class JavaTaskAndLogServiceServiceImpl implements JavaTaskAndLogService {
//保存到调度轮
JavaTaskJobTask javaTaskJobTask = new JavaTaskJobTask(javaTask);
WorkRoulette.addJob(javaTaskJobTask,javaTask.getTriggerTime());
System.out.println("-----------我执行结束了------------");
}
}
package com.byit.task;
import com.byit.conf.MythJobAutoConfigure;
import com.byit.conf.RpcResultHttpCallback;
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.ScheduleTypeEnum;
import com.byit.model.JobTaskRunLogWithBLOBs;
import com.byit.packet.request.PluginRpcRequestPacket;
import com.byit.packet.response.PluginRpcResponsePacket;
import com.byit.param.CommunicationParam;
import com.byit.param.defaultparam.DefaultResultCallback;
import com.byit.service.impl.JobTaskRunLogServiceImpl;
import com.byit.service.impl.RunJavaServiceImpl;
import com.byit.util.ServiceInfoUtil;
......
......@@ -21,7 +21,7 @@ public class JavaTaskThreadRunHelper extends BaseThreadRunHelper {
private final JavaTaskService javaTaskService;
private final JavaTaskAndLogService javaTaskAndLogService;
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) {
this.dataSource = dataSource;
......@@ -32,7 +32,7 @@ public class JavaTaskThreadRunHelper extends BaseThreadRunHelper {
@Override
public Long start() {
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)){
javaTasks.forEach(javaTask -> {
//查看该任务的剩余执行次数
......
......@@ -13,15 +13,11 @@ import org.springframework.stereotype.Component;
* @author huangfu
*/
@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
public class SentEmailServer extends BaseJobHandler {
@Autowired
private CglibServerTest cglibServerTest;
@Override
public ReturnResult<String> execute(CommunicationParam param) throws Exception {
cglibServerTest.test();
System.out.println("-------------SentEmailServer-被调度执行------------"+param);
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