job启动、停止、执行逻辑完善

This commit is contained in:
aoshiguchen
2022-09-14 19:36:57 +08:00
parent 6e3311ede3
commit ac0047f2db
2 changed files with 55 additions and 28 deletions
@@ -72,7 +72,6 @@ public class JobExecutor implements ApplicationRunner, IJobExecutor {
continue;
}
jobHandlerMap.put(jobHandler.name(), item);
runJobSet.add(jobHandler.name());
}
}
@@ -111,36 +110,47 @@ public class JobExecutor implements ApplicationRunner, IJobExecutor {
}
@Override
public synchronized void add(JobInfo jobInfo) throws JobException {
if (null == jobInfo || StringUtil.isEmpty(jobInfo.getName()) || StringUtil.isEmpty(jobInfo.getCron()) || jobInfoMap.containsKey(jobInfo.getName())) {
public void add(JobInfo jobInfo) throws JobException {
if (null == jobInfo || StringUtil.isEmpty(jobInfo.getId()) || StringUtil.isEmpty(jobInfo.getName()) ||
StringUtil.isEmpty(jobInfo.getCron()) || runJobSet.contains(jobInfo.getName())) {
return;
}
jobInfoMap.put(jobInfo.getId(), jobInfo);
synchronized (jobInfo.getId()) {
runJobSet.add(jobInfo.getName());
jobInfoMap.put(jobInfo.getId(), jobInfo);
TriggerKey triggerKey = TriggerKey.triggerKey(jobInfo.getId());
triggerKeyMap.put(jobInfo.getId(), triggerKey);
JobKey jobKey = new JobKey(jobInfo.getName());
TriggerKey triggerKey = TriggerKey.triggerKey(jobInfo.getId());
triggerKeyMap.put(jobInfo.getId(), triggerKey);
JobKey jobKey = new JobKey(jobInfo.getName());
CronScheduleBuilder cronScheduleBuilder = CronScheduleBuilder.cronSchedule(jobInfo.getCron()).withMisfireHandlingInstructionDoNothing();
CronTrigger cronTrigger = TriggerBuilder.newTrigger().withIdentity(triggerKey).withSchedule(cronScheduleBuilder).build();
JobDetail jobDetail = JobBuilder.newJob(JobBean.class).withIdentity(jobKey).build();
CronScheduleBuilder cronScheduleBuilder = CronScheduleBuilder.cronSchedule(jobInfo.getCron()).withMisfireHandlingInstructionDoNothing();
CronTrigger cronTrigger = TriggerBuilder.newTrigger().withIdentity(triggerKey).withSchedule(cronScheduleBuilder).build();
JobDetail jobDetail = JobBuilder.newJob(JobBean.class).withIdentity(jobKey).build();
try {
scheduler.scheduleJob(jobDetail, cronTrigger);
scheduler.start();
} catch (Exception e) {
throw new RuntimeException(String.format("新增job[name=%s]异常", jobInfo.getName()));
try {
scheduler.scheduleJob(jobDetail, cronTrigger);
scheduler.start();
} catch (Exception e) {
throw new RuntimeException(String.format("新增job[name=%s]异常", jobInfo.getName()));
}
}
}
@Override
public void remove(String jobName) {
runJobSet.remove(jobName);
public void remove(String jobId) {
JobInfo jobInfo = jobInfoMap.get(jobId);
if (null == jobInfo) {
return;
}
synchronized (jobId) {
runJobSet.remove(jobInfo.getName());
unscheduleJob(jobId);
}
}
@Override
public void trigger(String jobId, String param) {
doExecute(jobId, param);
}
public void execute(JobExecutionContext context) throws JobExecutionException {
@@ -150,19 +160,23 @@ public class JobExecutor implements ApplicationRunner, IJobExecutor {
String jobId = context.getTrigger().getKey().getName();
JobInfo jobInfo = jobInfoMap.get(jobId);
if (null == jobInfo) {
TriggerKey triggerKey = triggerKeyMap.get(jobId);
if (null != triggerKey) {
try {
scheduler.unscheduleJob(triggerKey);
} catch (Exception e) {
e.printStackTrace();
}
}
unscheduleJob(jobId);
return;
}
doExecute(jobId, jobInfo.getParam());
}
private void unscheduleJob(String jobId) {
TriggerKey triggerKey = triggerKeyMap.get(jobId);
if (null != triggerKey) {
try {
scheduler.unscheduleJob(triggerKey);
} catch (Exception e) {
e.printStackTrace();
}
}
}
private void doExecute(String jobId, String param) {
JobInfo jobInfo = jobInfoMap.get(jobId);
if (null == jobInfo) {
@@ -28,10 +28,13 @@ import fun.asgc.neutrino.core.annotation.NonIntercept;
import fun.asgc.neutrino.core.db.page.Page;
import fun.asgc.neutrino.core.db.page.PageQuery;
import fun.asgc.neutrino.core.quartz.IJobSource;
import fun.asgc.neutrino.core.quartz.JobExecutor;
import fun.asgc.neutrino.core.quartz.JobInfo;
import fun.asgc.neutrino.core.util.BeanManager;
import fun.asgc.neutrino.core.util.CollectionUtil;
import fun.asgc.neutrino.core.util.StringUtil;
import fun.asgc.neutrino.core.web.annotation.RequestBody;
import fun.asgc.neutrino.proxy.server.constant.EnableStatusEnum;
import fun.asgc.neutrino.proxy.server.constant.ExceptionConstant;
import fun.asgc.neutrino.proxy.server.controller.req.JobInfoExecuteReq;
import fun.asgc.neutrino.proxy.server.controller.req.JobInfoListReq;
@@ -70,12 +73,22 @@ public class JobInfoService implements IJobSource {
JobInfoDO jobInfoDO = jobInfoMapper.findById(req.getId());
ParamCheckUtil.checkNotNull(jobInfoDO, ExceptionConstant.JOB_INFO_NOT_EXIST);
jobInfoMapper.updateEnableStatus(req.getId(), req.getEnable(), new Date());
if (EnableStatusEnum.ENABLE.getStatus().equals(req.getEnable())) {
BeanManager.getBean(JobExecutor.class).add(new JobInfo()
.setId(String.valueOf(jobInfoDO.getId()))
.setName(jobInfoDO.getHandler())
.setDesc(jobInfoDO.getDesc())
.setCron(jobInfoDO.getCron())
.setParam(jobInfoDO.getParam())
);
} else {
BeanManager.getBean(JobExecutor.class).remove(String.valueOf(req.getId()));
}
return new JobInfoUpdateEnableStatusRes();
}
public JobInfoExecuteRes execute(JobInfoExecuteReq req) {
BeanManager.getBean(JobExecutor.class).trigger(String.valueOf(req.getId()), req.getParam());
return new JobInfoExecuteRes();
}