Commit fa446771 by huangfusuper

错误日志上传

parent d114f545
...@@ -6,6 +6,8 @@ import com.byit.dto.executor.RunParamWrapped; ...@@ -6,6 +6,8 @@ import com.byit.dto.executor.RunParamWrapped;
import com.byit.dto.plugin.RunLog; import com.byit.dto.plugin.RunLog;
import com.byit.enums.NodePropertyEnum; import com.byit.enums.NodePropertyEnum;
import com.byit.enums.task.RunResultEnum; import com.byit.enums.task.RunResultEnum;
import com.byit.filesystem.FastDfsFileSystem;
import com.byit.filesystem.FileSystem;
import com.byit.job.utils.MythLogUtils; import com.byit.job.utils.MythLogUtils;
import com.byit.model.JobTaskRunLogWithBLOBs; import com.byit.model.JobTaskRunLogWithBLOBs;
import com.byit.model.JobTaskSchedule; import com.byit.model.JobTaskSchedule;
...@@ -23,7 +25,9 @@ import io.netty.util.TimerTask; ...@@ -23,7 +25,9 @@ import io.netty.util.TimerTask;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.data.redis.core.StringRedisTemplate;
import java.nio.charset.StandardCharsets;
import java.util.Date; import java.util.Date;
import java.util.HashMap;
import java.util.List; import java.util.List;
import java.util.TreeSet; import java.util.TreeSet;
import java.util.stream.Collectors; import java.util.stream.Collectors;
...@@ -89,9 +93,9 @@ public class JavaNodeExecutorTask implements TimerTask { ...@@ -89,9 +93,9 @@ public class JavaNodeExecutorTask implements TimerTask {
GetRegConfig getRegConfig = SpringUtil.getBean(GetRegConfig.class); GetRegConfig getRegConfig = SpringUtil.getBean(GetRegConfig.class);
TreeSet<String> gatewayConf = getRegConfig.getGatewayConf(); TreeSet<String> gatewayConf = getRegConfig.getGatewayConf();
List<String> callUrlList = gatewayConf.stream().map(callbackIp -> { List<String> callUrlList = gatewayConf.stream().map(callbackIp -> {
if("2".equals(hasRelatedFlow)){ if ("2".equals(hasRelatedFlow)) {
return String.format("%s%s%s", HTTP_PRE, callbackIp, HTTP_SUFFIX_FLOW); return String.format("%s%s%s", HTTP_PRE, callbackIp, HTTP_SUFFIX_FLOW);
}else{ } else {
return String.format("%s%s%s", HTTP_PRE, callbackIp, HTTP_SUFFIX); return String.format("%s%s%s", HTTP_PRE, callbackIp, HTTP_SUFFIX);
} }
}).collect(Collectors.toList()); }).collect(Collectors.toList());
...@@ -129,20 +133,34 @@ public class JavaNodeExecutorTask implements TimerTask { ...@@ -129,20 +133,34 @@ public class JavaNodeExecutorTask implements TimerTask {
jobTaskRunLogById.setAlertEnd("1"); jobTaskRunLogById.setAlertEnd("1");
} catch (Exception e) { } catch (Exception e) {
String runErrorMessage = MythLogUtils.getMessage(e);
if (mythJobTaskSchedule.getScheduleType() == 4) { if (mythJobTaskSchedule.getScheduleType() == 4) {
StringRedisTemplate stringRedisTemplate = (StringRedisTemplate) SpringUtil.getBean("stringRedisTemplate"); StringRedisTemplate stringRedisTemplate = (StringRedisTemplate) SpringUtil.getBean("stringRedisTemplate");
RunLog runLog = RunLog.builder().isEnd(true).isSuccess(false).runLog("执行资源异常" + MythLogUtils.getMessage(e)).build(); RunLog runLog = RunLog.builder().isEnd(true).isSuccess(false).runLog("执行资源异常" + runErrorMessage).build();
//stringRedisTemplate.convertAndSend("REAL:EXEC:" + mythJobTaskSchedule.getLogId() , JSON.toJSONString(runLog, WriteClassName)); //stringRedisTemplate.convertAndSend("REAL:EXEC:" + mythJobTaskSchedule.getLogId() , JSON.toJSONString(runLog, WriteClassName));
//使用redis 想队尾push一个值 //使用redis 想队尾push一个值
stringRedisTemplate.opsForList().rightPush("REAL:EXEC:" + mythJobTaskSchedule.getLogId(), JSON.toJSONString(runLog, WriteClassName)); stringRedisTemplate.opsForList().rightPush("REAL:EXEC:" + mythJobTaskSchedule.getLogId(), JSON.toJSONString(runLog, WriteClassName));
} }
log.error("------执行机异常{}-----", MythLogUtils.getMessage(e)); log.error("------执行机异常{}-----", runErrorMessage);
assert jobTaskRunLogById != null; assert jobTaskRunLogById != null;
jobTaskRunLogById.setTriggerCode(RunResultEnum.TRIGGER_ERROR.getCode()); jobTaskRunLogById.setTriggerCode(RunResultEnum.TRIGGER_ERROR.getCode());
jobTaskRunLogById.setRunCode(RunResultEnum.RUN_ERROR.getCode()); jobTaskRunLogById.setRunCode(RunResultEnum.RUN_ERROR.getCode());
jobTaskRunLogById.setRunMsg(mythJobTaskSchedule.getNodeName() + ":" + MythLogUtils.getMessage(e)); jobTaskRunLogById.setRunMsg(mythJobTaskSchedule.getNodeName() + ":" + runErrorMessage);
jobTaskRunLogById.setTriggerMsg(mythJobTaskSchedule.getNodeName() + ":" + MythLogUtils.getMessage(e)); jobTaskRunLogById.setTriggerMsg(mythJobTaskSchedule.getNodeName() + ":" + runErrorMessage);
FileSystem fileSystem = SpringUtil.getBean(FastDfsFileSystem.class);
HashMap<String, String> fileNameMap = new HashMap<>(1);
fileNameMap.put("filename", String.format("%s_error.log", mythJobTaskSchedule.getNodeName()));
String uploadFilePath;
try {
uploadFilePath = fileSystem.uploadFile(runErrorMessage.getBytes(StandardCharsets.UTF_8), "log", fileNameMap);
jobTaskRunLogById.setLogRemotelyPath(uploadFilePath);
} catch (Exception ex) {
String logUploadErrorMessage = MythLogUtils.getMessage(ex);
jobTaskRunLogById.setRunMsg(jobTaskRunLogById.getRunMsg() + "\n 日志上传异常:" + logUploadErrorMessage);
jobTaskRunLogById.setTriggerMsg(jobTaskRunLogById.getTriggerMsg() + "\n 日志上传异常:" + logUploadErrorMessage);
}
} }
jobTaskRunLogService.updateJobTaskRunLogWithBLOBs(jobTaskRunLogById); jobTaskRunLogService.updateJobTaskRunLogWithBLOBs(jobTaskRunLogById);
......
...@@ -2,11 +2,12 @@ package com.byit.task; ...@@ -2,11 +2,12 @@ package com.byit.task;
import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSON;
import com.byit.conf.MythJobAutoConfigure; import com.byit.conf.MythJobAutoConfigure;
import com.byit.dto.executor.RunParamWrapped;
import com.byit.dto.personalise.JavaJobConfDTO; import com.byit.dto.personalise.JavaJobConfDTO;
import com.byit.dto.plugin.JavaTask; import com.byit.dto.plugin.JavaTask;
import com.byit.enums.ScheduleTypeEnum; import com.byit.enums.ScheduleTypeEnum;
import com.byit.enums.task.RunResultEnum; import com.byit.enums.task.RunResultEnum;
import com.byit.filesystem.FastDfsFileSystem;
import com.byit.filesystem.FileSystem;
import com.byit.job.utils.DateUtil; import com.byit.job.utils.DateUtil;
import com.byit.job.utils.MythLogUtils; import com.byit.job.utils.MythLogUtils;
import com.byit.model.JobTaskRunLogWithBLOBs; import com.byit.model.JobTaskRunLogWithBLOBs;
...@@ -21,7 +22,9 @@ import io.netty.util.Timeout; ...@@ -21,7 +22,9 @@ import io.netty.util.Timeout;
import io.netty.util.TimerTask; import io.netty.util.TimerTask;
import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.StringUtils;
import java.nio.charset.StandardCharsets;
import java.util.Date; import java.util.Date;
import java.util.HashMap;
import java.util.List; import java.util.List;
import java.util.TreeSet; import java.util.TreeSet;
import java.util.stream.Collectors; import java.util.stream.Collectors;
...@@ -55,7 +58,7 @@ public class JavaTaskJobTask implements TimerTask { ...@@ -55,7 +58,7 @@ public class JavaTaskJobTask implements TimerTask {
String personalise = javaTask.getPersonalise(); String personalise = javaTask.getPersonalise();
String hasRelatedFlow = "1"; String hasRelatedFlow = "1";
JavaJobConfDTO javaJobConfDTO = JSON.parseObject(personalise, JavaJobConfDTO.class); JavaJobConfDTO javaJobConfDTO = JSON.parseObject(personalise, JavaJobConfDTO.class);
if(javaJobConfDTO != null && StringUtils.isNoneBlank(javaJobConfDTO.getHasRelatedFlow())){ if (javaJobConfDTO != null && StringUtils.isNoneBlank(javaJobConfDTO.getHasRelatedFlow())) {
hasRelatedFlow = javaJobConfDTO.getHasRelatedFlow(); hasRelatedFlow = javaJobConfDTO.getHasRelatedFlow();
} }
...@@ -65,9 +68,9 @@ public class JavaTaskJobTask implements TimerTask { ...@@ -65,9 +68,9 @@ public class JavaTaskJobTask implements TimerTask {
TreeSet<String> gatewayConf = getRegConfig.getGatewayConf(); TreeSet<String> gatewayConf = getRegConfig.getGatewayConf();
String finalHasRelatedFlow = hasRelatedFlow; String finalHasRelatedFlow = hasRelatedFlow;
List<String> callUrlList = gatewayConf.stream().map(callbackIp -> { List<String> callUrlList = gatewayConf.stream().map(callbackIp -> {
if("2".equals(finalHasRelatedFlow)){ if ("2".equals(finalHasRelatedFlow)) {
return String.format("%s%s%s", HTTP_PRE, callbackIp, HTTP_SUFFIX_FLOW); return String.format("%s%s%s", HTTP_PRE, callbackIp, HTTP_SUFFIX_FLOW);
}else{ } else {
return String.format("%s%s%s", HTTP_PRE, callbackIp, HTTP_SUFFIX); return String.format("%s%s%s", HTTP_PRE, callbackIp, HTTP_SUFFIX);
} }
}).collect(Collectors.toList()); }).collect(Collectors.toList());
...@@ -102,8 +105,23 @@ public class JavaTaskJobTask implements TimerTask { ...@@ -102,8 +105,23 @@ public class JavaTaskJobTask implements TimerTask {
log.setTriggerCode(RunResultEnum.TRIGGER_ERROR.getCode()); log.setTriggerCode(RunResultEnum.TRIGGER_ERROR.getCode());
log.setRunCode(RunResultEnum.RUN_ERROR.getCode()); log.setRunCode(RunResultEnum.RUN_ERROR.getCode());
log.setRunMsg(javaTask.getTaskName() + ":" + MythLogUtils.getMessage(e));
log.setTriggerMsg(javaTask.getTaskName() + ":" + MythLogUtils.getMessage(e)); String runErrorMsgger = MythLogUtils.getMessage(e);
log.setRunMsg(javaTask.getTaskName() + ":" + runErrorMsgger);
log.setTriggerMsg(javaTask.getTaskName() + ":" + runErrorMsgger);
FileSystem fileSystem = SpringUtil.getBean(FastDfsFileSystem.class);
HashMap<String, String> fileNameMap = new HashMap<>(1);
fileNameMap.put("filename", String.format("%s_error.log", javaTask.getTaskName()));
String uploadFilePath;
try {
uploadFilePath = fileSystem.uploadFile(runErrorMsgger.getBytes(StandardCharsets.UTF_8), "log", fileNameMap);
log.setLogRemotelyPath(uploadFilePath);
} catch (Exception ex) {
String uploadLogErrorMessage = MythLogUtils.getMessage(ex);
log.setRunMsg(log.getRunMsg() + "\n 日志上传异常:" + uploadLogErrorMessage);
log.setTriggerMsg(log.getTriggerMsg() + "\n 日志上传异常:" + uploadLogErrorMessage);
}
} }
JobTaskRunLogServiceImpl jobTaskRunLogService = SpringUtil.getBean(JobTaskRunLogServiceImpl.class); JobTaskRunLogServiceImpl jobTaskRunLogService = SpringUtil.getBean(JobTaskRunLogServiceImpl.class);
jobTaskRunLogService.updateJobTaskRunLogWithBLOBs(log); jobTaskRunLogService.updateJobTaskRunLogWithBLOBs(log);
......
...@@ -11,6 +11,8 @@ import com.byit.dto.plugin.RunLog; ...@@ -11,6 +11,8 @@ import com.byit.dto.plugin.RunLog;
import com.byit.enums.*; import com.byit.enums.*;
import com.byit.enums.task.RunResultEnum; import com.byit.enums.task.RunResultEnum;
import com.byit.enums.task.RunTypeEnum; import com.byit.enums.task.RunTypeEnum;
import com.byit.filesystem.FastDfsFileSystem;
import com.byit.filesystem.FileSystem;
import com.byit.job.utils.DateUtil; import com.byit.job.utils.DateUtil;
import com.byit.job.utils.MythLogUtils; import com.byit.job.utils.MythLogUtils;
import com.byit.job.utils.PlaceholderUtils; import com.byit.job.utils.PlaceholderUtils;
...@@ -28,10 +30,8 @@ import lombok.extern.slf4j.Slf4j; ...@@ -28,10 +30,8 @@ import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.StringUtils;
import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.data.redis.core.StringRedisTemplate;
import java.util.Date; import java.nio.charset.StandardCharsets;
import java.util.List; import java.util.*;
import java.util.Map;
import java.util.TreeSet;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import static com.alibaba.fastjson.serializer.SerializerFeature.WriteClassName; import static com.alibaba.fastjson.serializer.SerializerFeature.WriteClassName;
...@@ -137,11 +137,11 @@ public class ScriptExecutorJobTask implements TimerTask { ...@@ -137,11 +137,11 @@ public class ScriptExecutorJobTask implements TimerTask {
command = PlaceholderUtils.commandReplace(command, param); command = PlaceholderUtils.commandReplace(command, param);
} }
if (ScheduleTypeEnum.REPAIR.getCode().equals(mythJobTaskSchedule.getScheduleType())) { if (ScheduleTypeEnum.REPAIR.getCode().equals(mythJobTaskSchedule.getScheduleType())) {
if(command.startsWith(JAVA_TYPE)) { if (command.startsWith(JAVA_TYPE)) {
if(ScheduleTypeEnum.REPAIR.getCode().equals(mythJobTaskSchedule.getScheduleType())){ if (ScheduleTypeEnum.REPAIR.getCode().equals(mythJobTaskSchedule.getScheduleType())) {
command = PlaceholderUtils.commandDateReplace(param, command, param.get(PlaceholderEnum.DATE_PLACEHOLDER.getName()), command = PlaceholderUtils.commandDateReplace(param, command, param.get(PlaceholderEnum.DATE_PLACEHOLDER.getName()),
param.get(PlaceholderEnum.NOW_DATE_PLACEHOLDER.getName())); param.get(PlaceholderEnum.NOW_DATE_PLACEHOLDER.getName()));
}else{ } else {
command = PlaceholderUtils.commandDateReplace(param, command); command = PlaceholderUtils.commandDateReplace(param, command);
} }
} }
...@@ -163,6 +163,7 @@ public class ScriptExecutorJobTask implements TimerTask { ...@@ -163,6 +163,7 @@ public class ScriptExecutorJobTask implements TimerTask {
*/ */
private void runJob(JobTaskSchedule mythJobTaskSchedule) { private void runJob(JobTaskSchedule mythJobTaskSchedule) {
DispatchResponseDto dispatchResponseDto = new DispatchResponseDto(); DispatchResponseDto dispatchResponseDto = new DispatchResponseDto();
String uploadFilePath = null;
try { try {
paramBuild(mythJobTaskSchedule); paramBuild(mythJobTaskSchedule);
mythJobTaskSchedule.setRunCommand(mythJobTaskSchedule.getRunCommand()); mythJobTaskSchedule.setRunCommand(mythJobTaskSchedule.getRunCommand());
...@@ -179,7 +180,7 @@ public class ScriptExecutorJobTask implements TimerTask { ...@@ -179,7 +180,7 @@ public class ScriptExecutorJobTask implements TimerTask {
scriptDto.setRemotePath(mythJobTaskSchedule.getScriptUrls()); scriptDto.setRemotePath(mythJobTaskSchedule.getScriptUrls());
GetRegConfig getRegConfig = SpringUtil.getBean(GetRegConfig.class); GetRegConfig getRegConfig = SpringUtil.getBean(GetRegConfig.class);
TreeSet<String> gatewayConf = getRegConfig.getGatewayConf(); TreeSet<String> gatewayConf = getRegConfig.getGatewayConf();
if(CollectionUtil.isEmpty(gatewayConf)) { if (CollectionUtil.isEmpty(gatewayConf)) {
throw new RuntimeException("网关地址为空!"); throw new RuntimeException("网关地址为空!");
} }
List<String> callUrlList = gatewayConf.stream().map(callbackIp -> String.format("%s%s%s", HTTP_PRE, callbackIp, HTTP_SUFFIX)).collect(Collectors.toList()); List<String> callUrlList = gatewayConf.stream().map(callbackIp -> String.format("%s%s%s", HTTP_PRE, callbackIp, HTTP_SUFFIX)).collect(Collectors.toList());
...@@ -188,18 +189,30 @@ public class ScriptExecutorJobTask implements TimerTask { ...@@ -188,18 +189,30 @@ public class ScriptExecutorJobTask implements TimerTask {
scriptDto.setLogRemotePath(jobTaskRunLogById.getLogRemotelyPath()); scriptDto.setLogRemotePath(jobTaskRunLogById.getLogRemotelyPath());
dispatchResponseDto = runScriptService.runScript(scriptDto); dispatchResponseDto = runScriptService.runScript(scriptDto);
} catch (Exception e) { } catch (Exception e) {
//错误信息
String errorMessage = MythLogUtils.getMessage(e);
if (mythJobTaskSchedule.getScheduleType().equals(4)) { if (mythJobTaskSchedule.getScheduleType().equals(4)) {
StringRedisTemplate stringRedisTemplate = (StringRedisTemplate) SpringUtil.getBean("stringRedisTemplate"); StringRedisTemplate stringRedisTemplate = (StringRedisTemplate) SpringUtil.getBean("stringRedisTemplate");
RunLog runLog = RunLog.builder().isEnd(true).isSuccess(false).runLog("执行资源异常" + MythLogUtils.getMessage(e)).build(); RunLog runLog = RunLog.builder().isEnd(true).isSuccess(false).runLog("执行资源异常" + errorMessage).build();
stringRedisTemplate.opsForList().rightPush("REAL:EXEC:" + mythJobTaskSchedule.getLogId(), JSON.toJSONString(runLog, WriteClassName)); stringRedisTemplate.opsForList().rightPush("REAL:EXEC:" + mythJobTaskSchedule.getLogId(), JSON.toJSONString(runLog, WriteClassName));
} }
log.error("------执行机异常{}-----", MythLogUtils.getMessage(e)); log.error("------执行机异常{}-----", errorMessage);
dispatchResponseDto.setMsg(MythLogUtils.getMessage(e)); dispatchResponseDto.setMsg(errorMessage);
dispatchResponseDto.setCode(JobResultEnum.DISPATCH_FAIL.getCode()); dispatchResponseDto.setCode(JobResultEnum.DISPATCH_FAIL.getCode());
FileSystem fileSystem = SpringUtil.getBean(FastDfsFileSystem.class);
HashMap<String, String> fileNameMap = new HashMap<>(1);
fileNameMap.put("filename", String.format("%s_error.log", mythJobTaskSchedule.getNodeName()));
try {
uploadFilePath = fileSystem.uploadFile(errorMessage.getBytes(StandardCharsets.UTF_8), "log", fileNameMap);
} catch (Exception ex) {
String logUploadErrorMessage = MythLogUtils.getMessage(ex);
dispatchResponseDto.setMsg(dispatchResponseDto.getMsg() + "\n 日志上传异常:" + logUploadErrorMessage);
}
} }
saveLog(mythJobTaskSchedule, dispatchResponseDto); saveLog(mythJobTaskSchedule, dispatchResponseDto, uploadFilePath);
} }
...@@ -231,12 +244,16 @@ public class ScriptExecutorJobTask implements TimerTask { ...@@ -231,12 +244,16 @@ public class ScriptExecutorJobTask implements TimerTask {
throw new RuntimeException("服务器分配资源失败"); throw new RuntimeException("服务器分配资源失败");
} }
private void saveLog(JobTaskSchedule mythJobTaskSchedule, DispatchResponseDto dispatchResponseDto) { private void saveLog(JobTaskSchedule mythJobTaskSchedule, DispatchResponseDto dispatchResponseDto, String uploadFilePath) {
log.debug("-----------saveLog--保存脚本调度日志开始------------"); log.debug("-----------saveLog--保存脚本调度日志开始------------");
JobTaskRunLogServiceImpl jobTaskRunLogService = SpringUtil.getBean(JobTaskRunLogServiceImpl.class); JobTaskRunLogServiceImpl jobTaskRunLogService = SpringUtil.getBean(JobTaskRunLogServiceImpl.class);
JobTaskRunLogWithBLOBs jobTaskRunLogById = jobTaskRunLogService.findJobTaskRunLogById(mythJobTaskSchedule.getLogId()); JobTaskRunLogWithBLOBs jobTaskRunLogById = jobTaskRunLogService.findJobTaskRunLogById(mythJobTaskSchedule.getLogId());
if (StringUtils.isNoneBlank(uploadFilePath)) {
jobTaskRunLogById.setLogRemotelyPath(uploadFilePath);
}
JobTaskRunLogWithBLOBs jobTaskRunLog = new JobTaskRunLogWithBLOBs(); JobTaskRunLogWithBLOBs jobTaskRunLog = new JobTaskRunLogWithBLOBs();
//放置调度记录 //放置调度记录
if (jobTaskRunLogById.getRunCount() > 1) { if (jobTaskRunLogById.getRunCount() > 1) {
......
...@@ -435,6 +435,12 @@ public class JobUtils { ...@@ -435,6 +435,12 @@ public class JobUtils {
return JSON.parseObject(response, ResponseResult.class); return JSON.parseObject(response, ResponseResult.class);
} }
public static void main(String[] args) {
JobUtils.setTOKEN("huangfu");
JobUtils.setRequestUrl("http://10.0.120.208:8998/myth-job-admin");
System.out.println(stopFlow("皇甫", "科星"));
}
/** /**
* 删除工作流 * 删除工作流
* *
......
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