Commit 02f63551 by huangfusuper

Netty客户端 资源类与资源准备类开发

parent db67db6e
package com.byit.factory;
import com.byit.future.PluginFutureResponse;
import com.byit.param.PluginCallback;
import com.byit.registry.PluginServiceRegistry;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
/**
* 插件端
* @author huangfu
*/
public class PluginClientFactory {
private Class<? extends PluginServiceRegistry> pluginServiceRegistryClass;
private String registryUrl;
private String biz;
private String env;
private PluginServiceRegistry pluginServiceRegistry;
private Map<String, PluginFutureResponse> pluginFutureResponseMap = new ConcurrentHashMap<>(8);
private List<PluginCallback> stopCallbacks = new ArrayList<>(8);
/**
* 停止回调集
* @param pluginCallback
*/
public void addStopCallback(PluginCallback pluginCallback){
this.stopCallbacks.add(pluginCallback);
}
public void setPluginFutureResponseMap(String requestId,PluginFutureResponse pluginFutureResponse){
this.pluginFutureResponseMap.put(requestId,pluginFutureResponse);
}
public PluginFutureResponse getPluginFutureResponseMap(String requestId){
return this.pluginFutureResponseMap.get(requestId);
}
public void removePluginFutureResponseMap(String requestId){
this.pluginFutureResponseMap.remove(requestId);
}
/**
* 注册中心初始化
* @throws Exception
*/
public void start() throws Exception{
if (pluginServiceRegistryClass!=null) {
pluginServiceRegistry = pluginServiceRegistryClass.newInstance();
pluginServiceRegistry.init(biz,env,registryUrl);
}
}
public void close(){
pluginServiceRegistry.stop();
stopCallbacks.forEach(value ->{
try {
value.run();
} catch (Exception e) {
e.printStackTrace();
}
});
}
public PluginClientFactory(Class<? extends PluginServiceRegistry> pluginServiceRegistryClass, String registryUrl, String biz, String env) {
this.pluginServiceRegistryClass = pluginServiceRegistryClass;
this.registryUrl = registryUrl;
this.biz = biz;
this.env = env;
}
public Class<? extends PluginServiceRegistry> getPluginServiceRegistryClass() {
return pluginServiceRegistryClass;
}
public void setPluginServiceRegistryClass(Class<? extends PluginServiceRegistry> pluginServiceRegistryClass) {
this.pluginServiceRegistryClass = pluginServiceRegistryClass;
}
public String getRegistryUrl() {
return registryUrl;
}
public void setRegistryUrl(String registryUrl) {
this.registryUrl = registryUrl;
}
public String getBiz() {
return biz;
}
public void setBiz(String biz) {
this.biz = biz;
}
public String getEnv() {
return env;
}
public void setEnv(String env) {
this.env = env;
}
public PluginServiceRegistry getPluginServiceRegistry() {
return pluginServiceRegistry;
}
public void setPluginServiceRegistry(PluginServiceRegistry pluginServiceRegistry) {
this.pluginServiceRegistry = pluginServiceRegistry;
}
}
package com.byit.init;
import cn.hutool.core.collection.CollectionUtil;
import com.byit.client.PluginClient;
import com.byit.factory.PluginClientFactory;
import com.byit.future.PluginFutureResponse;
import com.byit.packet.request.PluginRpcRequestPacket;
import com.byit.registry.PluginServiceRegistry;
import com.byit.rpc.remoting.invoker.route.LoadBalance;
import org.apache.commons.lang3.StringUtils;
import java.lang.reflect.Proxy;
import java.util.TreeSet;
import java.util.UUID;
/**
* 插件端 客户端的数据初始化
* @author huangfu
*/
public class PluginClientInitialization {
private String callbackUrl;
private long timeout;
private LoadBalance loadBalance;
private PluginClientFactory pluginClientFactory;
private Class<?> iFace;
private Class<? extends PluginClient> pluginClientClass;
private PluginClient pluginClient;
public PluginClientInitialization(String callbackUrl, long timeout, LoadBalance loadBalance, PluginClientFactory pluginClientFactory, Class<?> iFace,Class<? extends PluginClient> pluginClientClass) {
this.callbackUrl = callbackUrl;
this.timeout = timeout;
this.loadBalance = loadBalance;
this.pluginClientFactory = pluginClientFactory;
this.iFace = iFace;
this.pluginClientClass = pluginClientClass;
initClient();
}
private void initClient(){
try {
pluginClient = this.pluginClientClass.newInstance();
pluginClient.init(this);
} catch (InstantiationException | IllegalAccessException e) {
throw new RuntimeException(e);
}
}
/**
* 动态代理 生成动态代理的方法实现 ,并将动态代理生成的实现返还回调用方!
* @return 动态代理生成的实现类
*/
public Object getObject(){
return Proxy.newProxyInstance(Thread.currentThread().getContextClassLoader(),
new Class[]{this.iFace},
(proxy, method, args) ->{
PluginRpcRequestPacket pluginRpcRequestPacket = (PluginRpcRequestPacket)args[0];
if(pluginRpcRequestPacket == null){
throw new RuntimeException("服务参数为null:【com.byit.packet.request.PluginRpcRequestPacket】!");
}
PluginServiceRegistry pluginServiceRegistry = pluginClientFactory.getPluginServiceRegistry();
//服务器的注册ip:port
String address = null;
TreeSet<String> discovery = pluginServiceRegistry.discovery(pluginRpcRequestPacket.getJobName());
if(CollectionUtil.isNotEmpty(discovery)){
if(discovery.size() ==1){
address = discovery.first();
}else{
address = loadBalance.rpcInvokerRouter.route(pluginRpcRequestPacket.getJobName(), discovery);
}
}
if(StringUtils.isBlank(address)){
throw new RuntimeException("该服务在注册表中不存在!");
}
pluginRpcRequestPacket.setRequestId(UUID.randomUUID().toString());
PluginFutureResponse pluginFutureResponse = new PluginFutureResponse(pluginClientFactory,pluginRpcRequestPacket,null);
try {
//发送请求
pluginClient.send(address,pluginRpcRequestPacket);
//获取结果
return pluginFutureResponse.get();
}catch (Exception e){
throw new RuntimeException(e);
}finally {
pluginFutureResponse.remove(pluginRpcRequestPacket.getRequestId());
}
});
}
public String getCallbackUrl() {
return callbackUrl;
}
public void setCallbackUrl(String callbackUrl) {
this.callbackUrl = callbackUrl;
}
public long getTimeout() {
return timeout;
}
public void setTimeout(long timeout) {
this.timeout = timeout;
}
public LoadBalance getLoadBalance() {
return loadBalance;
}
public void setLoadBalance(LoadBalance loadBalance) {
this.loadBalance = loadBalance;
}
public PluginClientFactory getPluginClientFactory() {
return pluginClientFactory;
}
public void setPluginClientFactory(PluginClientFactory pluginClientFactory) {
this.pluginClientFactory = pluginClientFactory;
}
public Class<?> getiFace() {
return iFace;
}
public void setiFace(Class<?> iFace) {
this.iFace = iFace;
}
public Class<? extends PluginClient> getPluginClientClass() {
return pluginClientClass;
}
public void setPluginClientClass(Class<? extends PluginClient> pluginClientClass) {
this.pluginClientClass = pluginClientClass;
}
public PluginClient getPluginClient() {
return pluginClient;
}
public void setPluginClient(PluginClient pluginClient) {
this.pluginClient = pluginClient;
}
}
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