refactor(task):重构任务调度与执行逻辑

- 统一使用 TaskJobKeys 常量替代硬编码字符串参数
- 优化线程变量命名 subMainThread -> workThread
- 完善任务创建时的参数校验逻辑
- 优化任务模板解析异常处理
- 调整任务调度构建逻辑为独立方法 buildAndSchedule
- 修复任务暂停逻辑 resumeTrigger -> pauseTrigger
- 增强获取任务模板空值判断逻辑
- 清理冗余代码与无用注释
- 使用 Lombok 注解简化依赖注入方式
This commit is contained in:
2025-10-27 13:29:34 +08:00
parent 1c5e1e34e0
commit 78725a4c25
9 changed files with 249 additions and 390 deletions
@@ -1,19 +1,20 @@
/*
* Copyright 2021-2025 Odboy
* Copyright 2021-2025 Odboy
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package cn.odboy.framework.exception;
package cn.odboy.framework.exception.web;
import lombok.Getter;
import org.springframework.http.HttpStatus;
@@ -36,4 +37,10 @@ public class BadRequestException extends RuntimeException {
super(msg);
this.status = status.value();
}
public BadRequestException(Throwable cause) {
// FIXME 20250827 修复完整异常链暴露给用户的问题
super(cause.getMessage());
this.status = BAD_REQUEST.value();
}
}
@@ -0,0 +1,2 @@
package cn.odboy.framework.exception.web;public class ServerException {
}
@@ -1,46 +0,0 @@
/*
* Copyright 2021-2025 Odboy
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package cn.odboy.framework.exception.web;
import lombok.Getter;
import org.springframework.http.HttpStatus;
import static org.springframework.http.HttpStatus.BAD_REQUEST;
/**
* 统一异常处理
*/
@Getter
public class BadRequestException extends RuntimeException {
private Integer status = BAD_REQUEST.value();
public BadRequestException(String msg) {
super(msg);
}
public BadRequestException(HttpStatus status, String msg) {
super(msg);
this.status = status.value();
}
public BadRequestException(Throwable cause) {
// FIXME 20250827 修复完整异常链暴露给用户的问题
super(cause.getMessage());
this.status = BAD_REQUEST.value();
}
}
@@ -41,16 +41,16 @@ public class QuartzManage {
JobDetail jobDetail = JobBuilder.newJob(ExecutionJobBean.class).withIdentity(JOB_NAME + quartzJob.getId()).build();
// 通过触发器名和cron 表达式创建 Trigger
Trigger cronTrigger = TriggerBuilder.newTrigger().withIdentity(JOB_NAME + quartzJob.getId()).startNow().withSchedule(CronScheduleBuilder.cronSchedule(quartzJob.getCronExpression())).build();
Trigger trigger = TriggerBuilder.newTrigger().withIdentity(JOB_NAME + quartzJob.getId()).startNow().withSchedule(CronScheduleBuilder.cronSchedule(quartzJob.getCronExpression())).build();
cronTrigger.getJobDataMap().put(SystemQuartzJobTb.JOB_KEY, quartzJob);
trigger.getJobDataMap().put(SystemQuartzJobTb.JOB_KEY, quartzJob);
// 重置启动时间
((CronTriggerImpl) cronTrigger).setStartTime(new Date());
((CronTriggerImpl) trigger).setStartTime(new Date());
// 执行定时任务,如果是持久化的,这里会报错,捕获输出
try {
scheduler.scheduleJob(jobDetail, cronTrigger);
scheduler.scheduleJob(jobDetail, trigger);
} catch (ObjectAlreadyExistsException e) {
log.warn("定时任务已存在,跳过加载");
}
@@ -0,0 +1,16 @@
package cn.odboy.task.constant;
/**
* 任务共用Keys
*
* @author odboy
* @date 2025-10-25
*/
public interface TaskJobKeys {
String ID = "id";
String CONTEXT_NAME = "contextName";
String LANGUAGE = "language";
String ENV_ALIAS = "envAlias";
String CHANGE_TYPE = "changeType";
String RETRY_NODE_CODE = "retryNodeCode";
}
@@ -19,6 +19,7 @@ package cn.odboy.task.core;
import cn.hutool.core.util.StrUtil;
import cn.odboy.framework.context.CsSpringBeanHolder;
import cn.odboy.framework.exception.web.BadRequestException;
import cn.odboy.task.constant.TaskJobKeys;
import cn.odboy.task.dal.dataobject.TaskInstanceDetailTb;
import cn.odboy.task.dal.dataobject.TaskInstanceInfoTb;
import cn.odboy.task.dal.model.TaskTemplateNodeVo;
@@ -40,7 +41,7 @@ import java.util.stream.Collectors;
@Slf4j
public class TaskJobBean extends QuartzJobBean implements InterruptableJob {
private Thread subMainThread = null;
private Thread workThread = null;
private String getBeanAlias(String code) {
// node_init
@@ -50,19 +51,19 @@ public class TaskJobBean extends QuartzJobBean implements InterruptableJob {
@Override
public void executeInternal(JobExecutionContext context) {
this.subMainThread = Thread.currentThread();
this.workThread = Thread.currentThread();
// ========================== 获取代理类 ==========================
TaskInstanceInfoService taskInstanceInfoService = CsSpringBeanHolder.getBean(TaskInstanceInfoService.class);
TaskInstanceDetailService taskInstanceDetailService = CsSpringBeanHolder.getBean(TaskInstanceDetailService.class);
// ========================== 获取参数 ==========================
JobDataMap dataMap = context.getMergedJobDataMap();
long id = dataMap.getLong("id");
long id = dataMap.getLong(TaskJobKeys.ID);
// ========================== 获取任务编排模板 ==========================
TaskInstanceInfoTb taskInstanceInfoVo = taskInstanceInfoService.getById(id);
String templateInfo = taskInstanceInfoVo.getTemplate();
List<TaskTemplateNodeVo> taskTemplateNodeVos = JSON.parseArray(templateInfo, TaskTemplateNodeVo.class);
// ========================== 判断是否重试任务 ==========================
String retryNodeCode = dataMap.getString("retryNodeCode");
String retryNodeCode = dataMap.getString(TaskJobKeys.RETRY_NODE_CODE);
if (StrUtil.isBlank(retryNodeCode)) {
executeNormalTask(taskInstanceInfoService, taskInstanceDetailService, id, dataMap, taskTemplateNodeVos);
} else {
@@ -191,151 +192,8 @@ public class TaskJobBean extends QuartzJobBean implements InterruptableJob {
@Override
public void interrupt() {
if (subMainThread != null) {
subMainThread.stop();
if (workThread != null) {
workThread.stop();
}
}
}
//@Slf4j
//public class TaskJobBean extends QuartzJobBean implements InterruptableJob {
// private Thread subMainThread = null;
//
// private String getBeanAlias(String code) {
// // node_init
// String[] s = code.split("_");
// return Arrays.stream(s).map(StrUtil::upperFirst).collect(Collectors.joining());
// }
//
// @Override
// public void executeInternal(JobExecutionContext context) {
// this.subMainThread = Thread.currentThread();
// // ========================== 获取代理类 ==========================
// TaskInstanceInfoService taskInstanceInfoService = CsSpringBeanHolder.getBean(TaskInstanceInfoService.class);
// TaskInstanceDetailService taskInstanceDetailService = CsSpringBeanHolder.getBean(TaskInstanceDetailService.class);
// // ========================== 获取参数 ==========================
// JobDataMap dataMap = context.getMergedJobDataMap();
// long id = dataMap.getLong("id");
// // ========================== 获取任务编排模板 ==========================
// TaskInstanceInfoTb taskInstanceInfoVo = taskInstanceInfoService.getById(id);
// String templateInfo = taskInstanceInfoVo.getTemplate();
// List<TaskTemplateNodeVo> taskTemplateNodeVos = JSON.parseArray(templateInfo, TaskTemplateNodeVo.class);
// // ========================== 判断是否重试任务 ==========================
// String retryNodeCode = dataMap.getString("retryNodeCode");
// if (StrUtil.isBlank(retryNodeCode)) {
// // ========================== 初始化执行明细 ==========================
// List<TaskInstanceDetailTb> taskInstanceDetails = new ArrayList<>();
// for (TaskTemplateNodeVo taskTemplateNodeVo : taskTemplateNodeVos) {
// TaskInstanceDetailTb taskInstanceDetail = new TaskInstanceDetailTb();
// taskInstanceDetail.setInstanceId(id);
// taskInstanceDetail.setFinishTime(null);
// taskInstanceDetail.setBizCode(taskTemplateNodeVo.getCode());
// taskInstanceDetail.setBizName(taskTemplateNodeVo.getName());
// taskInstanceDetail.setExecuteInfo("未开始");
// taskInstanceDetail.setExecuteStatus("pending");
// taskInstanceDetails.add(taskInstanceDetail);
// }
// taskInstanceDetailService.saveBatch(taskInstanceDetails);
// Map<String, Long> codeIdMap = taskInstanceDetails.stream().collect(Collectors.toMap(TaskInstanceDetailTb::getBizCode, TaskInstanceDetailTb::getId));
// // ========================== 顺序执行 ==========================
// try {
// for (TaskTemplateNodeVo taskTemplateNodeVo : taskTemplateNodeVos) {
// String code = taskTemplateNodeVo.getCode();
// try {
// TaskStepExecutor executor = CsSpringBeanHolder.getBean("taskStep" + getBeanAlias(code));
// executor.execute(codeIdMap.getOrDefault(code, null), dataMap, taskTemplateNodeVo, new TaskStepCallback() {
// @Override
// public void onStart() {
// taskInstanceDetailService.fastStart(id, code, dataMap);
// }
//
// @Override
// public void onFinish(String executeInfo) {
// taskInstanceDetailService.fastSuccessWithInfo(id, code, executeInfo);
// }
// });
// } catch (Exception e) {
// taskInstanceDetailService.fastFailWithInfo(id, code, e.getMessage());
// throw new RuntimeException(e);
// }
// }
// taskInstanceInfoService.fastSuccessWithData(id, dataMap);
// } catch (Exception e) {
// log.error("任务执行失败", e);
// taskInstanceInfoService.fastFailWithMessageData(id, e.getMessage(), dataMap);
// }
// } else {
// boolean isFound = false;
// // ========================== 初始化执行明细 ==========================
// List<TaskInstanceDetailTb> taskInstanceDetails = new ArrayList<>();
// List<TaskTemplateNodeVo> taskTemplateNodeRetrys = new ArrayList<>();
// for (TaskTemplateNodeVo taskTemplateNodeVo : taskTemplateNodeVos) {
// if (isFound) {
// TaskInstanceDetailTb taskInstanceDetail = new TaskInstanceDetailTb();
// taskInstanceDetail.setInstanceId(id);
// taskInstanceDetail.setFinishTime(null);
// taskInstanceDetail.setBizCode(taskTemplateNodeVo.getCode());
// taskInstanceDetail.setBizName(taskTemplateNodeVo.getName());
// taskInstanceDetail.setExecuteInfo("未开始");
// taskInstanceDetail.setExecuteStatus("pending");
// taskInstanceDetails.add(taskInstanceDetail);
// taskTemplateNodeRetrys.add(taskTemplateNodeVo);
// continue;
// }
// if (taskTemplateNodeVo.getCode().equals(retryNodeCode)) {
// if (!taskTemplateNodeVo.getRetry()) {
// throw new BadRequestException("节点 " + retryNodeCode + "不支持重试");
// }
// TaskInstanceDetailTb taskInstanceDetail = new TaskInstanceDetailTb();
// taskInstanceDetail.setInstanceId(id);
// taskInstanceDetail.setFinishTime(null);
// taskInstanceDetail.setBizCode(taskTemplateNodeVo.getCode());
// taskInstanceDetail.setBizName(taskTemplateNodeVo.getName());
// taskInstanceDetail.setExecuteInfo("未开始");
// taskInstanceDetail.setExecuteStatus("pending");
// taskInstanceDetails.add(taskInstanceDetail);
// taskTemplateNodeRetrys.add(taskTemplateNodeVo);
// isFound = true;
// }
// }
// taskInstanceDetailService.removeByInstanceIdAndBizCodeList(id, taskInstanceDetails.stream().map(TaskInstanceDetailTb::getBizCode).distinct().collect(Collectors.toList()));
// taskInstanceDetailService.saveBatch(taskInstanceDetails);
// Map<String, Long> codeIdMap = taskInstanceDetails.stream().collect(Collectors.toMap(TaskInstanceDetailTb::getBizCode, TaskInstanceDetailTb::getId));
// // ========================== 顺序执行 ==========================
// try {
// for (TaskTemplateNodeVo taskTemplateNodeVo : taskTemplateNodeRetrys) {
// String code = taskTemplateNodeVo.getCode();
// try {
// TaskStepExecutor executor = CsSpringBeanHolder.getBean("taskStep" + getBeanAlias(code));
// executor.execute(codeIdMap.getOrDefault(code, null), dataMap, taskTemplateNodeVo, new TaskStepCallback() {
// @Override
// public void onStart() {
// taskInstanceDetailService.fastStart(id, code, dataMap);
// }
//
// @Override
// public void onFinish(String executeInfo) {
// taskInstanceDetailService.fastSuccessWithInfo(id, code, executeInfo);
// }
// });
// } catch (Exception e) {
// taskInstanceDetailService.fastFailWithInfo(id, code, e.getMessage());
// throw new RuntimeException(e);
// }
// }
// taskInstanceInfoService.fastSuccessWithData(id, dataMap);
// } catch (Exception e) {
// log.error("任务执行失败", e);
// taskInstanceInfoService.fastFailWithMessageData(id, e.getMessage(), dataMap);
// }
// }
// }
//
// @Override
// public void interrupt() {
// if (subMainThread != null) {
// subMainThread.stop();
// }
// }
//}
@@ -16,8 +16,11 @@
package cn.odboy.task.core;
import cn.hutool.core.bean.BeanUtil;
import cn.hutool.core.util.StrUtil;
import cn.odboy.framework.context.CsSpringBeanHolder;
import cn.odboy.framework.exception.web.BadRequestException;
import cn.odboy.task.constant.TaskChangeTypeEnum;
import cn.odboy.task.constant.TaskJobKeys;
import cn.odboy.task.constant.TaskStatusEnum;
import cn.odboy.task.dal.dataobject.TaskInstanceDetailTb;
import cn.odboy.task.dal.dataobject.TaskInstanceInfoTb;
@@ -30,34 +33,29 @@ import cn.odboy.task.service.TaskInstanceInfoService;
import cn.odboy.task.service.TaskTemplateInfoService;
import cn.odboy.util.CsDateUtil;
import com.alibaba.fastjson2.JSON;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.quartz.*;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import javax.annotation.Resource;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;
@Slf4j
@Component
@RequiredArgsConstructor
public class TaskManage {
@Resource
private Scheduler scheduler;
@Autowired
private TaskTemplateInfoService taskTemplateInfoService;
@Autowired
private TaskInstanceInfoService taskInstanceInfoService;
@Autowired
private TaskInstanceDetailService taskInstanceDetailService;
private final Scheduler scheduler;
private final TaskTemplateInfoService taskTemplateInfoService;
private final TaskInstanceInfoService taskInstanceInfoService;
private final TaskInstanceDetailService taskInstanceDetailService;
/**
* 创建任务单
* 创建任务单(注:同一应用同一变更类型无法并发多实例)
*
* @param contextName 上下文名称,这里特指应用名
* @param changeTypeEnum 变更类型
* @param language 开发语言、资源版本
* @param envAlias 环境别名
* @param source 来源
* @param reason 变更原因
@@ -65,6 +63,21 @@ public class TaskManage {
*/
@Transactional(rollbackFor = Exception.class)
public TaskInstanceInfoTb createJob(String contextName, TaskChangeTypeEnum changeTypeEnum, String language, String envAlias, String source, String reason, JobDataMap dataMap) {
if (StrUtil.isBlank(contextName)) {
throw new BadRequestException("参数contextName必填");
}
if (changeTypeEnum == null) {
throw new BadRequestException("参数changeTypeEnum必填");
}
if (StrUtil.isBlank(language)) {
throw new BadRequestException("参数language必填");
}
if (StrUtil.isBlank(envAlias)) {
throw new BadRequestException("参数envAlias必填");
}
if (StrUtil.isBlank(source)) {
throw new BadRequestException("参数source必填");
}
if (dataMap == null) {
dataMap = new JobDataMap();
}
@@ -74,19 +87,33 @@ public class TaskManage {
throw new BadRequestException("没有查询到任务编排模板");
}
String templateInfo = taskInstanceInfoVo.getTemplateInfo();
List<TaskTemplateNodeVo> taskTemplateNodeVos = JSON.parseArray(templateInfo, TaskTemplateNodeVo.class);
List<TaskTemplateNodeVo> taskTemplateNodeVos;
try {
taskTemplateNodeVos = JSON.parseArray(templateInfo, TaskTemplateNodeVo.class);
} catch (Exception e) {
log.error("任务编排模板解析失败", e);
throw new BadRequestException("任务编排模板解析失败");
}
if (taskTemplateNodeVos == null || taskTemplateNodeVos.isEmpty()) {
throw new BadRequestException("没有查询到任务编排明细");
}
// ========================== 传递参数 ==========================
// 应用名、资源类型
dataMap.put("contextName", contextName);
// 开发语言、资源版本
dataMap.put("language", language);
dataMap.put("envAlias", envAlias);
// 变更类型
dataMap.put("changeType", changeTypeEnum.getCode());
// ========================== 创建任务 ==========================
TaskManage taskManage = CsSpringBeanHolder.getBean(TaskManage.class);
TaskInstanceInfoTb newInstance = taskManage.saveTaskInstanceInfoTb(contextName, changeTypeEnum, language, envAlias, source, reason, dataMap, templateInfo);
dataMap.put(TaskJobKeys.ID, newInstance.getId());
// 应用名、资源类型
dataMap.put(TaskJobKeys.CONTEXT_NAME, contextName);
// 开发语言、资源版本
dataMap.put(TaskJobKeys.LANGUAGE, language);
dataMap.put(TaskJobKeys.ENV_ALIAS, envAlias);
// 变更类型
dataMap.put(TaskJobKeys.CHANGE_TYPE, changeTypeEnum.getCode());
buildAndSchedule(changeTypeEnum.getCode(), contextName, dataMap);
return newInstance;
}
@Transactional(rollbackFor = Exception.class)
public TaskInstanceInfoTb saveTaskInstanceInfoTb(String contextName, TaskChangeTypeEnum changeTypeEnum, String language, String envAlias, String source, String reason, JobDataMap dataMap, String templateInfo) {
TaskInstanceInfoTb newInstance = new TaskInstanceInfoTb();
newInstance.setContextName(contextName);
newInstance.setLanguage(language);
@@ -99,21 +126,24 @@ public class TaskManage {
newInstance.setTemplate(templateInfo);
newInstance.setJobData(JSON.toJSONString(dataMap));
taskInstanceInfoService.save(newInstance);
dataMap.put("id", newInstance.getId());
return newInstance;
}
private void buildAndSchedule(String changeTypeEnum, String contextName, JobDataMap dataMap) {
// ========================== 执行任务 ==========================
JobKey jobKey = JobKey.jobKey(changeTypeEnum.getCode(), contextName);
TriggerKey triggerKey = TriggerKey.triggerKey(changeTypeEnum.getCode(), contextName);
JobKey jobKey = JobKey.jobKey(changeTypeEnum, contextName);
TriggerKey triggerKey = TriggerKey.triggerKey(changeTypeEnum, contextName);
JobDetail jobDetail = JobBuilder.newJob(TaskJobBean.class).withIdentity(jobKey).usingJobData(dataMap).build();
Trigger cronTrigger = TriggerBuilder.newTrigger().withIdentity(triggerKey).startNow().build();
Trigger trigger = TriggerBuilder.newTrigger().withIdentity(triggerKey).startNow().build();
try {
scheduler.scheduleJob(jobDetail, cronTrigger);
scheduler.scheduleJob(jobDetail, trigger);
} catch (ObjectAlreadyExistsException e) {
log.error("任务jobKey={},triggerKey={}任务已存在,跳过加载", jobKey.getName(), triggerKey.getName(), e);
throw new BadRequestException("任务已存在,跳过加载");
} catch (SchedulerException e) {
log.error("任务执行失败", e);
log.error("任务jobKey={},triggerKey={}执行失败", jobKey.getName(), triggerKey.getName(), e);
throw new BadRequestException(e);
}
return newInstance;
}
@Transactional(rollbackFor = Exception.class)
@@ -132,7 +162,7 @@ public class TaskManage {
JobKey jobKey = JobKey.jobKey(changeType, contextName);
TriggerKey triggerKey = TriggerKey.triggerKey(changeType, contextName);
// 停止触发器
scheduler.resumeTrigger(triggerKey);
scheduler.pauseTrigger(triggerKey);
// 移除触发器
scheduler.unscheduleJob(triggerKey);
// 删除并中断任务
@@ -164,21 +194,10 @@ public class TaskManage {
taskInstanceInfoTb.setFinishTime(null);
taskInstanceInfoService.updateById(taskInstanceInfoTb);
// ========================== 执行任务 ==========================
dataMap.put("retryNodeCode", retryNodeCode);
dataMap.put(TaskJobKeys.RETRY_NODE_CODE, retryNodeCode);
String changeType = taskInstanceInfoTb.getChangeType();
String contextName = taskInstanceInfoTb.getContextName();
JobKey jobKey = JobKey.jobKey(changeType, contextName);
TriggerKey triggerKey = TriggerKey.triggerKey(changeType, contextName);
JobDetail jobDetail = JobBuilder.newJob(TaskJobBean.class).withIdentity(jobKey).usingJobData(dataMap).build();
Trigger cronTrigger = TriggerBuilder.newTrigger().withIdentity(triggerKey).startNow().build();
try {
scheduler.scheduleJob(jobDetail, cronTrigger);
} catch (ObjectAlreadyExistsException e) {
throw new BadRequestException("任务已存在,跳过加载");
} catch (SchedulerException e) {
log.error("任务执行失败", e);
throw new BadRequestException(e);
}
buildAndSchedule(changeType, contextName, dataMap);
return taskInstanceInfoTb;
}
@@ -189,6 +208,9 @@ public class TaskManage {
if (historyInstance == null) {
// 仅返回模板
TaskTemplateInfoVo templateInfo = taskTemplateInfoService.getTemplateInfoByECL(envAlias, contextName, language, changeType);
if (templateInfo == null) {
throw new BadRequestException("没有查询到任务编排模板");
}
record.setTemplate(templateInfo.getTemplateInfo());
return record;
}