Commit 2a9f1517 by guo_minglei@163.com

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

parents 05ca36fe 70be1540
...@@ -21,6 +21,14 @@ public interface JobTaskMapper { ...@@ -21,6 +21,14 @@ public interface JobTaskMapper {
List<JobTask> findJobTaskByTriggerNextTimeLessThanEqual(@Param("triggerNextTime") Long triggerNextTime); List<JobTask> findJobTaskByTriggerNextTimeLessThanEqual(@Param("triggerNextTime") Long triggerNextTime);
/** /**
*
* @param runId
* @param flowId
* @return
*/
List<JobTask> findAllByRunId(@Param("runId")String runId,@Param("flowId")Integer flowId);
/**
* 根据id查询 * 根据id查询
* @param id * @param id
* @return * @return
...@@ -58,7 +66,11 @@ public interface JobTaskMapper { ...@@ -58,7 +66,11 @@ public interface JobTaskMapper {
* 根据id删除多个 * 根据id删除多个
* @param jobTaskSchedules * @param jobTaskSchedules
*/ */
void deleteInId(@Param("jobTaskSchedules") List<JobTaskSchedule> jobTaskSchedules); void deleteInId(@Param("jobTaskSchedules") List<JobTaskSchedule> jobTaskSchedules);/**
* 根据id删除多个
* @param ids
*/
void deleteInIds(@Param("ids") List<Integer> ids);
/** /**
* 根据runid删除运行信息 * 根据runid删除运行信息
......
...@@ -2,6 +2,7 @@ package com.byit.service; ...@@ -2,6 +2,7 @@ package com.byit.service;
import com.byit.model.JobTask; import com.byit.model.JobTask;
import com.byit.model.JobTaskSchedule; import com.byit.model.JobTaskSchedule;
import org.apache.ibatis.annotations.Param;
import java.util.List; import java.util.List;
...@@ -20,6 +21,13 @@ public interface JobTaskService { ...@@ -20,6 +21,13 @@ public interface JobTaskService {
List<JobTask> findJobTaskByTriggerNextTimeLessThanEqual(long maxNextTime); List<JobTask> findJobTaskByTriggerNextTimeLessThanEqual(long maxNextTime);
/** /**
* 查询一批节点
* @param runId
* @return
*/
List<JobTask> findJobTaskByRunId(String runId,Integer flowId);
/**
* 添加单个任务节点 * 添加单个任务节点
* @param jobTask * @param jobTask
*/ */
...@@ -43,4 +51,6 @@ public interface JobTaskService { ...@@ -43,4 +51,6 @@ public interface JobTaskService {
* @param jobTasks * @param jobTasks
*/ */
void removeMythJobTaskInIds(List<JobTaskSchedule> jobTasks); void removeMythJobTaskInIds(List<JobTaskSchedule> jobTasks);
void deleteInIds(@Param("ids") List<Integer> ids);
} }
...@@ -38,6 +38,11 @@ public class JobTaskServiceImpl implements JobTaskService { ...@@ -38,6 +38,11 @@ public class JobTaskServiceImpl implements JobTaskService {
return jobTaskMapper.findJobTaskByTriggerNextTimeLessThanEqual(maxNextTime); return jobTaskMapper.findJobTaskByTriggerNextTimeLessThanEqual(maxNextTime);
} }
@Override
public List<JobTask> findJobTaskByRunId(String runId,Integer flowId) {
return jobTaskMapper.findAllByRunId(runId,flowId);
}
/** /**
* 添加一个任务 * 添加一个任务
* @param jobTask 任务实体 * @param jobTask 任务实体
...@@ -69,4 +74,9 @@ public class JobTaskServiceImpl implements JobTaskService { ...@@ -69,4 +74,9 @@ public class JobTaskServiceImpl implements JobTaskService {
public void removeMythJobTaskInIds(List<JobTaskSchedule> jobTaskSchedules) { public void removeMythJobTaskInIds(List<JobTaskSchedule> jobTaskSchedules) {
jobTaskMapper.deleteInId(jobTaskSchedules); jobTaskMapper.deleteInId(jobTaskSchedules);
} }
@Override
public void deleteInIds(List<Integer> ids) {
jobTaskMapper.deleteInIds(ids);
}
} }
package com.byit.service.mapservice; package com.byit.service.mapservice;
import com.byit.model.JobTaskRunLog; import com.byit.model.JobTaskRunLog;
import com.byit.model.JobTaskSchedule;
import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.annotation.Transactional;
/** /**
...@@ -13,4 +14,10 @@ public interface JobTaskRunLogAndJobTaskService { ...@@ -13,4 +14,10 @@ public interface JobTaskRunLogAndJobTaskService {
* @param jobTaskRunLog * @param jobTaskRunLog
*/ */
void updateLogAndSaveJobTask(JobTaskRunLog jobTaskRunLog); void updateLogAndSaveJobTask(JobTaskRunLog jobTaskRunLog);
/**
* 删除task节点 执行节点的快速失败
* @param jobTaskSchedule
*/
void removeTaskNodeAndSaveRunLog(JobTaskSchedule jobTaskSchedule);
} }
package com.byit.service.mapservice.impl; package com.byit.service.mapservice.impl;
import com.byit.model.JobTask; import cn.hutool.core.collection.CollectionUtil;
import com.byit.model.JobTaskRunLog; import com.byit.dto.executor.DispatchResponseDto;
import com.byit.model.Node; import com.byit.enums.EmailEnum;
import com.byit.enums.JobResultEnum;
import com.byit.enums.NodeRunStatusPropertyEnum;
import com.byit.model.*;
import com.byit.service.impl.JobTaskRunLogServiceImpl;
import com.byit.service.mapservice.JobTaskRunLogAndJobTaskService; import com.byit.service.mapservice.JobTaskRunLogAndJobTaskService;
import com.byit.service.JobTaskRunLogService; import com.byit.service.JobTaskRunLogService;
import com.byit.service.JobTaskService; import com.byit.service.JobTaskService;
import com.byit.service.NodeService; import com.byit.service.NodeService;
import com.byit.util.SpringUtil;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.BeanUtils; import org.springframework.beans.BeanUtils;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.annotation.Transactional;
import java.util.Date;
import java.util.List;
import java.util.stream.Collectors;
/** /**
* @author huangfu * @author huangfu
* 日志表和任务表的事务映射 * 日志表和任务表的事务映射
...@@ -63,4 +72,50 @@ public class JobTaskRunLogAndJobTaskServiceImpl implements JobTaskRunLogAndJobTa ...@@ -63,4 +72,50 @@ public class JobTaskRunLogAndJobTaskServiceImpl implements JobTaskRunLogAndJobTa
jobTask.setFlowName(jobTaskRunLog.getFlowName()); jobTask.setFlowName(jobTaskRunLog.getFlowName());
jobTaskService.addMythJobTask(jobTask); jobTaskService.addMythJobTask(jobTask);
} }
@Override
public void removeTaskNodeAndSaveRunLog(JobTaskSchedule jobTaskSchedule) {
//保存当前节点为
saveErrorLog(jobTaskSchedule);
}
private void saveErrorLog(JobTaskSchedule mythJobTaskSchedule){
log.debug("-----------saveLog--保存脚本调度日志开始------------");
JobTaskRunLogServiceImpl jobTaskRunLogService = SpringUtil.getBean(JobTaskRunLogServiceImpl.class);
JobTaskRunLogWithBLOBs jobTaskRunLogById = jobTaskRunLogService.findJobTaskRunLogById(mythJobTaskSchedule.getLogId());
JobTaskRunLogWithBLOBs jobTaskRunLog = new JobTaskRunLogWithBLOBs();
//放置调度记录
if (jobTaskRunLogById.getRunCount()>1) {
//第一次调度的日志
String triggerMsg = jobTaskRunLogById.getTriggerMsg();
jobTaskRunLog.setTriggerMsg(triggerMsg+"| 节点被kill");
}else{
jobTaskRunLog.setTriggerMsg("节点被kill");
}
jobTaskRunLog.setLogId(mythJobTaskSchedule.getLogId());
jobTaskRunLog.setVersionName(mythJobTaskSchedule.getVersionName());
jobTaskRunLog.setRunType("1");
jobTaskRunLog.setJobType(mythJobTaskSchedule.getJobType());
jobTaskRunLog.setHandlerName(mythJobTaskSchedule.getHandlerName());
jobTaskRunLog.setTriggerTime(new Date());
jobTaskRunLog.setTriggerCode("2");
Date thisTime = new Date();
jobTaskRunLog.setStartTime(thisTime);
jobTaskRunLog.setEndTime(thisTime);
jobTaskRunLog.setRunCode(NodeRunStatusPropertyEnum.RUN_FAILURE.getCode());
jobTaskRunLog.setAlertEnd(EmailEnum.IS_ALARM_NO.getCode());
if (jobTaskRunLogById.getRunCount()>1) {
//上一次的执行日志
String runMsg = jobTaskRunLogById.getRunMsg();
jobTaskRunLog.setRunMsg(runMsg+"|任务节点被kill!");
}
jobTaskRunLogService.updateJobTaskRunLogWithBLOBs(jobTaskRunLog);
log.info("-----------saveLog--保存脚本调度日志结束------------");
}
} }
...@@ -46,7 +46,8 @@ public class ScriptExecutorJobTask implements TimerTask { ...@@ -46,7 +46,8 @@ public class ScriptExecutorJobTask implements TimerTask {
public void run(Timeout timeout) throws Exception { public void run(Timeout timeout) throws Exception {
log.debug("---------开始交验工作流时否正在运行中------------"); log.debug("---------开始交验工作流时否正在运行中------------");
if(checkFlowStatusIsKill(mythJobTaskSchedule.getFlowId(),mythJobTaskSchedule.getRunId())){ if(checkFlowStatusIsKill(mythJobTaskSchedule.getFlowId(),mythJobTaskSchedule.getRunId())){
log.warn("--------------该工作流已经被杀死,不执行-------------------"); log.warn("--------------该工作流已经被杀死,执行快速失败!-------------------");
saveErrorLog(mythJobTaskSchedule);
return; return;
} }
log.debug("-----------------工作流校验完成-------------"); log.debug("-----------------工作流校验完成-------------");
...@@ -158,4 +159,42 @@ public class ScriptExecutorJobTask implements TimerTask { ...@@ -158,4 +159,42 @@ public class ScriptExecutorJobTask implements TimerTask {
jobTaskRunLogService.updateJobTaskRunLogWithBLOBs(jobTaskRunLog); jobTaskRunLogService.updateJobTaskRunLogWithBLOBs(jobTaskRunLog);
log.info("-----------saveLog--保存脚本调度日志结束------------"); log.info("-----------saveLog--保存脚本调度日志结束------------");
} }
private void saveErrorLog(JobTaskSchedule mythJobTaskSchedule){
log.debug("-----------saveLog--保存脚本调度日志开始------------");
JobTaskRunLogServiceImpl jobTaskRunLogService = SpringUtil.getBean(JobTaskRunLogServiceImpl.class);
JobTaskRunLogWithBLOBs jobTaskRunLogById = jobTaskRunLogService.findJobTaskRunLogById(mythJobTaskSchedule.getLogId());
JobTaskRunLogWithBLOBs jobTaskRunLog = new JobTaskRunLogWithBLOBs();
//放置调度记录
if (jobTaskRunLogById.getRunCount()>1) {
//第一次调度的日志
String triggerMsg = jobTaskRunLogById.getTriggerMsg();
jobTaskRunLog.setTriggerMsg(triggerMsg+"| 节点被kill");
}else{
jobTaskRunLog.setTriggerMsg("节点被kill");
}
jobTaskRunLog.setLogId(mythJobTaskSchedule.getLogId());
jobTaskRunLog.setVersionName(mythJobTaskSchedule.getVersionName());
jobTaskRunLog.setRunType("1");
jobTaskRunLog.setJobType(mythJobTaskSchedule.getJobType());
jobTaskRunLog.setHandlerName(mythJobTaskSchedule.getHandlerName());
jobTaskRunLog.setTriggerTime(new Date());
jobTaskRunLog.setTriggerCode("2");
Date thisTime = new Date();
jobTaskRunLog.setStartTime(thisTime);
jobTaskRunLog.setEndTime(thisTime);
jobTaskRunLog.setRunCode(NodeRunStatusPropertyEnum.RUN_FAILURE.getCode());
jobTaskRunLog.setAlertEnd(EmailEnum.IS_ALARM_NO.getCode());
if (jobTaskRunLogById.getRunCount()>1) {
//上一次的执行日志
String runMsg = jobTaskRunLogById.getRunMsg();
jobTaskRunLog.setRunMsg(runMsg+"|任务节点被kill!");
}
jobTaskRunLogService.updateJobTaskRunLogWithBLOBs(jobTaskRunLog);
log.info("-----------saveLog--保存脚本调度日志结束------------");
}
} }
...@@ -63,6 +63,15 @@ ...@@ -63,6 +63,15 @@
from job_task from job_task
where trigger_time <![CDATA[ <= ]]> #{triggerNextTime,jdbcType=BIGINT} and trigger_status != '0' where trigger_time <![CDATA[ <= ]]> #{triggerNextTime,jdbcType=BIGINT} and trigger_status != '0'
</select> </select>
<!--根据RunId查询一批节点-->
<select id="findAllByRunId" resultMap="BaseResultMap">
select
<include refid="Base_Column_List" />
,
<include refid="Blob_Column_List" />
from job_task
where run_id=#{runId,jdbcType=VARCHAR} and flow_Id = #{flowId,jdbcType = INTEGER}
</select>
<!--根据id查询--> <!--根据id查询-->
<select id="findJobTaskById" parameterType="java.lang.Integer" resultMap="ResultMapWithBLOBs"> <select id="findJobTaskById" parameterType="java.lang.Integer" resultMap="ResultMapWithBLOBs">
select select
...@@ -449,6 +458,14 @@ ...@@ -449,6 +458,14 @@
delete from job_task delete from job_task
where id = #{id,jdbcType=INTEGER} where id = #{id,jdbcType=INTEGER}
</delete> </delete>
<!--多个删除-->
<delete id="deleteInIds" parameterType="com.byit.model.JobTaskSchedule">
delete from job_task
where id in
<foreach collection="jobTaskSchedules" item="jobTaskSchedule" open="(" separator="," close=")">
#{id,jdbcType=INTEGER}
</foreach>
</delete>
<delete id="deleteByRunId" parameterType="java.lang.Integer"> <delete id="deleteByRunId" parameterType="java.lang.Integer">
delete from job_task delete from job_task
......
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