Commit 0eacacfe by huangfusuper

Netty 客户端开发与实现

parent 39bc2c5a
package com.byit.client.handler;
import com.byit.future.PluginFutureResponse;
import com.byit.init.PluginClientInitialization;
import com.byit.packet.response.PluginRpcResponsePacket;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
/**
* Netty客户端业务处理类
* @author huangfu
*/
public class NettyClientHandler extends SimpleChannelInboundHandler<PluginRpcResponsePacket> {
private PluginClientInitialization pluginClientInitialization;
private PluginConnectClient pluginConnectClient;
public NettyClientHandler(PluginClientInitialization pluginClientInitialization, PluginConnectClient pluginConnectClient) {
this.pluginClientInitialization = pluginClientInitialization;
this.pluginConnectClient = pluginConnectClient;
}
@Override
protected void channelRead0(ChannelHandlerContext ctx, PluginRpcResponsePacket msg) throws Exception {
PluginFutureResponse pluginFutureResponseMap = pluginClientInitialization.getPluginClientFactory().getPluginFutureResponseMap(msg.getRequestId());
pluginFutureResponseMap.setRpcResponsePacket(msg);
}
}
package com.byit.client.handler;
import com.byit.handler.PackerSpliterHandler;
import com.byit.handler.PacketDecodeHandler;
import com.byit.handler.PacketEncodeHandler;
import com.byit.init.PluginClientInitialization;
import com.byit.packet.request.PluginRpcRequestPacket;
import com.byit.utils.IpUtil;
import io.netty.bootstrap.Bootstrap;
import io.netty.channel.*;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioSocketChannel;
/**
* netty客户端
* @author huangfu
*/
public class NettyPluginConnectionClient extends PluginConnectClient {
private EventLoopGroup group;
private Channel channel;
@Override
public void init(String address, PluginClientInitialization pluginClientInitialization) throws Exception {
//将这个传递到 handler的业务处理器 未来增加心跳检测是 需要调用这个发送心跳
final PluginConnectClient thisClient = this;
Object[] arrays = IpUtil.parseIpPort(address);
String ip = (String)arrays[0];
int port = (int)arrays[1];
this.group = new NioEventLoopGroup();
Bootstrap bootstrap = new Bootstrap();
bootstrap.group(group)
.channel(NioSocketChannel.class)
.option(ChannelOption.TCP_NODELAY, true)
.option(ChannelOption.SO_KEEPALIVE, true)
.option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 10000)
.handler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) throws Exception {
ChannelPipeline pipeline = ch.pipeline();
pipeline.addLast("packerSpliterHandler",new PackerSpliterHandler());
pipeline.addLast("packetDecodeHandler",new PacketDecodeHandler());
pipeline.addLast("nettyClientHandler",new NettyClientHandler(pluginClientInitialization,thisClient));
pipeline.addLast("packetEncodeHandler",new PacketEncodeHandler());
}
});
this.channel = bootstrap.connect(ip, port).sync().channel();
// valid
if (!isValidate()) {
close();
return;
}
}
@Override
public void close() {
if(this.channel !=null && isValidate()){
this.channel.close();
}
if(this.group != null && !group.isShutdown()){
group.shutdownGracefully();
}
}
@Override
public boolean isValidate() {
if (this.channel != null){
return this.channel.isActive();
}
return false;
}
@Override
public void send(PluginRpcRequestPacket pluginRpcRequestPacket) throws Exception {
this.channel.writeAndFlush(pluginRpcRequestPacket);
}
}
package com.byit.client.handler;
import cn.hutool.core.collection.CollectionUtil;
import com.byit.init.PluginClientInitialization;
import com.byit.packet.request.PluginRpcRequestPacket;
import com.byit.rpc.remoting.invoker.RpcInvokerFactory;
import com.byit.rpc.remoting.invoker.reference.RpcReferenceBean;
import com.byit.rpc.remoting.net.common.ConnectClient;
import com.byit.rpc.remoting.net.params.RpcRequest;
import com.byit.rpc.serialize.Serializer;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
/**
* 插件端 客户端 连接对象创建的基类
* @author huangfu
*/
public abstract class PluginConnectClient {
/**
* 初始化方法 将nETTY的客户端初始化完成
* @param address
* @param pluginClientInitialization
* @throws Exception
*/
public abstract void init(String address, final PluginClientInitialization pluginClientInitialization) throws Exception;
/**
* 关闭方法 关闭netty的客户端
*/
public abstract void close();
/**
* 校验客户端通道是否处于活跃状态
* @return
*/
public abstract boolean isValidate();
/**
* 发送数据的方法
* @param pluginRpcRequestPacket
* @throws Exception
*/
public abstract void send(PluginRpcRequestPacket pluginRpcRequestPacket) throws Exception ;
/**
* 发送的委托方法
* @param pluginRpcRequestPacket
* @param address
* @param connectClientImpl
* @param pluginClientInitialization
* @throws Exception
*/
public static void send(PluginRpcRequestPacket pluginRpcRequestPacket, String address,
Class<? extends PluginConnectClient> connectClientImpl,
final PluginClientInitialization pluginClientInitialization) throws Exception {
PluginConnectClient pluginConnectClient = PluginConnectClient.getPool(address, connectClientImpl, pluginClientInitialization);
pluginConnectClient.send(pluginRpcRequestPacket);
}
/**
* 存储所有的客户端连接对象 避免反复创建连接的消耗
*/
private static volatile ConcurrentMap<String, PluginConnectClient> connectClientMap;
/**
* 获取对象连接的锁对象 每一个地址对应一把锁
*/
private static volatile ConcurrentMap<String, Object> connectClientLockMap = new ConcurrentHashMap<>();
private static PluginConnectClient getPool(String address, Class<? extends PluginConnectClient> connectClientImpl,
final PluginClientInitialization pluginClientInitialization) throws Exception {
//初始化连接池 为了避免重复初始化 需要判断
if (connectClientMap == null) {
synchronized (PluginConnectClient.class){
if (connectClientMap == null) {
connectClientMap = new ConcurrentHashMap<>(8);
pluginClientInitialization.getPluginClientFactory().addStopCallback(() ->{
if (CollectionUtil.isNotEmpty(connectClientMap)) {
connectClientMap.forEach((key,value) ->{
value.close();
});
connectClientMap.clear();
}
});
}
}
}
//如果不为null 则将该链接获取出来并判断是否是有效连接
PluginConnectClient pluginConnectClient = connectClientMap.get(address);
if (pluginConnectClient!=null && pluginConnectClient.isValidate()) {
return pluginConnectClient;
}
//先获取该对象对应的锁对象
Object lock = connectClientLockMap.get(address);
if(lock == null){
connectClientLockMap.putIfAbsent(address,new Object());
lock = connectClientLockMap.get(address);
}
//如果该连接无效或者为null 则需要将该链接删除或者设置
synchronized (lock){
pluginConnectClient = connectClientMap.get(address);
if (pluginConnectClient!=null && pluginConnectClient.isValidate()) {
return pluginConnectClient;
}
if(pluginConnectClient != null){
connectClientMap.remove(address);
}
//实例化连接对象
PluginConnectClient newPluginConnectClient = connectClientImpl.newInstance();
newPluginConnectClient.init(address,pluginClientInitialization);
//重新在连接池里设置
connectClientLockMap.put(address,newPluginConnectClient);
return newPluginConnectClient;
}
}
}
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