Commit 51969212 by huangfusuper

Rpc 插件的启动方式更改为Spring的启动方式

parent dc97ac08
package com.byit.scan; package com.byit.scan;
import com.byit.annotations.JobHandler; import com.byit.task.annotations.JobHandler;
import com.byit.scan.base.IScanProject; import com.byit.scan.base.IScanProject;
import com.byit.utils.JobUtils; import com.byit.utils.JobUtils;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
......
package com.byit.scan; package com.byit.scan;
import com.byit.annotations.JobHandler; import com.byit.task.annotations.JobHandler;
import com.byit.executor.handler.interfaces.IJobHandler; import com.byit.executor.handler.interfaces.IJobHandler;
import com.byit.scan.base.IScanProject; import com.byit.scan.base.IScanProject;
import com.byit.utils.JobUtils; import com.byit.utils.JobUtils;
......
/* /*
package com.byit.scan; package com.byit.scan;
import com.byit.annotations.JobHandler; import com.byit.task.annotations.JobHandler;
import com.byit.executor.handler.interfaces.IJobHandler; import com.byit.executor.handler.interfaces.IJobHandler;
import com.byit.utils.JobUtils; import com.byit.utils.JobUtils;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
......
package com.byit.scan.base; package com.byit.scan.base;
import com.byit.annotations.JobHandler; import com.byit.task.annotations.JobHandler;
import com.byit.executor.handler.interfaces.IJobHandler; import com.byit.executor.handler.interfaces.IJobHandler;
/** /**
......
...@@ -23,4 +23,5 @@ public class ServerConfigurationModel { ...@@ -23,4 +23,5 @@ public class ServerConfigurationModel {
private Integer maxSize; private Integer maxSize;
private String pluginServiceRegistry; private String pluginServiceRegistry;
private String pluginServerClass; private String pluginServerClass;
private String remainingOperationsClass;
} }
package com.byit.task.handler;
import com.byit.task.handler.interfaces.IJobHandler;
/**
* @program: byit-myth-job->AbsIJobHandler
* @description: TODO
* @author: huangfu
* @date: 2019/11/13 16:06
**/
public abstract class BaseJobHandler implements IJobHandler {
@Override
public void init(String param) {
}
@Override
public void destroy(String param) {
}
}
package com.byit.task.handler.interfaces;
import com.byit.dto.web.ReturnResult;
/**
* @program: byit-myth-job->IJobHandler
* @description: 所有执行器的执行类的基类,即所有类型的执行器都是此接口的实现类
* @author: huangfu
* @date: 2019/11/13 15:31
**/
public interface IJobHandler {
/**
* 所有的任务执行程序必须回调此方法执行
* @param param
* @return
* @throws Exception
*/
ReturnResult<String> execute(String param) throws Exception;
/**
* 初始化时调用
* @param param
*/
void init(String param);
/**
* 销毁时调用
* @param param
*/
void destroy(String param);
}
package com.byit.callback;
import java.util.Map;
/**
* 默认的后续处理器
* @author huangfu
*/
public class DefaultRemainingOperationsCallBack implements RemainingOperationsCallBack {
@Override
public void serverRemainingOperations(Map<String, Object> map) {
System.out.println(map);
}
}
package com.byit.callback;
import java.util.Map;
/**
* 剩余操作的接口回调类
* @author huangfu
*/
public interface RemainingOperationsCallBack {
/**
* 服务的后续处理
* @param map 被注册的服务
*/
void serverRemainingOperations(Map<String,Object> map);
}
...@@ -2,11 +2,13 @@ package com.byit.factory; ...@@ -2,11 +2,13 @@ package com.byit.factory;
import cn.hutool.core.collection.CollectionUtil; import cn.hutool.core.collection.CollectionUtil;
import cn.hutool.core.net.NetUtil; import cn.hutool.core.net.NetUtil;
import com.byit.callback.DefaultRemainingOperationsCallBack;
import com.byit.callback.RemainingOperationsCallBack;
import com.byit.registry.DataSourceServiceRegistry; import com.byit.registry.DataSourceServiceRegistry;
import com.byit.registry.PluginServiceRegistry; import com.byit.registry.PluginServiceRegistry;
import com.byit.registry.client.model.RegistryDataParamVO; import com.byit.registry.client.model.RegistryDataParamVO;
import com.byit.server.netty.NettyPluginServer;
import com.byit.server.PluginServer; import com.byit.server.PluginServer;
import com.byit.server.netty.NettyPluginServer;
import com.byit.utils.IpUtil; import com.byit.utils.IpUtil;
import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.StringUtils;
...@@ -22,6 +24,7 @@ public abstract class PluginServerFactory { ...@@ -22,6 +24,7 @@ public abstract class PluginServerFactory {
private final static String DEFAULT_ENV_NAME="plugin-netty"; private final static String DEFAULT_ENV_NAME="plugin-netty";
private final static Class<? extends PluginServiceRegistry> DEFAULT_SERVICE_REGISTRY_CLASS = DataSourceServiceRegistry.class; private final static Class<? extends PluginServiceRegistry> DEFAULT_SERVICE_REGISTRY_CLASS = DataSourceServiceRegistry.class;
private final static Class<? extends PluginServer> DEFAULT_NETTY_PLUGIN_SERVER_CLASS = NettyPluginServer.class; private final static Class<? extends PluginServer> DEFAULT_NETTY_PLUGIN_SERVER_CLASS = NettyPluginServer.class;
private final static Class<? extends RemainingOperationsCallBack> DEFAULT_REMAINING_OPERATIONS_CLASS = DefaultRemainingOperationsCallBack.class;
/** /**
* 任务缓存处理器 * 任务缓存处理器
*/ */
...@@ -41,8 +44,10 @@ public abstract class PluginServerFactory { ...@@ -41,8 +44,10 @@ public abstract class PluginServerFactory {
private int port; private int port;
private Class<? extends PluginServiceRegistry> serviceRegistryClass; private Class<? extends PluginServiceRegistry> serviceRegistryClass;
private Class<? extends PluginServer> pluginServerClass; private Class<? extends PluginServer> pluginServerClass;
private Class<? extends RemainingOperationsCallBack> remainingOperationsCallBack;
private PluginServer pluginServer; private PluginServer pluginServer;
private PluginServiceRegistry pluginServiceRegistry; private PluginServiceRegistry pluginServiceRegistry;
private RemainingOperationsCallBack remainingOperationsCallBackObj;
private Set<String> serverKeys = new HashSet<>(8); private Set<String> serverKeys = new HashSet<>(8);
private List<RegistryDataParamVO> registryDataParamVOs = new ArrayList<>(2); private List<RegistryDataParamVO> registryDataParamVOs = new ArrayList<>(2);
...@@ -52,7 +57,8 @@ public abstract class PluginServerFactory { ...@@ -52,7 +57,8 @@ public abstract class PluginServerFactory {
public void init(int corePoolSize, int maxPoolSize,Integer port, String registryUrl, String env, String biz, public void init(int corePoolSize, int maxPoolSize,Integer port, String registryUrl, String env, String biz,
Class<? extends PluginServiceRegistry> serviceRegistryClass, Class<? extends PluginServiceRegistry> serviceRegistryClass,
Class<? extends PluginServer> pluginServerClass){ Class<? extends PluginServer> pluginServerClass,
Class<? extends RemainingOperationsCallBack> remainingOperationsCallBack){
if(!(corePoolSize>0 && maxPoolSize>0 && maxPoolSize>=corePoolSize)){ if(!(corePoolSize>0 && maxPoolSize>0 && maxPoolSize>=corePoolSize)){
this.corePoolSize = 60; this.corePoolSize = 60;
...@@ -97,6 +103,12 @@ public abstract class PluginServerFactory { ...@@ -97,6 +103,12 @@ public abstract class PluginServerFactory {
this.pluginServerClass = DEFAULT_NETTY_PLUGIN_SERVER_CLASS; this.pluginServerClass = DEFAULT_NETTY_PLUGIN_SERVER_CLASS;
} }
if (remainingOperationsCallBack != null) {
this.remainingOperationsCallBack = remainingOperationsCallBack;
}else{
this.remainingOperationsCallBack = DEFAULT_REMAINING_OPERATIONS_CLASS;
}
if(this.serviceRegistryClass == null){ if(this.serviceRegistryClass == null){
...@@ -106,6 +118,10 @@ public abstract class PluginServerFactory { ...@@ -106,6 +118,10 @@ public abstract class PluginServerFactory {
if(this.pluginServerClass == null){ if(this.pluginServerClass == null){
throw new RuntimeException("插件端使用的服务类型不能为空!"); throw new RuntimeException("插件端使用的服务类型不能为空!");
} }
if (this.remainingOperationsCallBack == null) {
throw new RuntimeException("服务的后续处理不能为空!");
}
} }
public void addService(String key,Object serverBean){ public void addService(String key,Object serverBean){
...@@ -120,14 +136,13 @@ public abstract class PluginServerFactory { ...@@ -120,14 +136,13 @@ public abstract class PluginServerFactory {
public void start() throws Exception { public void start() throws Exception {
pluginServer = pluginServerClass.newInstance(); pluginServer = pluginServerClass.newInstance();
pluginServiceRegistry = serviceRegistryClass.newInstance(); pluginServiceRegistry = serviceRegistryClass.newInstance();
remainingOperationsCallBackObj = remainingOperationsCallBack.newInstance();
//放置启动回调 启动成功后会调用注册服务的方法 //放置启动回调 启动成功后会调用注册服务的方法
pluginServer.setStartPluginCallback(()->{ pluginServer.setStartPluginCallback(()->{
//数据初始化 //数据初始化
pluginServiceRegistry.init(biz,env,registryUrl); pluginServiceRegistry.init(biz,env,registryUrl);
//服务注册 //服务注册
serverKeys.forEach(serverKey ->{ serverKeys.forEach(serverKey -> registryDataParamVOs.add(new RegistryDataParamVO(serverKey,ip+":"+port)));
registryDataParamVOs.add(new RegistryDataParamVO(serverKey,ip+":"+port));
});
pluginServiceRegistry.registry(registryDataParamVOs); pluginServiceRegistry.registry(registryDataParamVOs);
}); });
//设置停止回调 //设置停止回调
...@@ -138,10 +153,19 @@ public abstract class PluginServerFactory { ...@@ -138,10 +153,19 @@ public abstract class PluginServerFactory {
pluginServiceRegistry.stop(); pluginServiceRegistry.stop();
pluginServiceRegistry = null; pluginServiceRegistry = null;
}); });
//调用剩余方法的回调
remainingOperationsCallBackObj.serverRemainingOperations(serverPoll);
//启动服务器 //启动服务器
pluginServer.start(this); pluginServer.start(this);
} }
/**
* 停止方法
*/
public void stop(){
pluginServer.stop();
}
public int getCorePoolSize() { public int getCorePoolSize() {
return corePoolSize; return corePoolSize;
...@@ -255,4 +279,12 @@ public abstract class PluginServerFactory { ...@@ -255,4 +279,12 @@ public abstract class PluginServerFactory {
public void setRegistryDataParamVOs(List<RegistryDataParamVO> registryDataParamVOs) { public void setRegistryDataParamVOs(List<RegistryDataParamVO> registryDataParamVOs) {
this.registryDataParamVOs = registryDataParamVOs; this.registryDataParamVOs = registryDataParamVOs;
} }
public Class<? extends RemainingOperationsCallBack> getRemainingOperationsCallBack() {
return remainingOperationsCallBack;
}
public void setRemainingOperationsCallBack(Class<? extends RemainingOperationsCallBack> remainingOperationsCallBack) {
this.remainingOperationsCallBack = remainingOperationsCallBack;
}
} }
package com.byit.factory; package com.byit.factory;
import cn.hutool.core.collection.CollectionUtil; import cn.hutool.core.collection.CollectionUtil;
import com.byit.callback.RemainingOperationsCallBack;
import com.byit.model.ServerConfigurationModel; import com.byit.model.ServerConfigurationModel;
import com.byit.model.ServiceConfigModel; import com.byit.model.ServiceConfigModel;
import com.byit.registry.PluginServiceRegistry; import com.byit.registry.PluginServiceRegistry;
import com.byit.server.PluginServer; import com.byit.server.PluginServer;
import com.byit.task.annotations.TaskHandler;
import com.byit.utils.XmlParseUtil; import com.byit.utils.XmlParseUtil;
import org.apache.commons.lang3.StringUtils;
import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Set;
/** /**
* main方法启动服务器的实现工厂 * main方法启动服务器的实现工厂
* @author huangfu * @author huangfu
*/ */
public class RpcMainPluginServerFactory extends PluginServerFactory { public class RpcMainPluginServerFactory extends PluginServerFactory {
/**
* 全属性传递注入
* @param corePoolSize 核心线程池大小
* @param maxPoolSize 最大线程池大小
* @param port 端口号
* @param registryUrl 注册中心地址
* @param env 环境标识
* @param biz 工程名称
* @param serverKeys 服务key
* @param serviceRegistryClass 注册中心类
* @param pluginServerClass 插件服务类
*/
public RpcMainPluginServerFactory(int corePoolSize, int maxPoolSize,Integer port, String registryUrl, String env, String biz, Set<String> serverKeys, Class<? extends PluginServiceRegistry> serviceRegistryClass, Class<? extends PluginServer> pluginServerClass) {
init(corePoolSize,maxPoolSize,port,registryUrl,env,biz,serviceRegistryClass,pluginServerClass);
}
/** /**
* 配置文件注入 * 配置文件注入
* @param configLocation * @param configLocation
...@@ -43,19 +28,32 @@ public class RpcMainPluginServerFactory extends PluginServerFactory { ...@@ -43,19 +28,32 @@ public class RpcMainPluginServerFactory extends PluginServerFactory {
ServerConfigurationModel serverConfigurationModel = serviceConfigModel.getServerConfigurationModel(); ServerConfigurationModel serverConfigurationModel = serviceConfigModel.getServerConfigurationModel();
String pluginServiceRegistryStr = serverConfigurationModel.getPluginServiceRegistry(); String pluginServiceRegistryStr = serverConfigurationModel.getPluginServiceRegistry();
String pluginServerClassStr = serverConfigurationModel.getPluginServerClass(); String pluginServerClassStr = serverConfigurationModel.getPluginServerClass();
Class<? extends PluginServiceRegistry> pluginServiceRegistryClass = (Class<? extends PluginServiceRegistry>) Class.forName(pluginServiceRegistryStr); String remainingOperationsClassStr = serverConfigurationModel.getRemainingOperationsClass();
Class<? extends PluginServer> pluginServerClass = (Class<? extends PluginServer>) Class.forName(pluginServerClassStr); Class<? extends PluginServiceRegistry> pluginServiceRegistryClass = null;
Map<String, String> classNames = serviceConfigModel.getClassNames(); if(StringUtils.isNotBlank(pluginServiceRegistryStr)){
pluginServiceRegistryClass = (Class<? extends PluginServiceRegistry>) Class.forName(pluginServiceRegistryStr);
}
Class<? extends PluginServer> pluginServerClass = null;
if(StringUtils.isNotBlank(pluginServerClassStr)){
pluginServerClass = (Class<? extends PluginServer>) Class.forName(pluginServerClassStr);
}
Class<? extends RemainingOperationsCallBack> remainingOperationsClass = null;
if(StringUtils.isNotBlank(remainingOperationsClassStr)){
remainingOperationsClass = (Class<? extends RemainingOperationsCallBack>) Class.forName(remainingOperationsClassStr);
}
List<String> classNames = serviceConfigModel.getClassNames();
super.init(serverConfigurationModel.getCoreSize(),serverConfigurationModel.getMaxSize(), super.init(serverConfigurationModel.getCoreSize(),serverConfigurationModel.getMaxSize(),
serverConfigurationModel.getPort(), serverConfigurationModel.getRegistryUrl(), serverConfigurationModel.getPort(), serverConfigurationModel.getRegistryUrl(),
serverConfigurationModel.getServerEnv(), serverConfigurationModel.getServerBiz(), serverConfigurationModel.getServerEnv(), serverConfigurationModel.getServerBiz(),
pluginServiceRegistryClass,pluginServerClass); pluginServiceRegistryClass,pluginServerClass,remainingOperationsClass);
if (CollectionUtil.isNotEmpty(classNames)) { if (CollectionUtil.isNotEmpty(classNames)) {
classNames.forEach((key,value) ->{ classNames.forEach(className ->{
try { try {
Object o = Class.forName(value).newInstance(); Class<?> aClass = Class.forName(className);
super.addService(key,o); Object o = aClass.newInstance();
TaskHandler annotation = aClass.getAnnotation(TaskHandler.class);
super.addService(annotation.taskName(),o);
} catch (Exception e) { } catch (Exception e) {
e.printStackTrace(); e.printStackTrace();
} }
......
package com.byit.factory;
import com.byit.callback.RemainingOperationsCallBack;
import com.byit.registry.PluginServiceRegistry;
import com.byit.server.PluginServer;
import com.byit.task.annotations.TaskHandler;
import com.byit.task.handler.interfaces.IJobHandler;
import org.springframework.beans.BeansException;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.ApplicationContext;
import org.springframework.context.ApplicationContextAware;
import java.util.Map;
/**
* spring实现
* @author huangfu
*/
public class RpcSpringPluginServerFactory extends PluginServerFactory implements ApplicationContextAware, InitializingBean, DisposableBean {
@Value("${myth.register.url}")
private String address;
@Value("${myth.plugin.biz}")
private String biz;
@Value("${myth.plugin.env}")
private String env;
@Value("${myth.plugin.port}")
private int port;
public RpcSpringPluginServerFactory() {
}
@Override
public void destroy() throws Exception {
super.stop();
}
@Override
public void afterPropertiesSet() throws Exception {
super.init(-1,-1,port,address,env,biz,null,null,null);
super.start();
}
@Override
public void setApplicationContext(ApplicationContext applicationContext) throws BeansException {
Map<String, Object> beansWithAnnotation = applicationContext.getBeansWithAnnotation(TaskHandler.class);
beansWithAnnotation.forEach((key,value) ->{
if(value instanceof IJobHandler){
TaskHandler annotation = value.getClass().getAnnotation(TaskHandler.class);
String taskName = annotation.taskName();
super.addService(taskName,value);
}else{
System.err.println("警告!bean"+key+"不是【com.byit.task.handler.interfaces.IJobHandler】类型!忽略该bean!");
}
});
}
}
...@@ -6,11 +6,12 @@ ...@@ -6,11 +6,12 @@
<server-biz value="myth-job"/> <server-biz value="myth-job"/>
<server-env value="dev"/> <server-env value="dev"/>
<server-thread core-size="30" max-size="100"/> <server-thread core-size="30" max-size="100"/>
<pluginServiceRegistry value="com.byit.registry.DataSourceServiceRegistry"/>
<pluginServerClass value="com.byit.server.netty.MainNettyPluginServer"/> <pluginServerClass value="com.byit.server.netty.MainNettyPluginServer"/>
<!--<pluginServiceRegistry value="com.byit.registry.DataSourceServiceRegistry"/>
<remainingOperationsClass value="com.byit.callback.DefaultRemainingOperationsCallBack"/>-->
</serverConfiguration> </serverConfiguration>
<classNames> <classNames>
<className jobName="sendEmailTest" class="com.byit.server.SendEmailTest"/> <className class="com.byit.server.SendEmailTest"/>
</classNames> </classNames>
</executor-plugin> </executor-plugin>
\ No newline at end of file
package com.byit.job; package com.byit.job;
import com.byit.annotations.JobHandler; import com.byit.task.annotations.JobHandler;
import com.byit.dto.web.ReturnResult; import com.byit.dto.web.ReturnResult;
import com.byit.executor.handler.BaseJobHandler; import com.byit.executor.handler.BaseJobHandler;
......
package com.byit.job; package com.byit.job;
import com.byit.annotations.JobHandler; import com.byit.task.annotations.JobHandler;
import com.byit.dto.web.ReturnResult; import com.byit.dto.web.ReturnResult;
import com.byit.executor.handler.BaseJobHandler; import com.byit.executor.handler.BaseJobHandler;
......
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