Commit a4c54902 by huangfusuper

改变初始参数问题

parent a45be81b
package com.byit.task;
import com.byit.conf.MythJobAutoConfigure;
import com.byit.conf.RpcResultHttpCallback;
import com.byit.dto.plugin.JavaTask;
import com.byit.dto.plugin.JobTaskRunLog;
import com.byit.enums.FlowPropertyEnum;
import com.byit.enums.NodePropertyEnum;
import com.byit.enums.NodeRunStatusPropertyEnum;
import com.byit.enums.ScheduleTypeEnum;
import com.byit.model.JobTaskRunLogWithBLOBs;
import com.byit.packet.request.PluginRpcRequestPacket;
import com.byit.packet.response.PluginRpcResponsePacket;
import com.byit.param.CommunicationParam;
import com.byit.param.defaultparam.DefaultResultCallback;
import com.byit.service.impl.JobTaskRunLogServiceImpl;
import com.byit.service.impl.RunJavaServiceImpl;
import com.byit.util.ServiceInfoUtil;
import com.byit.util.SpringUtil;
import io.netty.util.Timeout;
import io.netty.util.TimerTask;
......@@ -29,31 +35,54 @@ public class JavaTaskJobTask implements TimerTask {
@Override
public void run(Timeout timeout) throws Exception {
MythJobAutoConfigure.LOW_LEVEL_JOB_THREAD_POOL.execute(()->{
runJob();
});
}
private void runJob(){
RunJavaServiceImpl service = SpringUtil.getBean(RunJavaServiceImpl.class);
PluginRpcRequestPacket request = new PluginRpcRequestPacket();
//保存到日志
Integer logId = saveLog(javaTask);
request.setExtension(logId+"");
request.setCallbackUrl("http://127.0.0.1:8081/myth-job-admin/api/callback/callbackRes");
request.setParam(javaTask.getParam());
request.setJobName(javaTask.getTaskName());
PluginRpcResponsePacket pluginRpcResponsePacket = service.runJava(request);
String ipAndPort = ServiceInfoUtil.getIpAndPort();
String callbackUrl = "http://" + ipAndPort + "/myth-job-admin/api/callback/callbackRes";
request.setCallbackUrl(callbackUrl);
CommunicationParam communicationParam = new CommunicationParam();
communicationParam.setLogId(logId+"");
communicationParam.setCallbackUrl(callbackUrl);
communicationParam.setBody(javaTask.getParam());
communicationParam.setBody("test");
request.setParam(communicationParam);
request.setJobName(javaTask.getTaskName());
JobTaskRunLogWithBLOBs log = new JobTaskRunLogWithBLOBs();
log.setLogId(logId);
if(pluginRpcResponsePacket.isStatus()){
log.setTriggerCode("1");
}else{
try {
PluginRpcResponsePacket pluginRpcResponsePacket = service.runJava(request);
if(pluginRpcResponsePacket.isStatus()){
log.setTriggerCode(NodeRunStatusPropertyEnum.RUN_SUCCESS.getCode());
}else{
log.setTriggerCode(NodeRunStatusPropertyEnum.RUN_FAILURE.getCode());
}
log.setTriggerMsg(pluginRpcResponsePacket.getMsg());
}catch (Exception e){
log.setTriggerCode("2");
log.setRunCode(NodeRunStatusPropertyEnum.RUN_FAILURE.getCode());
log.setRunMsg(javaTask.getTaskName()+":"+e.getMessage());
log.setTriggerMsg(javaTask.getTaskName()+":"+e.getMessage());
}
log.setTriggerMsg(pluginRpcResponsePacket.getMsg());
JobTaskRunLogServiceImpl jobTaskRunLogService = SpringUtil.getBean(JobTaskRunLogServiceImpl.class);
jobTaskRunLogService.updateJobTaskRunLogWithBLOBs(log);
}
private Integer saveLog(JavaTask javaTask){
JobTaskRunLogWithBLOBs log = new JobTaskRunLogWithBLOBs();
log.setStartTime(new Date());
log.setScheduleType(ScheduleTypeEnum.JAVA_SYNC.getCode());
log.setNodeId(javaTask.getId());
log.setNodeName(javaTask.getJobName());
......
......@@ -3,6 +3,7 @@ package com.byit.packet.request;
import com.byit.enums.Command;
import com.byit.enums.SerializerAlgorithm;
import com.byit.packet.BasePacketModel;
import com.byit.param.CommunicationParam;
import com.byit.param.ResultCallback;
import lombok.Data;
import lombok.EqualsAndHashCode;
......@@ -27,7 +28,7 @@ public class PluginRpcRequestPacket extends BasePacketModel {
*/
private String extension;
private String param;
private CommunicationParam param;
private String callbackUrl;
......
......@@ -25,6 +25,7 @@ public class PluginRpcResponsePacket extends BasePacketModel {
private String msg;
private Object result;
private String type;
private long runTime;
private boolean status = false;
@Override
public Command getCommand() {
......
package com.byit.param;
import lombok.Data;
import java.io.Serializable;
/**
* 执行参数
* @author huangfu
*/
@Data
public class CommunicationParam implements Serializable {
/**
* 日志id
*/
private String logId;
/**
* 回调的url
*/
private String callbackUrl;
/**
* 服务端参数
*/
private String body;
}
......@@ -2,6 +2,7 @@ package com.byit.task.handler.interfaces;
import com.byit.dto.web.ReturnResult;
import com.byit.param.CommunicationParam;
/**
* @program: byit-myth-job->IJobHandler
......@@ -17,7 +18,7 @@ public interface IJobHandler {
* @return
* @throws Exception
*/
ReturnResult<String> execute(String param) throws Exception;
ReturnResult<String> execute(CommunicationParam param) throws Exception;
/**
* 初始化时调用
......
......@@ -50,11 +50,14 @@ public class NettyPluginServerHandler extends SimpleChannelInboundHandler<Plugin
String jobName = msg.getJobName();
Object bean = serverPoll.get(jobName);
IJobHandler iJobHandler = (IJobHandler)bean;
long startTime = System.currentTimeMillis();
ReturnResult<String> execute = iJobHandler.execute(msg.getParam());
long endTime = System.currentTimeMillis();
rpcResponsePacket.setResult(execute);
rpcResponsePacket.setCode("000000");
rpcResponsePacket.setMsg("SUCCESS");
rpcResponsePacket.setStatus(true);
rpcResponsePacket.setRunTime(endTime-startTime);
rpcResponsePacket.setExtension(msg.getExtension());
}catch (Exception e){
rpcResponsePacket.setCode("500000");
......@@ -64,12 +67,6 @@ public class NettyPluginServerHandler extends SimpleChannelInboundHandler<Plugin
}
rpcResponsePacket.setType(ResponseTyEnum.RESPONSE.getType());
rpcResponsePacket.setRequestId(msg.getRequestId());
//ctx.channel().writeAndFlush(rpcResponsePacket);
// ResultCallback resultCallback = msg.getResultCallback();
// if (resultCallback != null) {
// System.out.println(resultCallback);
// resultCallback.resultCallback(ctx,rpcResponsePacket);
// }
String responseStr = JSON.toJSONString(rpcResponsePacket);
HttpRequest post = HttpRequest.post(msg.getCallbackUrl());
......
package com.byit.server;
import com.byit.dto.web.ReturnResult;
import com.byit.param.CommunicationParam;
import com.byit.task.annotations.TaskHandler;
import com.byit.task.handler.BaseJobHandler;
@TaskHandler(taskName = "sendEmailTest")
public class SendEmailTest extends BaseJobHandler {
@Override
public ReturnResult<String> execute(String param) throws Exception {
public ReturnResult<String> execute(CommunicationParam param) throws Exception {
return null;
}
......
package com.byit.server;
import com.byit.dto.web.ReturnResult;
import com.byit.param.CommunicationParam;
import com.byit.task.annotations.TaskHandler;
import com.byit.task.handler.BaseJobHandler;
import org.springframework.stereotype.Component;
......@@ -12,7 +13,7 @@ import org.springframework.stereotype.Component;
@TaskHandler(taskName = "01Test")
public class SenEmailServer extends BaseJobHandler {
@Override
public ReturnResult<String> execute(String param) throws Exception {
public ReturnResult<String> execute(CommunicationParam param) throws Exception {
return null;
}
}
package com.byit.server;
import com.byit.dto.web.ReturnResult;
import com.byit.param.CommunicationParam;
import com.byit.task.annotations.TaskHandler;
import com.byit.task.handler.BaseJobHandler;
import org.springframework.stereotype.Component;
......@@ -12,8 +13,8 @@ import org.springframework.stereotype.Component;
@TaskHandler(cron = "0 0/1 * * * ?",taskName = "sentEmailServer",autoPublish = true,publishUrl = "http://127.0.0.1:8081/myth-job-admin/api/node/autoAddJavaTask")
public class SentEmailServer extends BaseJobHandler {
@Override
public ReturnResult<String> execute(String param) throws Exception {
System.out.println("-------------SentEmailServer-被调度执行------------");
public ReturnResult<String> execute(CommunicationParam param) throws Exception {
System.out.println("-------------SentEmailServer-被调度执行------------"+param);
return ReturnResult.SUCCESS;
}
}
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