Commit 9521b2a3 by huangfusuper

优化plugin部分代码结构,调整任务服务器的运行方式

parent 2cdcde71
package com.byit.task;
import com.byit.api.ExecutorServer;
import com.byit.job.handler.interfaces.IJobHandler;
import com.byit.job.model.ReturnResult;
import com.byit.rpc.remoting.invoker.annotation.RpcReference;
import io.netty.util.Timeout;
import io.netty.util.TimerTask;
/**
* @program: byit-myth-job->JavaTask
* @description: 任务工作线程
* @author: huangfu
* @date: 2019/11/14 17:09
**/
public class JobTask implements TimerTask {
@RpcReference
private ExecutorServer executorServer;
private IJobHandler iJobHandler;
private String param;
public JobTask(IJobHandler iJobHandler, String param) {
this.iJobHandler = iJobHandler;
this.param = param;
}
@Override
public void run(Timeout timeout) throws Exception {
ReturnResult<String> result = executorServer.execute(iJobHandler, param);
}
}
package com.byit.job.enums.plugin;
import com.byit.job.enums.IEnum;
/**
* 插件枚举
* @author huangfu
*/
public enum PluginEnum implements IEnum {
REQUEST_PORT_OR_IP_IS_MISSING("请求ip或者port为null",10000);
private String msg;
private int code;
PluginEnum(String msg, int code) {
this.msg = msg;
this.code = code;
}
@Override
public int getCode() {
return this.code;
}
@Override
public String getMsg() {
return this.msg;
}
}
package com.byit.job.exceptions;
import com.byit.job.enums.IEnum;
/**
* @program: byit-myth-job->IException
* @description: 异常父类
* @author: huangfu
* @date: 2019/11/27 14:34
**/
public interface IException {
/**
* 返回错误枚举
* @return
*/
IEnum getIEnum();
}
package com.byit.job.exceptions.plugin;
import com.byit.job.enums.IEnum;
import com.byit.job.exceptions.IException;
/**
* @program: byit-myth-job->PluginException
* @description: 对于插件里面错误的异常封装
* @author: huangfu
* @date: 2019/11/27 14:37
**/
public class PluginException extends RuntimeException implements IException {
private IEnum iEnum;
public PluginException() {
}
public <E extends IEnum> PluginException(E e){
super(e.getMsg());
this.iEnum = e;
}
@Override
public IEnum getIEnum() {
return this.iEnum;
}
}
......@@ -42,27 +42,11 @@ public class JavaBeanJobInfo {
/**
* 添加任务的地址
*/
private String requestUrl;
private String requestIP;
private String param;
public JavaBeanJobInfo(String url, String requestUrl, String param, String routingStrategy, String blockingStrategy, String callbackToken, String gatewayToken) {
this.param = param;
this.requestUrl = requestUrl;
this.url = url;
this.routingStrategy = routingStrategy;
this.blockingStrategy = blockingStrategy;
this.callbackToken = callbackToken;
this.gatewayToken = gatewayToken;
}
public String getRequestUrl() {
return requestUrl;
}
private String requestPort;
public void setRequestUrl(String requestUrl) {
this.requestUrl = requestUrl;
}
private String param;
@Override
public String toString() {
......@@ -74,9 +58,41 @@ public class JavaBeanJobInfo {
", blockingStrategy='" + blockingStrategy + '\'' +
", callbackToken='" + callbackToken + '\'' +
", gatewayToken='" + gatewayToken + '\'' +
", requestIP='" + requestIP + '\'' +
", requestPort=" + requestPort +
", param='" + param + '\'' +
'}';
}
public String getRequestIP() {
return requestIP;
}
public void setRequestIP(String requestIP) {
this.requestIP = requestIP;
}
public String getRequestPort() {
return requestPort;
}
public void setRequestPort(String requestPort) {
this.requestPort = requestPort;
}
public JavaBeanJobInfo(String jobHandelName, String url, String mythCron, String routingStrategy, String blockingStrategy, String callbackToken, String gatewayToken, String requestIP, String requestPort, String param) {
this.jobHandelName = jobHandelName;
this.url = url;
this.mythCron = mythCron;
this.routingStrategy = routingStrategy;
this.blockingStrategy = blockingStrategy;
this.callbackToken = callbackToken;
this.gatewayToken = gatewayToken;
this.requestIP = requestIP;
this.requestPort = requestPort;
this.param = param;
}
public String getJobHandelName() {
return jobHandelName;
}
......
......@@ -15,11 +15,14 @@
<dependencies>
<!-- https://mvnrepository.com/artifact/com.alibaba/fastjson -->
<dependency>
<groupId>myth-job</groupId>
<artifactId>myth-executor-api</artifactId>
</dependency>
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>fastjson</artifactId>
<version>1.2.62</version>
</dependency>
<dependency>
......@@ -32,6 +35,21 @@
<artifactId>netty-all</artifactId>
</dependency>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-lang3</artifactId>
</dependency>
<!-- slf4j -->
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-api</artifactId>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-log4j12</artifactId>
</dependency>
</dependencies>
......
package com.byit.launcher;
import com.byit.rpc.ServerRunThread;
import lombok.extern.slf4j.Slf4j;
/**
* @program: byit-myth-job->JobRunLauncher
* @description: 运行job服务的启动器
* @author: huangfu
* @date: 2019/11/27 15:45
**/
@Slf4j
public class JobRunServerLauncher {
/**
* 线程是否已经被启动
*/
public static volatile String THREAD_RUN_MARK = null;
public JobRunServerLauncher(int port) {
runServer(port);
}
private void runServer(int port){
if(THREAD_RUN_MARK==null){
new Thread(new ServerRunThread(port)).start();
log.info("---------------------线程启动------------------");
}
}
}
package com.byit.rpc;
import cn.hutool.http.HttpUtil;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.byit.launcher.JobRunServerLauncher;
import com.byit.utils.JobUtils;
import com.byit.utils.ScanRootPackage;
import io.netty.bootstrap.ServerBootstrap;
......@@ -22,7 +20,7 @@ public class ServerRunThread implements Runnable {
public ServerRunThread(Integer port) {
new ScanRootPackage();
JobUtils.MARK = "SUCCESS";
JobRunServerLauncher.THREAD_RUN_MARK = "SUCCESS";
this.port = port;
}
......
......@@ -2,9 +2,13 @@ package com.byit.utils;
import cn.hutool.http.HttpUtil;
import com.alibaba.fastjson.JSON;
import com.byit.job.enums.plugin.PluginEnum;
import com.byit.job.exceptions.plugin.PluginException;
import com.byit.job.handler.interfaces.IJobHandler;
import com.byit.job.model.JavaBeanJobInfo;
import com.byit.rpc.ServerRunThread;
import lombok.extern.log4j.Log4j;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
......@@ -15,31 +19,34 @@ import java.util.concurrent.ConcurrentHashMap;
* @author: huangfu
* @date: 2019/11/18 12:30
**/
@Slf4j
public class JobUtils {
public static String MARK = null;
private static final String REQUEST_PREFIX = "http://";
private static final String REQUEST_ADD_JOB_RESOURCES_SUFFIX = "/job/addJob";
/**
* 任务的缓存
*/
public static final Map<String,Class<? extends IJobHandler>> jobCache = new ConcurrentHashMap<>();
/**
* 开始任务
*/
public static void jobServerStart(){
if(MARK==null){
new Thread(new ServerRunThread(9999)).start();
System.out.println("-------------线程启动----------------" );
}
}
/**
* 添加一个任务节点
* @param javaBeanJobInfo 任务节点的详尽配置
* @return 添加结果
*/
public static String addJob(JavaBeanJobInfo javaBeanJobInfo){
log.info("---------------开始添加一个任务,jobHandelName:{}---------------------",javaBeanJobInfo.getJobHandelName());
String requestIP = javaBeanJobInfo.getRequestIP();
String requestPort = javaBeanJobInfo.getRequestPort( );
if(StringUtils.isBlank(requestIP) || StringUtils.isBlank(requestPort)){
log.error("------------------任务添加失败----------------------");
throw new PluginException(PluginEnum.REQUEST_PORT_OR_IP_IS_MISSING);
}
//请求的路径
String requestUrl = REQUEST_PREFIX+requestIP+":"+requestPort+REQUEST_ADD_JOB_RESOURCES_SUFFIX;
//发送请求 添加任务
String post = HttpUtil.post(javaBeanJobInfo.getRequestUrl( ), JSON.toJSONString(javaBeanJobInfo));
return post;
String addRequestResult = HttpUtil.post(requestUrl, JSON.toJSONString(javaBeanJobInfo));
log.info("--------------------添加任务完成,添加结果为:{}------------------------",addRequestResult);
return addRequestResult;
}
}
......@@ -2,6 +2,7 @@ package com.byit.utils;
import com.byit.annotations.JobHandler;
import com.byit.job.handler.interfaces.IJobHandler;
import lombok.extern.slf4j.Slf4j;
import java.io.File;
import java.util.ArrayList;
......@@ -14,6 +15,7 @@ import java.util.stream.Collectors;
* @author: huangfu
* @date: 2019/11/19 17:31
**/
@Slf4j
public class ScanRootPackage {
public ScanRootPackage() {
addJobCache();
......
package com.byit.api;
import com.byit.job.handler.impl.GlueJobHandler;
import com.byit.job.handler.interfaces.IJobHandler;
import com.byit.job.model.ReturnResult;
/**
* @program: byit-myth-job->ExecutorServer
* @description: API外部调用接口
* @author: huangfu
* @date: 2019/11/13 16:52
**/
public interface ExecutorServer {
/**
* 执行接口
* @param iJobHandler 具体任务
* @param param 执行参数
* @return
*/
ReturnResult<String> execute(IJobHandler iJobHandler,String param) throws Exception;
}
package com.byit.api;
package com.byit.plugin.api;
import com.byit.job.model.JavaBeanJobInfo;
/**
* @program: byit-myth-job->JobOperating
* @description: 任务操作
* @description: 插件方 JobOperating 任务操作API
* @author: huangfu
* @date: 2019/11/18 11:21
* @date: 2019/11/27 13:44
**/
public interface JobOperating {
/**
......@@ -15,5 +14,4 @@ public interface JobOperating {
* @param javaBeanJobInfo 任务的详细配置
*/
void addJob(JavaBeanJobInfo javaBeanJobInfo);
}
......@@ -14,7 +14,7 @@ import com.byit.job.model.ReturnResult;
public class DemoJob extends BaseJobHandler {
@Override
public ReturnResult<String> execute(String s) throws Exception {
System.out.println("------------------------------" );
System.out.println("--------------任务就这样运行了----------------"+s );
return ReturnResult.SUCCESS;
}
}
......@@ -2,6 +2,7 @@ package com.byit.job;
import com.alibaba.fastjson.JSON;
import com.byit.job.model.JavaBeanJobInfo;
import com.byit.launcher.JobRunServerLauncher;
import com.byit.utils.JobUtils;
import java.io.IOException;
......@@ -20,16 +21,13 @@ public class Mains {
String blockingStrategy = "阻塞策略";
String callbackToken = "dsasadsad";
String gatewayToken="asdsadsa";
String requestUrl = "http://localhost:8080/job/addJob";
String requestIP = "127.0.0.1";
String requestPort="8080";
String param="sadsadsa";
String name="addJob";
JavaBeanJobInfo javaBeanJobInfo = new JavaBeanJobInfo(plServerUrl,requestUrl
,param,routingStrategy,blockingStrategy,callbackToken,gatewayToken);
javaBeanJobInfo.setJobHandelName(name);
javaBeanJobInfo.setMythCron(mythCron);
JavaBeanJobInfo javaBeanJobInfo = new JavaBeanJobInfo(name,plServerUrl,mythCron,routingStrategy,blockingStrategy,callbackToken,gatewayToken,requestIP,requestPort,param);
JobUtils.addJob(javaBeanJobInfo);
JobUtils.jobServerStart();
System.out.println(JobUtils.jobCache );
new JobRunServerLauncher(9999);
}
}
......@@ -52,6 +52,7 @@
<spring-cloud.version>Dalston.RELEASE</spring-cloud.version>
<lombok.version>1.18.4</lombok.version>
<hutool.version>4.5.11</hutool.version>
<commons.lang3.version>3.9</commons.lang3.version>
<maven-source-plugin.version>3.1.0</maven-source-plugin.version>
<maven-javadoc-plugin.version>3.1.1</maven-javadoc-plugin.version>
......@@ -139,6 +140,14 @@
<version>${hutool.version}</version>
</dependency>
<!-- https://mvnrepository.com/artifact/org.apache.commons/commons-lang3 -->
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-lang3</artifactId>
<version>${commons.lang3.version}</version>
</dependency>
</dependencies>
</dependencyManagement>
......
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