Commit f3e4d83d by huangfusuper

将结果回调修改为最大努力通知型

parent 7d0b920f
......@@ -14,6 +14,10 @@
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
</dependency>
<dependency>
<groupId>myth-job</groupId>
<artifactId>myth-admin-core</artifactId>
</dependency>
......
......@@ -27,7 +27,7 @@ mybatis:
myth-rpc:
registry:
address: http://localhost:8080/myth-register
address: http://10.0.120.208:8080/myth-register
env: huangfu
biz: byit-myth-job
logging:
......
......@@ -6,4 +6,6 @@ spring:
active: local
application:
name: myth-job-admin
gateway:
name: byit-myth-gateway
package com.byit;
import com.byit.util.GetRegConfig;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.test.context.junit4.SpringRunner;
import java.util.TreeSet;
@SpringBootTest
@RunWith(SpringRunner.class)
public class GatewayRegConfTest {
@Autowired
private GetRegConfig getRegConfig;
@Test
public void testReg(){
TreeSet<String> gatewayConf = getRegConfig.getGatewayConf();
gatewayConf.forEach(System.out::println);
}
}
......@@ -15,6 +15,7 @@ import com.byit.service.FlowStatusService;
import com.byit.service.impl.JobTaskRunLogServiceImpl;
import com.byit.service.impl.RunJavaServiceImpl;
import com.byit.job.utils.MythLogUtils;
import com.byit.util.GetRegConfig;
import com.byit.util.ServiceInfoUtil;
import com.byit.util.SpringUtil;
import io.netty.util.Timeout;
......@@ -23,6 +24,9 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.StringRedisTemplate;
import java.util.Date;
import java.util.List;
import java.util.TreeSet;
import java.util.stream.Collectors;
import static com.alibaba.fastjson.serializer.SerializerFeature.WriteClassName;
......@@ -33,6 +37,8 @@ import static com.alibaba.fastjson.serializer.SerializerFeature.WriteClassName;
@Slf4j
public class JavaNodeExecutorTask implements TimerTask {
private static final Integer INIT_SLEEP_TIME = 100;
public static final String HTTP_PRE = "http://";
public static final String HTTP_SUFFIX = "/myth-job-admin/api/callback/callbackRes";
private final JobTaskSchedule mythJobTaskSchedule;
......@@ -71,13 +77,15 @@ public class JavaNodeExecutorTask implements TimerTask {
jobTaskRunLogById = spinLock(jobTaskRunLogById);
//构建JAVA
PluginRpcRequestPacket request = new PluginRpcRequestPacket();
String ipAndPort = ServiceInfoUtil.getIpAndPort();
String callbackUrl = "http://" + ipAndPort + "/myth-job-admin/api/callback/callbackRes";
request.setCallbackUrl(callbackUrl);
GetRegConfig getRegConfig = SpringUtil.getBean(GetRegConfig.class);
TreeSet<String> gatewayConf = getRegConfig.getGatewayConf();
List<String> callUrlList = gatewayConf.stream().map(callbackIp -> String.format("%s%s%s", HTTP_PRE, callbackIp, HTTP_SUFFIX)).collect(Collectors.toList());
request.setCallbackUrl(JSON.toJSONString(callUrlList));
CommunicationParam communicationParam = new CommunicationParam();
communicationParam.setExpand1("REAL:EXEC:"+jobTaskRunLogById.getLogId());
communicationParam.setLogId(jobTaskRunLogById.getLogId()+"");
communicationParam.setCallbackUrl(callbackUrl);
communicationParam.setCallbackUrl(JSON.toJSONString(callUrlList));
communicationParam.setBody(mythJobTaskSchedule.getRunParam());
request.setParam(communicationParam);
......
package com.byit.task;
import com.alibaba.fastjson.JSON;
import com.byit.conf.MythJobAutoConfigure;
import com.byit.dto.plugin.JavaTask;
import com.byit.enums.ScheduleTypeEnum;
......@@ -12,19 +13,25 @@ import com.byit.param.CommunicationParam;
import com.byit.service.impl.JobTaskRunLogServiceImpl;
import com.byit.service.impl.RunJavaServiceImpl;
import com.byit.job.utils.MythLogUtils;
import com.byit.util.GetRegConfig;
import com.byit.util.ServiceInfoUtil;
import com.byit.util.SpringUtil;
import io.netty.util.Timeout;
import io.netty.util.TimerTask;
import java.util.Date;
import java.util.List;
import java.util.TreeSet;
import java.util.stream.Collectors;
/**
* java节点的执行器
* 立即运行 java单任务 java节点的执行器
* @author huangfu
*/
public class JavaTaskJobTask implements TimerTask {
public static final String JAVA_SYNC = "JAVA_SYNC";
public static final String HTTP_PRE = "http://";
public static final String HTTP_SUFFIX = "/myth-job-admin/api/callback/callbackRes";
private final JavaTask javaTask;
public JavaTaskJobTask(JavaTask javaTask) {
......@@ -42,12 +49,14 @@ public class JavaTaskJobTask implements TimerTask {
PluginRpcRequestPacket request = new PluginRpcRequestPacket();
//保存到日志
Integer logId = saveLog(javaTask);
String ipAndPort = ServiceInfoUtil.getIpAndPort();
String callbackUrl = "http://" + ipAndPort + "/myth-job-admin/api/callback/callbackRes";
request.setCallbackUrl(callbackUrl);
GetRegConfig getRegConfig = SpringUtil.getBean(GetRegConfig.class);
TreeSet<String> gatewayConf = getRegConfig.getGatewayConf();
List<String> callUrlList = gatewayConf.stream().map(callbackIp -> String.format("%s%s%s", HTTP_PRE, callbackIp, HTTP_SUFFIX)).collect(Collectors.toList());
request.setCallbackUrl(JSON.toJSONString(callUrlList));
CommunicationParam communicationParam = new CommunicationParam();
communicationParam.setLogId(logId+"");
communicationParam.setCallbackUrl(callbackUrl);
communicationParam.setCallbackUrl(JSON.toJSONString(callUrlList));
communicationParam.setBody(javaTask.getParam());
String dateFormat = DateUtil.dateFormat(new Date(javaTask.getTriggerTime()), DateUtil.FORMAT_DATE_TIME);
......
......@@ -8,13 +8,17 @@ import com.byit.dto.plugin.RunLog;
import com.byit.enums.*;
import com.byit.enums.task.RunResultEnum;
import com.byit.enums.task.RunTypeEnum;
import com.byit.factory.PluginClientFactory;
import com.byit.factory.PluginSpringClientFactory;
import com.byit.model.JobTaskRunLogWithBLOBs;
import com.byit.model.JobTaskSchedule;
import com.byit.registry.PluginServiceRegistry;
import com.byit.service.FastRunLogService;
import com.byit.service.FlowStatusService;
import com.byit.service.RunScriptService;
import com.byit.service.impl.JobTaskRunLogServiceImpl;
import com.byit.job.utils.MythLogUtils;
import com.byit.util.GetRegConfig;
import com.byit.util.ServiceInfoUtil;
import com.byit.util.SpringUtil;
import io.netty.util.Timeout;
......@@ -22,7 +26,11 @@ import io.netty.util.TimerTask;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.StringRedisTemplate;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;
import java.util.TreeSet;
import java.util.stream.Collectors;
import static com.alibaba.fastjson.serializer.SerializerFeature.WriteClassName;
......@@ -96,7 +104,10 @@ public class ScriptExecutorJobTask implements TimerTask {
scriptDto.setParam(mythJobTaskSchedule.getRunParam());
scriptDto.setRunId(mythJobTaskSchedule.getRunId());
scriptDto.setRemotePath(mythJobTaskSchedule.getScriptUrls());
scriptDto.setCallbackUrl(HTTP_PRE+ ServiceInfoUtil.getIpAndPort()+HTTP_SUFFIX);
GetRegConfig getRegConfig = SpringUtil.getBean(GetRegConfig.class);
TreeSet<String> gatewayConf = getRegConfig.getGatewayConf();
List<String> callUrlList = gatewayConf.stream().map(callbackIp -> String.format("%s%s%s", HTTP_PRE, callbackIp, HTTP_SUFFIX)).collect(Collectors.toList());
scriptDto.setCallbackUrl(JSON.toJSONString(callUrlList));
//二次执行的情况下 会有这个信息
scriptDto.setLogRemotePath(jobTaskRunLogById.getLogRemotelyPath());
dispatchResponseDto = runScriptService.runScript(scriptDto);
......
package com.byit.util;
import com.byit.utils.GetRegisteredDataUtil;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import java.util.TreeSet;
/**
* @author huangfu
*/
@Component
public class GetRegConfig {
@Value("${spring.gateway.name}")
private String gatewayName;
private final GetRegisteredDataUtil getRegisteredDataUtil;
public GetRegConfig(GetRegisteredDataUtil getRegisteredDataUtil) {
this.getRegisteredDataUtil = getRegisteredDataUtil;
}
/**
* 获取网关的ip+port
* @return
*/
public TreeSet<String> getGatewayConf() {
return getRegisteredDataUtil.getRegisteredIpAndPort(gatewayName);
}
}
......@@ -21,6 +21,8 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import java.util.Date;
import java.util.Iterator;
import java.util.List;
/**
* 命令式方式执行的执行机
......@@ -107,12 +109,19 @@ public class ImperativeExecutionImpl implements ScriptExecutorService {
jobRunResultDto.setLogId(scriptDto.getLogId());
//设置远程日志文件的路径
jobRunResultDto.setLogRemotelyPath(scriptDto.getLogRemotePath());
//最大努力通知性
String callbackUrls = scriptDto.getCallbackUrl();
List<String> callbackUrlList = JSON.parseArray(callbackUrls, String.class);
for (String callbackUrl : callbackUrlList) {
try {
RequestUtil.retryPost(scriptDto.getCallbackUrl(), JSON.toJSONString(jobRunResultDto), 3 , 1);
} catch (InterruptedException e) {
e.printStackTrace();
log.error("执行机异常:{}",ExecutorLogUtil.getMessage(e));
}
RequestUtil.retryPost(callbackUrl, JSON.toJSONString(jobRunResultDto), 3 , 1);
log.info("--------------runPythonScript,脚本调用结束-----------");
break;
}catch (Exception e) {
log.error("--------执行机节结果回复失败,失败原因是:{}--------",ExecutorLogUtil.getMessage(e));
}
}
}
}
......@@ -35,7 +35,7 @@ public class RequestUtil {
Thread.sleep(1000 * thisRetryCount);
retryPost(requestUrl,body,retryTotalCount,++thisRetryCount);
}else{
throw new RuntimeException(String.format("执行机回调调度中心失败,失败原因是:%s,重试了%s次", ExecutorLogUtil.getMessage(e),retryTotalCount));
throw new RuntimeException(String.format("执行机回调调度中心失败,地址为%s,失败原因是:%s,重试了%s次", requestUrl, ExecutorLogUtil.getMessage(e),retryTotalCount));
}
}
......
......@@ -20,7 +20,7 @@ import org.springframework.web.filter.CorsFilter;
* modified By
**/
@EnableZuulProxy
@RpcService(serverName = "[System]-myth-job-gateway")
@RpcService(http_type = true, serverName = "[System]-myth-job-gateway")
@SpringBootApplication
@Slf4j
public class GatewayApplication {
......
......@@ -42,7 +42,6 @@
<dependency>
<groupId>myth-job</groupId>
<artifactId>myth-register-client</artifactId>
<version>1.0-SNAPSHOT</version>
</dependency>
<dependency>
......
......@@ -15,7 +15,7 @@ public class CommunicationParam implements Serializable {
*/
private String logId;
/**
* 回调的url
* 回调的url 集合类型
*/
private String callbackUrl;
/**
......
......@@ -41,7 +41,7 @@ public class PluginRequestUtil {
Thread.sleep(1000 * thisRetryCount);
retryPost(requestUrl,body,retryTotalCount,++thisRetryCount);
}else{
throw new RuntimeException(String.format("执行机回调调度中心失败,失败原因是:%s,重试了%s次", PluginLogUtils.getMessage(e),retryTotalCount));
throw new RuntimeException(String.format("执行机回调调度中心失败,地址为:%s,失败原因是:%s,重试了%s次", requestUrl, PluginLogUtils.getMessage(e),retryTotalCount));
}
}
......
......@@ -117,5 +117,7 @@ public class PluginSpringClientFactory extends InstantiationAwareBeanPostProcess
pluginClientFactory.close();
}
public PluginClientFactory getPluginClientFactory() {
return pluginClientFactory;
}
}
......@@ -16,6 +16,7 @@ import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
import io.netty.handler.timeout.IdleStateEvent;
import java.util.List;
import java.util.Map;
import java.util.concurrent.Callable;
import java.util.concurrent.ThreadPoolExecutor;
......@@ -95,15 +96,17 @@ public class NettyPluginServerHandler extends SimpleChannelInboundHandler<Plugin
rpcResponsePacket.setType(ResponseTyEnum.RESPONSE.getType());
rpcResponsePacket.setRequestId(msg.getRequestId());
String responseStr = JSON.toJSONString(rpcResponsePacket);
String callbackUrl = msg.getCallbackUrl();
String callbackUrls = msg.getCallbackUrl();
List<String> callbackUrlList = JSON.parseArray(callbackUrls, String.class);
for (String callbackUrl : callbackUrlList) {
try {
PluginRequestUtil.retryPost(callbackUrl,responseStr,3,1);
} catch (InterruptedException e) {
e.printStackTrace();
System.err.println(String.format("执行机异常:%s", PluginLogUtils.getMessage(e)));
break;
}catch (Exception e) {
System.err.println(String.format("执行机节结果回复失败,失败原因是:%s", PluginLogUtils.getMessage(e)));
}
}
}
/**
......
......@@ -8,6 +8,7 @@ import com.byit.model.MythConfigModel;
import com.byit.model.PluginConfigurationModel;
import com.byit.model.RegisteredConfigurationModel;
import com.byit.server.netty.NettyPluginServer;
import com.byit.utils.GetRegisteredDataUtil;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.condition.ConditionalOnBean;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
......@@ -46,4 +47,10 @@ public class PluginAutoConfigure {
RegisteredConfigurationModel register = mythConfigModel.getRegister();
return new PluginSpringClientFactory(register.getUrl(),plugin.getEnv(),plugin.getBiz());
}
@Bean
@ConditionalOnBean(PluginClientMark.class)
public GetRegisteredDataUtil getRegisteredDataUtil(){
GetRegisteredDataUtil getRegisteredDataUtil = new GetRegisteredDataUtil(pluginSpringClientFactory());
return getRegisteredDataUtil;
}
}
package com.byit.utils;
import com.byit.factory.PluginClientFactory;
import com.byit.factory.PluginSpringClientFactory;
import com.byit.registry.PluginServiceRegistry;
import java.util.TreeSet;
/**
* 注册中心配置
* @author huangfu
*/
public class GetRegisteredDataUtil {
/**
* 插件客户端工厂
*/
private PluginSpringClientFactory pluginSpringClientFactory;
public GetRegisteredDataUtil(PluginSpringClientFactory pluginSpringClientFactory) {
this.pluginSpringClientFactory = pluginSpringClientFactory;
}
public TreeSet<String> getRegisteredIpAndPort(String serverName){
//获取插件工厂对象
PluginClientFactory pluginClientFactory = pluginSpringClientFactory.getPluginClientFactory();
//获取插件的注册中心使用管理器
PluginServiceRegistry pluginServiceRegistry = pluginClientFactory.getPluginServiceRegistry();
return pluginServiceRegistry.discovery(serverName);
}
}
......@@ -72,12 +72,18 @@
<org-apache-commons.version>1.3</org-apache-commons.version>
<fastdfs-client-java-version>1.27.0.0</fastdfs-client-java-version>
<fastdfs-client-version>1.26.1-RELEASE</fastdfs-client-version>
<register-client-version>1.0-SNAPSHOT</register-client-version>
</properties>
<dependencyManagement>
<dependencies>
<dependency>
<groupId>myth-job</groupId>
<artifactId>myth-register-client</artifactId>
<version>${register-client-version}</version>
</dependency>
<!-- fastdfs-client 依赖-->
<dependency>
<groupId>com.github.tobato</groupId>
......
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