Commit 64146270 by huangfusuper

完成执行机链式执行的代码冲重构

parent 26d1b60c
...@@ -161,8 +161,8 @@ ...@@ -161,8 +161,8 @@
select select
<include refid="Base_Column_List" /> <include refid="Base_Column_List" />
from run_recording from run_recording
where start_time &gt;= #{startDate} where triggerTime &gt;= #{startDate}
and start_time &lt;= #{endDate} and triggerTime &lt;= #{endDate}
<if test="flowIds != null"> <if test="flowIds != null">
and flow_id in ( and flow_id in (
<foreach collection="flowIds" item="flowId" separator=","> <foreach collection="flowIds" item="flowId" separator=",">
......
...@@ -38,4 +38,13 @@ public class ScriptDto implements Serializable { ...@@ -38,4 +38,13 @@ public class ScriptDto implements Serializable {
* 已有的日志路径,不存在就为null * 已有的日志路径,不存在就为null
*/ */
private String logRemotePath; private String logRemotePath;
/**
* 日志信心
*/
private String runLogData;
/**
* 执行结果
*/
private Integer exitCode;
} }
package com.byit.node.impl;
import com.byit.dto.executor.ScriptDto;
import com.byit.exceptions.ExecutorException;
import com.byit.filesystem.FileSystem;
import com.byit.node.ProcessingNodeChain;
import com.byit.node.machine.ProcessingMachine;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import java.nio.charset.StandardCharsets;
import java.util.HashMap;
import java.util.Map;
/**
* 日志信息处理节点
* @author huangfu
*/
@Component
@Slf4j
public class LogDataProcessingMachine implements ProcessingMachine {
private final FileSystem fileSystem;
private static final String FILE_NAME = "filename";
/**
* 日志后缀
*/
public static final String LOG_SUFFIX = "log";
/**
* 带点的日志后缀
*/
public static final String LOG_RETOUCH_SUFFIX = ".log";
public LogDataProcessingMachine(FileSystem fileSystem) {
this.fileSystem = fileSystem;
}
@Override
public void executor(ScriptDto scriptDto, ProcessingNodeChain processingNodeChain) {
try{
String logPath;
byte[] logDataByte = stringToByteArray(scriptDto.getRunLogData());
Map<String,String> fileMateData = new HashMap<>(2);
fileMateData.put(FILE_NAME,scriptDto.getRunId()+scriptDto.getRunId()+ LOG_RETOUCH_SUFFIX);
//不为空 则追加
if(scriptDto.getLogRemotePath() != null){
byte[] sourceLogByte = fileSystem.downloaderFile(scriptDto.getLogRemotePath());
byte[] resultLogByte = mergeFile(sourceLogByte, logDataByte);
//上传日志文件
logPath = fileSystem.uploadFile(resultLogByte, LOG_SUFFIX,fileMateData);
//删除原有的日志文件
fileSystem.fileRemove(scriptDto.getLogRemotePath());
}else{
//上传日志文件
logPath = fileSystem.uploadFile(logDataByte,LOG_SUFFIX,fileMateData);
}
scriptDto.setLogRemotePath(logPath);
}catch (Exception e) {
e.printStackTrace();
throw new ExecutorException(e);
}
}
/**
* 将字符串转换我数组
* @param logData 日志字符串
* @return 转换的字节数组
*/
private byte[] stringToByteArray(String logData) {
return logData.getBytes(StandardCharsets.UTF_8);
}
/**
* 合并两个字节数组
* @param sourceByte 字节数组1
* @param targetByte 字节数组2
* @return 合并后的字节数组
*/
private byte[] mergeFile(byte[] sourceByte, byte[] targetByte){
byte[] result = new byte[sourceByte.length+targetByte.length];
System.arraycopy(sourceByte,0,result,0,sourceByte.length);
System.arraycopy(targetByte,0,result,sourceByte.length,targetByte.length);
return result;
}
}
...@@ -6,17 +6,21 @@ import com.byit.enums.PlaceholderEnum; ...@@ -6,17 +6,21 @@ import com.byit.enums.PlaceholderEnum;
import com.byit.exceptions.ExecutorException; import com.byit.exceptions.ExecutorException;
import com.byit.node.ProcessingNodeChain; import com.byit.node.ProcessingNodeChain;
import com.byit.node.machine.ProcessingMachine; import com.byit.node.machine.ProcessingMachine;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.StringUtils;
import org.springframework.stereotype.Component;
/** /**
* 1.数据校验处理节点 * 1.数据校验处理节点
* 主要就是校验 scriptDto对象内部参数的合法性,若不合法则直接异常 * 主要就是校验 scriptDto对象内部参数的合法性,若不合法则直接异常
* @author huangfu * @author huangfu
*/ */
@Component
@Slf4j
public class ParameterVerificationProcessingMachine implements ProcessingMachine { public class ParameterVerificationProcessingMachine implements ProcessingMachine {
@Override @Override
public void executor(ScriptDto scriptDto, ProcessingNodeChain processingNodeChain) { public void executor(ScriptDto scriptDto, ProcessingNodeChain processingNodeChain) {
log.debug("----------开始交验执行参数-------------");
//参数为null //参数为null
if (null == scriptDto) { if (null == scriptDto) {
throw new ExecutorException(ExecutorExceptionEnum.REQUIRED_PARAMETERS_ARE_EMPTY); throw new ExecutorException(ExecutorExceptionEnum.REQUIRED_PARAMETERS_ARE_EMPTY);
...@@ -51,5 +55,6 @@ public class ParameterVerificationProcessingMachine implements ProcessingMachine ...@@ -51,5 +55,6 @@ public class ParameterVerificationProcessingMachine implements ProcessingMachine
// 回调处理链上的下一个节点 // 回调处理链上的下一个节点
processingNodeChain.doProcessing(scriptDto,processingNodeChain); processingNodeChain.doProcessing(scriptDto,processingNodeChain);
log.debug("---------------校验执行参数完毕---------------");
} }
} }
package com.byit.node.impl;
import com.byit.dto.executor.ScriptDto;
import com.byit.node.ProcessingNodeChain;
import com.byit.node.machine.ProcessingMachine;
import com.byit.process.MythJobProcess;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;
import java.util.Arrays;
import java.util.List;
/**
* 执行脚本的处理节点
* @author huangfu
*/
@Component
@Slf4j
public class RunScriptProcessingMachine implements ProcessingMachine {
/**
* 节点被杀死
*/
private static final Integer NODE_KILL_CODE = -100;
private static final String KILL_MESSAGE = "该任务节点已被强制kill!";
private final StringRedisTemplate stringRedisTemplate;
public RunScriptProcessingMachine(StringRedisTemplate stringRedisTemplate) {
this.stringRedisTemplate = stringRedisTemplate;
}
@Override
public void executor(ScriptDto scriptDto, ProcessingNodeChain processingNodeChain) {
List<String> cmdList = Arrays.asList(scriptDto.getCommand().split(" "));
MythJobProcess mythJobProcess = new MythJobProcess(cmdList, null, null, scriptDto.getLogId(), stringRedisTemplate);
//保存日志
String logData = mythJobProcess.call();
//成功 是0 失败是其他的 杀死是-100
int exitCode = mythJobProcess.getExitCode();
if(exitCode == NODE_KILL_CODE){
logData = KILL_MESSAGE;
}
scriptDto.setRunLogData(logData);
scriptDto.setExitCode(exitCode);
processingNodeChain.doProcessing(scriptDto,processingNodeChain);
}
}
package com.byit.service;
import com.alibaba.fastjson.JSON;
import com.byit.dto.executor.DispatchResponseDto;
import com.byit.dto.executor.JobRunResultDto;
import com.byit.dto.executor.ScriptDto;
import com.byit.dto.web.ReturnResult;
import com.byit.enums.JobResultEnum;
import com.byit.executor.api.ScriptExecutorService;
import com.byit.node.ProcessingNodeChain;
import com.byit.node.impl.CommandAndScriptProcessingMachine;
import com.byit.node.impl.LogDataProcessingMachine;
import com.byit.node.impl.ParameterVerificationProcessingMachine;
import com.byit.node.impl.RunScriptProcessingMachine;
import com.byit.pool.RunThreadPool;
import com.byit.rpc.remoting.provider.annotation.RpcService;
import com.byit.utils.ServiceInfoUtil;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import java.util.Date;
/**
* 命令式方式执行的执行机
* @author huangfu
*/
@Component
@RpcService
@Slf4j
public class ImperativeExecutionImpl implements ScriptExecutorService {
private final CommandAndScriptProcessingMachine commandAndScriptProcessingMachine;
private final LogDataProcessingMachine logDataProcessingMachine;
private final RunScriptProcessingMachine runScriptProcessingMachine;
private final ParameterVerificationProcessingMachine parameterVerificationProcessingMachine;
private static final Integer NODE_RUN_SUCCESS_CODE = 0;
/**
* 节点被杀死
*/
private static final Integer NODE_KILL_CODE = -100;
public ImperativeExecutionImpl(CommandAndScriptProcessingMachine commandAndScriptProcessingMachine, LogDataProcessingMachine logDataProcessingMachine, RunScriptProcessingMachine runScriptProcessingMachine, ParameterVerificationProcessingMachine parameterVerificationProcessingMachine) {
this.commandAndScriptProcessingMachine = commandAndScriptProcessingMachine;
this.logDataProcessingMachine = logDataProcessingMachine;
this.runScriptProcessingMachine = runScriptProcessingMachine;
this.parameterVerificationProcessingMachine = parameterVerificationProcessingMachine;
}
@Override
public DispatchResponseDto runScript(ScriptDto scriptDto) {
log.info("----------------runScript 开始调用脚本执行机 start{} ------------------------",scriptDto);
DispatchResponseDto dispatchResponseDto;
try {
RunThreadPool.SCRIPT_RUN_THREAD_POOL.execute(()-> System.out.println("1"));
dispatchResponseDto = DispatchResponseDto.builder()
.code(JobResultEnum.DISPATCH_SUCCESS.getCode())
.msg(JobResultEnum.DISPATCH_SUCCESS.getMsg())
.build();
}catch (Exception e){
dispatchResponseDto = DispatchResponseDto.builder()
.code(JobResultEnum.DISPATCH_FAIL.getCode())
.msg(JobResultEnum.DISPATCH_FAIL.getMsg())
.build();
}
dispatchResponseDto.setUrl(ServiceInfoUtil.getIpAndPort());
log.info("----------------runScript 调用脚本执行机结束 end---------------------");
return dispatchResponseDto;
}
public void startScript(ScriptDto scriptDto){
log.info("--------------runPythonScript,脚本调用开始-----------");
//创建回复对象
JobRunResultDto jobRunResultDto = new JobRunResultDto();
try {
//功能链添加
ProcessingNodeChain processingNodeChain = new ProcessingNodeChain().addProcessingNode(parameterVerificationProcessingMachine)
.addProcessingNode(commandAndScriptProcessingMachine)
.addProcessingNode(runScriptProcessingMachine)
.addProcessingNode(logDataProcessingMachine);
//功能链执行
processingNodeChain.doProcessing(scriptDto,processingNodeChain);
//判断任务的执行状态 0成功 -100kill 其他 失败
if(NODE_RUN_SUCCESS_CODE.equals(scriptDto.getExitCode())){
jobRunResultDto.setReturnResult(ReturnResult.SUCCESS);
}else if(NODE_KILL_CODE.equals(scriptDto.getExitCode())){
ReturnResult<String> returnResult = new ReturnResult<>();
returnResult.setCode(JobResultEnum.KILL_SUCCESS.getCode());
returnResult.setMsg(JobResultEnum.KILL_SUCCESS.getMsg());
jobRunResultDto.setReturnResult(returnResult);
}else{
jobRunResultDto.setReturnResult(ReturnResult.FAIL);
}
}catch (Exception e) {
log.error("调度执行机出现异常:{}",e.getMessage());
ReturnResult<String> returnResult = new ReturnResult<>();
returnResult.setCode(ReturnResult.FAIL.getCode());
returnResult.setMsg(e.getMessage());
jobRunResultDto.setReturnResult(returnResult);
}
//设置结束时间
jobRunResultDto.setEndTime(new Date());
jobRunResultDto.setLogId(scriptDto.getLogId());
//设置远程日志文件的路径
jobRunResultDto.setLogRemotelyPath(scriptDto.getLogRemotePath());
cn.hutool.http.HttpUtil.post(scriptDto.getCallbackUrl(), JSON.toJSONString(jobRunResultDto));
log.info("--------------runPythonScript,脚本调用结束-----------");
}
}
...@@ -31,8 +31,6 @@ import java.util.*; ...@@ -31,8 +31,6 @@ import java.util.*;
* 脚本执行的实现 * 脚本执行的实现
* @author huangfu * @author huangfu
*/ */
@Service
@RpcService
@Slf4j @Slf4j
public class ScriptExecutorServiceImpl implements ScriptExecutorService { public class ScriptExecutorServiceImpl implements ScriptExecutorService {
/** /**
......
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