refactor(task): 重构任务模块代码结构

This commit is contained in:
2025-12-16 09:12:15 +08:00
parent 6e3061cf04
commit 09393e4636
76 changed files with 205 additions and 135 deletions
@@ -14,7 +14,7 @@
<dependencies>
<dependency>
<groupId>cn.odboy</groupId>
<artifactId>cutejava-module-task</artifactId>
<artifactId>cutejava-module-system</artifactId>
<version>1.4.1</version>
</dependency>
</dependencies>
-1
View File
@@ -11,7 +11,6 @@
<!-- 基础设施 -->
<module>cutejava-framework</module>
<module>cutejava-module-system</module>
<module>cutejava-module-task</module>
<module>cutejava-starter</module>
</modules>
@@ -8,8 +8,8 @@
</parent>
<modelVersion>4.0.0</modelVersion>
<artifactId>cutejava-module-task</artifactId>
<name>任务模块</name>
<artifactId>cutejava-module-task-v1</artifactId>
<name>任务模块:强制停止任务线程</name>
<dependencies>
<dependency>
@@ -27,7 +27,8 @@ import lombok.Getter;
@Getter
@AllArgsConstructor
public enum TaskStatusEnum {
Running("running", "运行中"), Success("success", "执行成功"), Fail("fail", "执行失败");
Pending("pending", "未开始"), Running("running", "运行中"), Success("success", "执行成功"),
Fail("fail", "执行失败");
private final String code;
private final String name;
@@ -21,6 +21,7 @@ import cn.odboy.framework.context.KitSpringBeanHolder;
import cn.odboy.framework.exception.BadRequestException;
import cn.odboy.framework.exception.ServerException;
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;
import cn.odboy.task.dal.model.TaskTemplateNodeVo;
@@ -153,8 +154,8 @@ public class TaskJobBean extends QuartzJobBean implements InterruptableJob {
taskInstanceDetail.setFinishTime(null);
taskInstanceDetail.setBizCode(taskTemplateNodeVo.getCode());
taskInstanceDetail.setBizName(taskTemplateNodeVo.getName());
taskInstanceDetail.setExecuteInfo("未开始");
taskInstanceDetail.setExecuteStatus("pending");
taskInstanceDetail.setExecuteInfo(TaskStatusEnum.Pending.getName());
taskInstanceDetail.setExecuteStatus(TaskStatusEnum.Pending.getCode());
return taskInstanceDetail;
}
@@ -8,8 +8,8 @@
</parent>
<modelVersion>4.0.0</modelVersion>
<artifactId>cutejava-module-task</artifactId>
<name>任务模块</name>
<artifactId>cutejava-module-task-v2</artifactId>
<name>任务模块V2:根据任务标记停止轮询检查</name>
<dependencies>
<dependency>
@@ -27,7 +27,8 @@ import lombok.Getter;
@Getter
@AllArgsConstructor
public enum TaskStatusEnum {
Running("running", "运行中"), Success("success", "执行成功"), Fail("fail", "执行失败");
Pending("pending", "未开始"), Running("running", "运行中"), Success("success", "执行成功"),
Fail("fail", "执行失败");
private final String code;
private final String name;
@@ -52,7 +52,7 @@ public class TaskTestsController {
TaskInstanceInfoTb instanceInfo =
taskManage.createJob("cutejava", TaskChangeTypeEnum.AppContainerDeploy, "java", "daily", "cutejava", "功能测试", null);
log.info("任务创建成功,实例为:{}", JSON.toJSONString(instanceInfo));
Thread.startVirtualThread(() -> {
ThreadUtil.execAsync(() -> {
ThreadUtil.safeSleep(5000);
taskManage.stopJob(instanceInfo.getId());
});
@@ -16,11 +16,13 @@
package cn.odboy.task.core;
import cn.hutool.core.thread.ThreadUtil;
import cn.hutool.core.util.StrUtil;
import cn.odboy.framework.context.KitSpringBeanHolder;
import cn.odboy.framework.exception.BadRequestException;
import cn.odboy.framework.exception.ServerException;
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;
import cn.odboy.task.dal.model.TaskTemplateNodeVo;
@@ -28,22 +30,18 @@ import cn.odboy.task.service.TaskInstanceDetailService;
import cn.odboy.task.service.TaskInstanceInfoService;
import cn.odboy.task.service.TaskInstanceStepDetailService;
import com.alibaba.fastjson2.JSON;
import lombok.extern.slf4j.Slf4j;
import org.quartz.InterruptableJob;
import org.quartz.JobDataMap;
import org.quartz.JobExecutionContext;
import org.springframework.scheduling.quartz.QuartzJobBean;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
import lombok.extern.slf4j.Slf4j;
import org.quartz.JobDataMap;
import org.quartz.JobExecutionContext;
import org.springframework.scheduling.quartz.QuartzJobBean;
@Slf4j
public class TaskJobBean extends QuartzJobBean implements InterruptableJob {
private Thread workThread = null;
public class TaskJobBean extends QuartzJobBean {
private String getBeanAlias(String code) {
if (StrUtil.isBlank(code)) {
throw new ServerException("参数code必填");
@@ -55,10 +53,10 @@ public class TaskJobBean extends QuartzJobBean implements InterruptableJob {
@Override
public void executeInternal(JobExecutionContext context) {
this.workThread = Thread.currentThread();
// ========================== 获取代理类 ==========================
TaskInstanceInfoService taskInstanceInfoService = KitSpringBeanHolder.getBean(TaskInstanceInfoService.class);
TaskInstanceDetailService taskInstanceDetailService = KitSpringBeanHolder.getBean(TaskInstanceDetailService.class);
TaskInstanceDetailService taskInstanceDetailService =
KitSpringBeanHolder.getBean(TaskInstanceDetailService.class);
// ========================== 获取参数 ==========================
JobDataMap dataMap = context.getMergedJobDataMap();
long id = dataMap.getLong(TaskJobKeys.ID);
@@ -71,7 +69,8 @@ public class TaskJobBean extends QuartzJobBean implements InterruptableJob {
if (StrUtil.isBlank(retryNodeCode)) {
executeNormalTask(taskInstanceInfoService, taskInstanceDetailService, id, dataMap, taskTemplateNodeVos);
} else {
executeRetryTask(taskInstanceInfoService, taskInstanceDetailService, id, dataMap, taskTemplateNodeVos, retryNodeCode);
executeRetryTask(taskInstanceInfoService, taskInstanceDetailService, id, dataMap, taskTemplateNodeVos,
retryNodeCode);
}
}
@@ -84,15 +83,19 @@ public class TaskJobBean extends QuartzJobBean implements InterruptableJob {
* @param dataMap 数据映射
* @param taskTemplateNodeVos 模板节点列表
*/
private void executeNormalTask(TaskInstanceInfoService taskInstanceInfoService, TaskInstanceDetailService taskInstanceDetailService, long id,
JobDataMap dataMap, List<TaskTemplateNodeVo> taskTemplateNodeVos) {
private void executeNormalTask(TaskInstanceInfoService taskInstanceInfoService,
TaskInstanceDetailService taskInstanceDetailService, long id, JobDataMap dataMap,
List<TaskTemplateNodeVo> taskTemplateNodeVos) {
// ========================== 初始化执行明细 ==========================
List<TaskInstanceDetailTb> taskInstanceDetails =
taskTemplateNodeVos.stream().map(taskTemplateNodeVo -> buildTaskInstanceDetail(id, taskTemplateNodeVo)).toList();
taskTemplateNodeVos.stream().map(taskTemplateNodeVo -> buildTaskInstanceDetail(id, taskTemplateNodeVo))
.collect(Collectors.toList());
taskInstanceDetailService.saveBatch(taskInstanceDetails);
Map<String, Long> codeIdMap = taskInstanceDetails.stream().collect(Collectors.toMap(TaskInstanceDetailTb::getBizCode, TaskInstanceDetailTb::getId));
Map<String, Long> codeIdMap = taskInstanceDetails.stream()
.collect(Collectors.toMap(TaskInstanceDetailTb::getBizCode, TaskInstanceDetailTb::getId));
// ========================== 顺序执行 ==========================
executeTaskSteps(taskInstanceInfoService, taskInstanceDetailService, id, dataMap, taskTemplateNodeVos, codeIdMap);
executeTaskSteps(taskInstanceInfoService, taskInstanceDetailService, id, dataMap, taskTemplateNodeVos,
codeIdMap);
}
/**
@@ -105,8 +108,9 @@ public class TaskJobBean extends QuartzJobBean implements InterruptableJob {
* @param taskTemplateNodeVos 模板节点列表
* @param retryNodeCode 重试节点编码
*/
private void executeRetryTask(TaskInstanceInfoService taskInstanceInfoService, TaskInstanceDetailService taskInstanceDetailService, long id,
JobDataMap dataMap, List<TaskTemplateNodeVo> taskTemplateNodeVos, String retryNodeCode) {
private void executeRetryTask(TaskInstanceInfoService taskInstanceInfoService,
TaskInstanceDetailService taskInstanceDetailService, long id, JobDataMap dataMap,
List<TaskTemplateNodeVo> taskTemplateNodeVos, String retryNodeCode) {
boolean isFound = false;
// ========================== 初始化执行明细 ==========================
List<TaskInstanceDetailTb> taskInstanceDetails = new ArrayList<>();
@@ -126,18 +130,24 @@ public class TaskJobBean extends QuartzJobBean implements InterruptableJob {
isFound = true;
}
}
TaskInstanceStepDetailService taskInstanceStepDetailService = KitSpringBeanHolder.getBean(TaskInstanceStepDetailService.class);
List<String> bizCodeList = taskInstanceDetails.stream().map(TaskInstanceDetailTb::getBizCode).distinct().toList();
List<TaskInstanceDetailTb> taskInstanceDetailTbs = taskInstanceDetailService.queryByInstanceIdAndBizCodeList(id, bizCodeList);
List<Long> taskInstanceDetailIds = taskInstanceDetailTbs.stream().map(TaskInstanceDetailTb::getId).toList();
TaskInstanceStepDetailService taskInstanceStepDetailService =
KitSpringBeanHolder.getBean(TaskInstanceStepDetailService.class);
List<String> bizCodeList =
taskInstanceDetails.stream().map(TaskInstanceDetailTb::getBizCode).distinct().collect(Collectors.toList());
List<TaskInstanceDetailTb> taskInstanceDetailTbs =
taskInstanceDetailService.queryByInstanceIdAndBizCodeList(id, bizCodeList);
List<Long> taskInstanceDetailIds =
taskInstanceDetailTbs.stream().map(TaskInstanceDetailTb::getId).collect(Collectors.toList());
// 根据ID删除明细
taskInstanceDetailService.removeByIds(taskInstanceDetailIds);
// 根据明细删除步骤
taskInstanceStepDetailService.removeByInstanceDetailIds(taskInstanceDetailIds);
taskInstanceDetailService.saveBatch(taskInstanceDetails);
Map<String, Long> codeIdMap = taskInstanceDetails.stream().collect(Collectors.toMap(TaskInstanceDetailTb::getBizCode, TaskInstanceDetailTb::getId));
Map<String, Long> codeIdMap = taskInstanceDetails.stream()
.collect(Collectors.toMap(TaskInstanceDetailTb::getBizCode, TaskInstanceDetailTb::getId));
// ========================== 顺序执行 ==========================
executeTaskSteps(taskInstanceInfoService, taskInstanceDetailService, id, dataMap, taskTemplateNodeRetryList, codeIdMap);
executeTaskSteps(taskInstanceInfoService, taskInstanceDetailService, id, dataMap, taskTemplateNodeRetryList,
codeIdMap);
}
/**
@@ -153,8 +163,8 @@ public class TaskJobBean extends QuartzJobBean implements InterruptableJob {
taskInstanceDetail.setFinishTime(null);
taskInstanceDetail.setBizCode(taskTemplateNodeVo.getCode());
taskInstanceDetail.setBizName(taskTemplateNodeVo.getName());
taskInstanceDetail.setExecuteInfo("未开始");
taskInstanceDetail.setExecuteStatus("pending");
taskInstanceDetail.setExecuteInfo(TaskStatusEnum.Pending.getName());
taskInstanceDetail.setExecuteStatus(TaskStatusEnum.Pending.getCode());
return taskInstanceDetail;
}
@@ -168,33 +178,57 @@ public class TaskJobBean extends QuartzJobBean implements InterruptableJob {
* @param taskTemplateNodes 需要执行的任务节点列表
* @param codeIdMap 节点编码与明细ID映射
*/
private void executeTaskSteps(TaskInstanceInfoService taskInstanceInfoService, TaskInstanceDetailService taskInstanceDetailService, Long id,
JobDataMap dataMap, List<TaskTemplateNodeVo> taskTemplateNodes, Map<String, Long> codeIdMap) {
private void executeTaskSteps(TaskInstanceInfoService taskInstanceInfoService,
TaskInstanceDetailService taskInstanceDetailService, Long id, JobDataMap dataMap,
List<TaskTemplateNodeVo> taskTemplateNodes, Map<String, Long> codeIdMap) {
try {
boolean isTerminate = false;
for (TaskTemplateNodeVo taskTemplateNodeVo : taskTemplateNodes) {
String code = taskTemplateNodeVo.getCode();
try {
TaskStepExecutor executor = KitSpringBeanHolder.getBean("taskStep" + getBeanAlias(code));
taskInstanceDetailService.fastStart(id, code, dataMap);
executor.execute(codeIdMap.getOrDefault(code, null), dataMap, taskTemplateNodeVo);
taskInstanceDetailService.fastSuccessWithInfo(id, code, null);
} catch (Exception e) {
log.error("任务节点执行失败", e);
taskInstanceDetailService.fastFailWithInfo(id, code, e.getMessage());
throw new ServerException(e);
while (true) {
TaskInstanceDetailTb currentTaskInstanceNode =
taskInstanceDetailService.getOneByInstanceIdAndCode(id, code);
if (TaskStatusEnum.Fail.getCode().equals(currentTaskInstanceNode.getExecuteStatus())) {
log.info("任务节点执行失败, instanceId={}, bizCode={}", id, code);
isTerminate = true;
break;
} else if (TaskStatusEnum.Success.getCode().equals(currentTaskInstanceNode.getExecuteStatus())) {
log.info("任务节点执行成功, instanceId={}, bizCode={}", id, code);
break;
} else if (TaskStatusEnum.Pending.getCode().equals(currentTaskInstanceNode.getExecuteStatus())) {
log.info("任务节点由待执行变更为执行中, instanceId={}, bizCode={}", id, code);
try {
TaskStepExecutor executor = KitSpringBeanHolder.getBean("taskStep" + getBeanAlias(code));
taskInstanceDetailService.fastStart(id, code, dataMap);
executor.execute(codeIdMap.getOrDefault(code, null), dataMap, taskTemplateNodeVo);
taskInstanceDetailService.fastSuccessWithInfo(id, code, null);
} catch (Exception e) {
log.error("任务节点执行失败", e);
taskInstanceDetailService.fastFailWithInfo(id, code, e.getMessage());
throw new ServerException(e);
}
break;
} else if (TaskStatusEnum.Running.getCode().equals(currentTaskInstanceNode.getExecuteStatus())) {
ThreadUtil.safeSleep(5000);
} else {
// 处理未知状态
log.warn("任务节点处于未知状态, instanceId={}, bizCode={}, status={}", id, code,
currentTaskInstanceNode.getExecuteStatus());
ThreadUtil.safeSleep(5000);
}
}
if (isTerminate) {
break;
}
}
taskInstanceInfoService.fastSuccessWithData(id, dataMap);
if (isTerminate) {
taskInstanceInfoService.fastFailWithMessageData(id, "任务执行失败", dataMap);
} else {
taskInstanceInfoService.fastSuccessWithData(id, dataMap);
}
} catch (Exception e) {
log.error("任务执行失败", e);
taskInstanceInfoService.fastFailWithMessageData(id, e.getMessage(), dataMap);
}
}
@Override
public void interrupt() {
if (workThread != null) {
workThread.stop();
}
}
}
@@ -33,15 +33,23 @@ import cn.odboy.task.service.TaskInstanceInfoService;
import cn.odboy.task.service.TaskTemplateInfoService;
import cn.odboy.util.KitDateUtil;
import com.alibaba.fastjson2.JSON;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.quartz.*;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.quartz.JobBuilder;
import org.quartz.JobDataMap;
import org.quartz.JobDetail;
import org.quartz.JobKey;
import org.quartz.ObjectAlreadyExistsException;
import org.quartz.Scheduler;
import org.quartz.SchedulerException;
import org.quartz.Trigger;
import org.quartz.TriggerBuilder;
import org.quartz.TriggerKey;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
@Slf4j
@Component
@@ -64,8 +72,8 @@ public class TaskManage {
* @param dataMap 任务参数
*/
@Transactional(rollbackFor = Exception.class)
public TaskInstanceInfoTb createJob(String contextName, TaskChangeTypeEnum changeTypeEnum, String language, String envAlias, String source, String reason,
JobDataMap dataMap) {
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必填");
}
@@ -85,7 +93,8 @@ public class TaskManage {
dataMap = new JobDataMap();
}
// ========================== 获取任务编排模板 ==========================
TaskTemplateInfoVo taskInstanceInfoVo = taskTemplateInfoService.getTemplateInfoByECL(envAlias, contextName, language, changeTypeEnum.getCode());
TaskTemplateInfoVo taskInstanceInfoVo =
taskTemplateInfoService.getTemplateInfoByECL(envAlias, contextName, language, changeTypeEnum.getCode());
if (taskInstanceInfoVo == null) {
throw new BadRequestException("没有查询到任务编排模板");
}
@@ -103,7 +112,8 @@ public class TaskManage {
// ========================== 创建任务 ==========================
TaskManage taskManage = KitSpringBeanHolder.getBean(TaskManage.class);
TaskInstanceInfoTb newInstance =
taskManage.saveTaskInstanceInfoTb(contextName, changeTypeEnum, language, envAlias, source, reason, dataMap, templateInfo);
taskManage.saveTaskInstanceInfoTb(contextName, changeTypeEnum, language, envAlias, source, reason, dataMap,
templateInfo);
dataMap.put(TaskJobKeys.ID, newInstance.getId());
// 应用名、资源类型
dataMap.put(TaskJobKeys.CONTEXT_NAME, contextName);
@@ -117,8 +127,8 @@ public class TaskManage {
}
@Transactional(rollbackFor = Exception.class)
public TaskInstanceInfoTb saveTaskInstanceInfoTb(String contextName, TaskChangeTypeEnum changeTypeEnum, String language, String envAlias, String source,
String reason, JobDataMap dataMap, String templateInfo) {
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);
@@ -170,9 +180,9 @@ public class TaskManage {
scheduler.pauseTrigger(triggerKey);
// 移除触发器
scheduler.unscheduleJob(triggerKey);
// 删除并中断任务
// 删除任务
scheduler.deleteJob(jobKey);
scheduler.interrupt(jobKey);
// scheduler.interrupt(jobKey);
} catch (SchedulerException e) {
log.error("任务停止失败", e);
throw new RuntimeException(e);
@@ -195,7 +205,7 @@ public class TaskManage {
if (dataMap == null) {
dataMap = new JobDataMap();
}
taskInstanceInfoTb.setStatus(TaskStatusEnum.Running.getCode());
taskInstanceInfoTb.setStatus(TaskStatusEnum.Pending.getCode());
taskInstanceInfoTb.setFinishTime(null);
taskInstanceInfoService.updateById(taskInstanceInfoTb);
// ========================== 执行任务 ==========================
@@ -209,10 +219,12 @@ public class TaskManage {
public TaskInstanceInfoVo getLastInfo(String contextName, String language, String envAlias, String changeType) {
TaskInstanceInfoVo record = new TaskInstanceInfoVo();
TaskInstanceInfoTb historyInstance = taskInstanceInfoService.getLastHistoryInstance(contextName, language, envAlias, changeType);
TaskInstanceInfoTb historyInstance =
taskInstanceInfoService.getLastHistoryInstance(contextName, language, envAlias, changeType);
if (historyInstance == null) {
// 仅返回模板
TaskTemplateInfoVo templateInfo = taskTemplateInfoService.getTemplateInfoByECL(envAlias, contextName, language, changeType);
TaskTemplateInfoVo templateInfo =
taskTemplateInfoService.getTemplateInfoByECL(envAlias, contextName, language, changeType);
if (templateInfo == null) {
throw new BadRequestException("没有查询到任务编排模板");
}
@@ -223,7 +235,8 @@ public class TaskManage {
record = BeanUtil.copyProperties(historyInstance, TaskInstanceInfoVo.class);
record.setHistory(buildNodeList(record));
TaskInstanceInfoTb runningInstance = taskInstanceInfoService.getLastRunningInstance(contextName, language, envAlias, changeType);
TaskInstanceInfoTb runningInstance =
taskInstanceInfoService.getLastRunningInstance(contextName, language, envAlias, changeType);
if (runningInstance != null) {
TaskInstanceInfoVo taskInstanceInfoVo = BeanUtil.copyProperties(runningInstance, TaskInstanceInfoVo.class);
record.setCurrent(buildNodeList(taskInstanceInfoVo));
@@ -243,9 +256,11 @@ public class TaskManage {
instanceNodeVo.setStartTime(taskInstanceDetail.getStartTime());
instanceNodeVo.setFinishTime(taskInstanceDetail.getFinishTime());
if (taskInstanceDetail.getFinishTime() == null) {
instanceNodeVo.setDurationDesc(KitDateUtil.formatSecondsDuration(taskInstanceDetail.getStartTime(), new Date()));
instanceNodeVo.setDurationDesc(
KitDateUtil.formatSecondsDuration(taskInstanceDetail.getStartTime(), new Date()));
} else {
instanceNodeVo.setDurationDesc(KitDateUtil.formatSecondsDuration(taskInstanceDetail.getStartTime(), taskInstanceDetail.getFinishTime()));
instanceNodeVo.setDurationDesc(KitDateUtil.formatSecondsDuration(taskInstanceDetail.getStartTime(),
taskInstanceDetail.getFinishTime()));
}
instanceNodeVo.setRunningDesc(taskInstanceDetail.getExecuteInfo());
instanceNodeVo.setStatus(taskInstanceDetail.getExecuteStatus());
@@ -20,7 +20,8 @@ import com.alibaba.fastjson2.annotation.JSONField;
import com.baomidou.mybatisplus.annotation.*;
import com.fasterxml.jackson.databind.annotation.JsonSerialize;
import com.fasterxml.jackson.databind.ser.std.ToStringSerializer;
import io.swagger.v3.oas.annotations.media.Schema;
import io.swagger.annotations.ApiModel;
import io.swagger.annotations.ApiModelProperty;
import lombok.Getter;
import lombok.Setter;
import lombok.ToString;
@@ -39,16 +40,16 @@ import java.util.Date;
@Setter
@ToString
@TableName("task_instance_detail")
@Schema(name = "TaskInstanceDetailTb对象", description = "任务实例明细")
@ApiModel(value = "TaskInstanceDetailTb对象", description = "任务实例明细")
public class TaskInstanceDetailTb extends KitObject {
@TableField(fill = FieldFill.INSERT)
@Schema(name = "创建时间: yyyy-MM-dd HH:mm:ss", hidden = true)
@ApiModelProperty(value = "创建时间: yyyy-MM-dd HH:mm:ss", hidden = true)
private Date createTime;
/**
* id
*/
@Schema(name = "id")
@ApiModelProperty("id")
@TableId(value = "id", type = IdType.ASSIGN_ID)
@JsonSerialize(using = ToStringSerializer.class)
@JSONField(format = "string")
@@ -58,20 +59,20 @@ public class TaskInstanceDetailTb extends KitObject {
* 任务实例id
*/
@TableField("instance_id")
@Schema(name = "任务实例id")
@ApiModelProperty("任务实例id")
private Long instanceId;
/**
* 开始时间
*/
@Schema(name = "开始时间")
@ApiModelProperty("开始时间")
@TableField("start_time")
private Date startTime;
/**
* 结束时间
*/
@Schema(name = "结束时间")
@ApiModelProperty("结束时间")
@TableField("finish_time")
private Date finishTime;
@@ -79,27 +80,27 @@ public class TaskInstanceDetailTb extends KitObject {
* 业务编码
*/
@TableField("biz_code")
@Schema(name = "业务编码")
@ApiModelProperty("业务编码")
private String bizCode;
/**
* 业务名称(步骤)
*/
@TableField("biz_name")
@Schema(name = "业务名称(步骤)")
@ApiModelProperty("业务名称(步骤)")
private String bizName;
/**
* 执行参数
*/
@Schema(name = "执行参数")
@ApiModelProperty("执行参数")
@TableField("execute_params")
private String executeParams;
/**
* 执行信息
*/
@Schema(name = "执行信息")
@ApiModelProperty("执行信息")
@TableField("execute_info")
private String executeInfo;
@@ -107,6 +108,6 @@ public class TaskInstanceDetailTb extends KitObject {
* 执行状态(running进行中 success成功 fail失败)
*/
@TableField("execute_status")
@Schema(name = "执行状态(running进行中 success成功 fail失败)")
@ApiModelProperty("执行状态(running进行中 success成功 fail失败)")
private String executeStatus;
}
@@ -20,7 +20,8 @@ import com.alibaba.fastjson2.annotation.JSONField;
import com.baomidou.mybatisplus.annotation.*;
import com.fasterxml.jackson.databind.annotation.JsonSerialize;
import com.fasterxml.jackson.databind.ser.std.ToStringSerializer;
import io.swagger.v3.oas.annotations.media.Schema;
import io.swagger.annotations.ApiModel;
import io.swagger.annotations.ApiModelProperty;
import lombok.Getter;
import lombok.Setter;
import lombok.ToString;
@@ -40,25 +41,25 @@ import java.util.Date;
@Setter
@ToString
@TableName("task_instance_info")
@Schema(name = "TaskInstanceInfoTb对象", description = "任务实例")
@ApiModel(value = "TaskInstanceInfoTb对象", description = "任务实例")
public class TaskInstanceInfoTb extends KitObject {
@CreatedBy
@TableField(fill = FieldFill.INSERT)
@Schema(name = "创建人", hidden = true)
@ApiModelProperty(value = "创建人", hidden = true)
private String createBy;
@TableField(fill = FieldFill.INSERT)
@Schema(name = "创建时间: yyyy-MM-dd HH:mm:ss", hidden = true)
@ApiModelProperty(value = "创建时间: yyyy-MM-dd HH:mm:ss", hidden = true)
private Date createTime;
@TableField(fill = FieldFill.INSERT_UPDATE)
@Schema(name = "更新时间: yyyy-MM-dd HH:mm:ss", hidden = true)
@ApiModelProperty(value = "更新时间: yyyy-MM-dd HH:mm:ss", hidden = true)
private Date updateTime;
/**
* id, @JsonSerialize和@JSONField,用于处理大数精度丢失问题
*/
@Schema(name = "id")
@ApiModelProperty("id")
@TableId(value = "id", type = IdType.ASSIGN_ID)
@JsonSerialize(using = ToStringSerializer.class)
@JSONField(format = "string")
@@ -67,21 +68,21 @@ public class TaskInstanceInfoTb extends KitObject {
/**
* 名称
*/
@Schema(name = "名称")
@ApiModelProperty("名称")
@TableField("context_name")
private String contextName;
/**
* 语言
*/
@Schema(name = "语言")
@ApiModelProperty("语言")
@TableField("`language`")
private String language;
/**
* 变更类型
*/
@Schema(name = "变更类型")
@ApiModelProperty("变更类型")
@TableField("change_type")
private String changeType;
@@ -89,27 +90,27 @@ public class TaskInstanceInfoTb extends KitObject {
* 环境别名
*/
@TableField("env_alias")
@Schema(name = "环境别名")
@ApiModelProperty("环境别名")
private String envAlias;
/**
* 状态(running进行中 success成功 fail失败)
*/
@TableField("`status`")
@Schema(name = "状态(running进行中 success成功 fail失败)")
@ApiModelProperty("状态(running进行中 success成功 fail失败)")
private String status;
/**
* 完成时间
*/
@Schema(name = "完成时间")
@ApiModelProperty("完成时间")
@TableField("finish_time")
private Date finishTime;
/**
* 来源
*/
@Schema(name = "来源")
@ApiModelProperty("来源")
@TableField("`source`")
private String source;
@@ -117,27 +118,27 @@ public class TaskInstanceInfoTb extends KitObject {
* 变更原因
*/
@TableField("reason")
@Schema(name = "变更原因")
@ApiModelProperty("变更原因")
private String reason;
/**
* 任务模板
*/
@TableField("template")
@Schema(name = "template")
@ApiModelProperty("template")
private String template;
/**
* QuartzJob参数
*/
@TableField("job_data")
@Schema(name = "QuartzJob参数")
@ApiModelProperty("QuartzJob参数")
private String jobData;
/**
* 异常信息
*/
@TableField("error_message")
@Schema(name = "异常信息")
@ApiModelProperty("异常信息")
private String errorMessage;
}
@@ -20,7 +20,8 @@ import com.alibaba.fastjson2.annotation.JSONField;
import com.baomidou.mybatisplus.annotation.*;
import com.fasterxml.jackson.databind.annotation.JsonSerialize;
import com.fasterxml.jackson.databind.ser.std.ToStringSerializer;
import io.swagger.v3.oas.annotations.media.Schema;
import io.swagger.annotations.ApiModel;
import io.swagger.annotations.ApiModelProperty;
import lombok.Getter;
import lombok.Setter;
import lombok.ToString;
@@ -39,16 +40,16 @@ import java.util.Date;
@Setter
@ToString
@TableName("task_instance_step_detail")
@Schema(name = "TaskInstanceStepDetailTb对象")
@ApiModel(value = "TaskInstanceStepDetailTb对象", description = "")
public class TaskInstanceStepDetailTb extends KitObject {
@TableField(fill = FieldFill.INSERT)
@Schema(name = "创建时间: yyyy-MM-dd HH:mm:ss", hidden = true)
@ApiModelProperty(value = "创建时间: yyyy-MM-dd HH:mm:ss", hidden = true)
private Date createTime;
/**
* id
*/
@Schema(name = "id")
@ApiModelProperty("id")
@TableId(value = "id", type = IdType.ASSIGN_ID)
@JsonSerialize(using = ToStringSerializer.class)
@JSONField(format = "string")
@@ -58,20 +59,20 @@ public class TaskInstanceStepDetailTb extends KitObject {
* task_instance_detail表id
*/
@TableField("instance_detail_id")
@Schema(name = "task_instance_detail表id")
@ApiModelProperty("task_instance_detail表id")
private Long instanceDetailId;
/**
* 步骤描述
*/
@TableField("step_desc")
@Schema(name = "步骤描述")
@ApiModelProperty("步骤描述")
private String stepDesc;
/**
* 状态(success成功 fail失败)
*/
@TableField("step_status")
@Schema(name = "状态(success成功 fail失败)")
@ApiModelProperty("状态(success成功 fail失败)")
private String stepStatus;
}
@@ -20,7 +20,8 @@ import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableField;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import io.swagger.v3.oas.annotations.media.Schema;
import io.swagger.annotations.ApiModel;
import io.swagger.annotations.ApiModelProperty;
import lombok.Getter;
import lombok.Setter;
import lombok.ToString;
@@ -37,13 +38,13 @@ import lombok.ToString;
@Setter
@ToString
@TableName("task_template_info")
@Schema(name = "TaskTemplateInfoTb对象", description = "任务模板")
@ApiModel(value = "TaskTemplateInfoTb对象", description = "任务模板")
public class TaskTemplateInfoTb extends KitBaseUserTimeTb {
/**
* id
*/
@Schema(name = "id")
@ApiModelProperty("id")
@TableId(value = "id", type = IdType.AUTO)
private Long id;
@@ -51,13 +52,13 @@ public class TaskTemplateInfoTb extends KitBaseUserTimeTb {
* 数据有效性
*/
@TableField("available")
@Schema(name = "数据有效性")
@ApiModelProperty("数据有效性")
private Boolean available;
/**
* 变更类型
*/
@Schema(name = "变更类型")
@ApiModelProperty("变更类型")
@TableField("change_type")
private String changeType;
@@ -65,13 +66,13 @@ public class TaskTemplateInfoTb extends KitBaseUserTimeTb {
* 流水线名称
*/
@TableField("`name`")
@Schema(name = "流水线名称")
@ApiModelProperty("流水线名称")
private String name;
/**
* 流水线描述
*/
@Schema(name = "流水线描述")
@ApiModelProperty("流水线描述")
@TableField("`description`")
private String description;
@@ -79,20 +80,20 @@ public class TaskTemplateInfoTb extends KitBaseUserTimeTb {
* 流水线模板内容
*/
@TableField("template")
@Schema(name = "流水线模板内容")
@ApiModelProperty("流水线模板内容")
private String template;
/**
* 环境别名
*/
@TableField("env_alias")
@Schema(name = "环境别名")
@ApiModelProperty("环境别名")
private String envAlias;
/**
* 应用为语言,资源为类型
*/
@TableField("`language`")
@Schema(name = "应用为语言,资源为类型")
@ApiModelProperty("应用为语言,资源为类型")
private String language;
}
@@ -30,6 +30,8 @@ import java.util.List;
* @since 2025-09-26
*/
public interface TaskInstanceDetailService extends IService<TaskInstanceDetailTb> {
TaskInstanceDetailTb getOneByInstanceIdAndCode(Long id, String code);
void fastFailWithInfo(Long instanceId, String code, String message);
void fastSuccess(Long instanceId, String code);
@@ -22,11 +22,10 @@ import cn.odboy.task.dal.mysql.TaskInstanceDetailMapper;
import cn.odboy.task.service.TaskInstanceDetailService;
import com.alibaba.fastjson2.JSON;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import org.quartz.JobDataMap;
import org.springframework.stereotype.Service;
import java.util.Date;
import java.util.List;
import org.quartz.JobDataMap;
import org.springframework.stereotype.Service;
/**
* <p>
@@ -37,7 +36,14 @@ import java.util.List;
* @since 2025-09-26
*/
@Service
public class TaskInstanceDetailServiceImpl extends ServiceImpl<TaskInstanceDetailMapper, TaskInstanceDetailTb> implements TaskInstanceDetailService {
public class TaskInstanceDetailServiceImpl extends ServiceImpl<TaskInstanceDetailMapper, TaskInstanceDetailTb>
implements TaskInstanceDetailService {
@Override
public TaskInstanceDetailTb getOneByInstanceIdAndCode(Long id, String code) {
return lambdaQuery().eq(TaskInstanceDetailTb::getInstanceId, id).eq(TaskInstanceDetailTb::getBizCode, code)
.one();
}
@Override
public void fastFailWithInfo(Long instanceId, String code, String executeInfo) {
@@ -45,7 +51,8 @@ public class TaskInstanceDetailServiceImpl extends ServiceImpl<TaskInstanceDetai
updRecord.setFinishTime(new Date());
updRecord.setExecuteInfo(executeInfo);
updRecord.setExecuteStatus(TaskStatusEnum.Fail.getCode());
lambdaUpdate().eq(TaskInstanceDetailTb::getInstanceId, instanceId).eq(TaskInstanceDetailTb::getBizCode, code).update(updRecord);
lambdaUpdate().eq(TaskInstanceDetailTb::getInstanceId, instanceId).eq(TaskInstanceDetailTb::getBizCode, code)
.update(updRecord);
}
@Override
@@ -54,7 +61,8 @@ public class TaskInstanceDetailServiceImpl extends ServiceImpl<TaskInstanceDetai
updRecord.setFinishTime(new Date());
updRecord.setExecuteInfo("执行成功");
updRecord.setExecuteStatus(TaskStatusEnum.Success.getCode());
lambdaUpdate().eq(TaskInstanceDetailTb::getInstanceId, instanceId).eq(TaskInstanceDetailTb::getBizCode, code).update(updRecord);
lambdaUpdate().eq(TaskInstanceDetailTb::getInstanceId, instanceId).eq(TaskInstanceDetailTb::getBizCode, code)
.update(updRecord);
}
@Override
@@ -63,7 +71,8 @@ public class TaskInstanceDetailServiceImpl extends ServiceImpl<TaskInstanceDetai
updRecord.setFinishTime(new Date());
updRecord.setExecuteInfo(executeInfo == null ? "执行成功" : executeInfo);
updRecord.setExecuteStatus(TaskStatusEnum.Success.getCode());
lambdaUpdate().eq(TaskInstanceDetailTb::getInstanceId, instanceId).eq(TaskInstanceDetailTb::getBizCode, code).update(updRecord);
lambdaUpdate().eq(TaskInstanceDetailTb::getInstanceId, instanceId).eq(TaskInstanceDetailTb::getBizCode, code)
.update(updRecord);
}
@Override
@@ -73,7 +82,8 @@ public class TaskInstanceDetailServiceImpl extends ServiceImpl<TaskInstanceDetai
updRecord.setExecuteInfo("运行中");
updRecord.setExecuteParams(JSON.toJSONString(dataMap));
updRecord.setExecuteStatus(TaskStatusEnum.Running.getCode());
lambdaUpdate().eq(TaskInstanceDetailTb::getInstanceId, instanceId).eq(TaskInstanceDetailTb::getBizCode, code).update(updRecord);
lambdaUpdate().eq(TaskInstanceDetailTb::getInstanceId, instanceId).eq(TaskInstanceDetailTb::getBizCode, code)
.update(updRecord);
}
@Override
@@ -81,11 +91,13 @@ public class TaskInstanceDetailServiceImpl extends ServiceImpl<TaskInstanceDetai
if (CollUtil.isEmpty(bizCodeList)) {
return CollUtil.newArrayList();
}
return lambdaQuery().eq(TaskInstanceDetailTb::getInstanceId, instanceId).in(TaskInstanceDetailTb::getBizCode, bizCodeList).list();
return lambdaQuery().eq(TaskInstanceDetailTb::getInstanceId, instanceId)
.in(TaskInstanceDetailTb::getBizCode, bizCodeList).list();
}
@Override
public List<TaskInstanceDetailTb> queryByInstanceId(Long instanceId) {
return lambdaQuery().eq(TaskInstanceDetailTb::getInstanceId, instanceId).orderByAsc(TaskInstanceDetailTb::getId).list();
return lambdaQuery().eq(TaskInstanceDetailTb::getInstanceId, instanceId).orderByAsc(TaskInstanceDetailTb::getId)
.list();
}
}
+1 -1
View File
@@ -14,7 +14,7 @@
<dependencies>
<dependency>
<groupId>cn.odboy</groupId>
<artifactId>cutejava-module-task</artifactId>
<artifactId>cutejava-module-task-v1</artifactId>
<version>1.4.1</version>
</dependency>
</dependencies>
+2 -1
View File
@@ -11,7 +11,8 @@
<!-- 基础设施 -->
<module>cutejava-framework</module>
<module>cutejava-module-system</module>
<module>cutejava-module-task</module>
<module>cutejava-module-task-v1</module>
<module>cutejava-module-task-v2</module>
<module>cutejava-starter</module>
</modules>