diff --git a/net/pom.xml b/net/pom.xml index 61f8c106..65196d07 100644 --- a/net/pom.xml +++ b/net/pom.xml @@ -72,6 +72,14 @@ ${zfoo.version} + + + + com.google.code.findbugs + jsr305 + 3.0.2 + + diff --git a/net/src/main/java/com/zfoo/net/consumer/Consumer.java b/net/src/main/java/com/zfoo/net/consumer/Consumer.java index d153c4d3..bb3dbe94 100644 --- a/net/src/main/java/com/zfoo/net/consumer/Consumer.java +++ b/net/src/main/java/com/zfoo/net/consumer/Consumer.java @@ -109,6 +109,7 @@ public class Consumer implements IConsumer { if (answerClass != null && answerClass != responsePacket.getClass()) { throw new UnexpectedProtocolException("client expect protocol:[{}], but found protocol:[{}]", answerClass, responsePacket.getClass().getName()); } + @SuppressWarnings("unchecked") var syncAnswer = new SyncAnswer<>((T) responsePacket, clientSignalAttachment); // load balancer之后调用 diff --git a/net/src/main/java/com/zfoo/net/handler/codec/jprotobuf/JProtobufTcpCodecHandler.java b/net/src/main/java/com/zfoo/net/handler/codec/jprotobuf/JProtobufTcpCodecHandler.java index 53d47c40..fc0622e7 100644 --- a/net/src/main/java/com/zfoo/net/handler/codec/jprotobuf/JProtobufTcpCodecHandler.java +++ b/net/src/main/java/com/zfoo/net/handler/codec/jprotobuf/JProtobufTcpCodecHandler.java @@ -83,6 +83,7 @@ public class JProtobufTcpCodecHandler extends ByteToMessageCodec) ProtobufProxy.create(packet.getClass()); byte[] bytes = protobufCodec.encode(packet); // header(4byte) + protocolId(2byte) diff --git a/net/src/main/java/com/zfoo/net/router/Router.java b/net/src/main/java/com/zfoo/net/router/Router.java index a8deb3ee..0b2acf64 100644 --- a/net/src/main/java/com/zfoo/net/router/Router.java +++ b/net/src/main/java/com/zfoo/net/router/Router.java @@ -188,7 +188,9 @@ public class Router implements IRouter { throw new UnexpectedProtocolException("client expect protocol:[{}], but found protocol:[{}]", answerClass, responsePacket.getClass().getName()); } - return new SyncAnswer<>((T) responsePacket, clientSignalAttachment); + @SuppressWarnings("unchecked") + var syncAnswer = new SyncAnswer<>((T) responsePacket, clientSignalAttachment); + return syncAnswer; } catch (TimeoutException e) { throw new NetTimeOutException("syncAsk timeout exception, ask:[{}], attachment:[{}]", JsonUtils.object2String(packet), JsonUtils.object2String(clientSignalAttachment)); } finally { @@ -253,7 +255,9 @@ public class Router implements IRouter { } // 异步返回,回调业务逻辑 - asyncAnswer.setFuturePacket((T) answer); + @SuppressWarnings("unchecked") + var answerPacket = (T) answer; + asyncAnswer.setFuturePacket(answerPacket); asyncAnswer.consume(); } catch (Throwable throwable1) { logger.error("Asynchronous callback method [ask:{}][answer:{}] error", packet.getClass().getSimpleName(), answer.getClass().getSimpleName(), throwable1); diff --git a/net/src/test/java/com/zfoo/net/base/netty/subscribe/ClientHandler.java b/net/src/test/java/com/zfoo/net/base/netty/subscribe/ClientHandler.java deleted file mode 100644 index 55879afe..00000000 --- a/net/src/test/java/com/zfoo/net/base/netty/subscribe/ClientHandler.java +++ /dev/null @@ -1,36 +0,0 @@ -package com.zfoo.net.base.netty.subscribe; - -import io.netty.channel.ChannelHandlerContext; -import io.netty.channel.ChannelInboundHandlerAdapter; - -/** - * @author godotg - * @version 1.0 - * @since 2017 05.23 18:02 - */ -public class ClientHandler extends ChannelInboundHandlerAdapter { - - @Override - public void channelActive(ChannelHandlerContext ctx) throws Exception { - for (int i = 0; i < 10; i++) { - //*****一定要实现Serializable,不实现不会抛异常,很难查找bug - ctx.write(new SubscribeReq(i, "sun", "nettyBook", "1885630", "shanghai")); - } - ctx.flush(); - } - - @Override - public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception { - System.out.println(msg); - } - - @Override - public void channelReadComplete(ChannelHandlerContext ctx) throws Exception { - ctx.flush(); - } - - @Override - public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception { - ctx.close(); - } -} diff --git a/net/src/test/java/com/zfoo/net/base/netty/subscribe/Serverhandler.java b/net/src/test/java/com/zfoo/net/base/netty/subscribe/Serverhandler.java deleted file mode 100644 index 894a5aa4..00000000 --- a/net/src/test/java/com/zfoo/net/base/netty/subscribe/Serverhandler.java +++ /dev/null @@ -1,32 +0,0 @@ -package com.zfoo.net.base.netty.subscribe; - -import io.netty.channel.ChannelHandlerContext; -import io.netty.channel.ChannelInboundHandlerAdapter; - -/** - * @author godotg - * @version 1.0 - * @since 2017 05.23 17:13 - */ -public class Serverhandler extends ChannelInboundHandlerAdapter { - - @Override - public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception { - SubscribeReq req = (SubscribeReq) msg; - System.out.println("Server accept client subscribe req: " + req.toString()); - SubscribeResp resp = new SubscribeResp(req.getReqID(), 0, "subscribe successfully!"); - //*****一定要实现Serializable,不实现不会抛异常,很难查找bug - ctx.writeAndFlush(resp); - } - - @Override - public void channelReadComplete(ChannelHandlerContext ctx) throws Exception { - ctx.flush(); - } - - @Override - public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception { - cause.printStackTrace(); - ctx.close(); - } -} diff --git a/net/src/test/java/com/zfoo/net/base/netty/subscribe/SubscribeClientTest.java b/net/src/test/java/com/zfoo/net/base/netty/subscribe/SubscribeClientTest.java deleted file mode 100644 index 188dae66..00000000 --- a/net/src/test/java/com/zfoo/net/base/netty/subscribe/SubscribeClientTest.java +++ /dev/null @@ -1,61 +0,0 @@ -package com.zfoo.net.base.netty.subscribe; - -import com.zfoo.protocol.util.ThreadUtils; -import io.netty.bootstrap.Bootstrap; -import io.netty.channel.ChannelFuture; -import io.netty.channel.ChannelInitializer; -import io.netty.channel.ChannelOption; -import io.netty.channel.EventLoopGroup; -import io.netty.channel.nio.NioEventLoopGroup; -import io.netty.channel.socket.SocketChannel; -import io.netty.channel.socket.nio.NioSocketChannel; -import io.netty.handler.codec.serialization.ClassResolvers; -import io.netty.handler.codec.serialization.ObjectDecoder; -import io.netty.handler.codec.serialization.ObjectEncoder; -import org.junit.Ignore; -import org.junit.Test; - -/** - * @author godotg - * @version 1.0 - * @since 2017 05.23 18:02 - */ -@Ignore -public class SubscribeClientTest { - - @Test - public void clientTest() { - var client = new SubscribeClientTest(); - client.connect(9999, "127.0.0.1"); - ThreadUtils.sleep(Long.MAX_VALUE); - } - - - public void connect(int port, String host) { - EventLoopGroup group = new NioEventLoopGroup(); - Bootstrap bootstrap = new Bootstrap(); - bootstrap.group(group).channel(NioSocketChannel.class) - .option(ChannelOption.TCP_NODELAY, true) - .handler(new ChildChannelHandler()); - ChannelFuture future = null; - try { - future = bootstrap.connect(host, port).sync(); - future.channel().closeFuture().sync(); - } catch (InterruptedException e) { - throw new RuntimeException(e); - } finally { - group.shutdownGracefully(); - } - } - - private class ChildChannelHandler extends ChannelInitializer { - @Override - protected void initChannel(SocketChannel channel) throws Exception { - channel.pipeline().addLast(new ObjectDecoder(1024 * 1024 - , ClassResolvers.cacheDisabled(this.getClass().getClassLoader()))); - channel.pipeline().addLast(new ObjectEncoder()); - channel.pipeline().addLast(new ClientHandler()); - } - } - -} diff --git a/net/src/test/java/com/zfoo/net/base/netty/subscribe/SubscribeReq.java b/net/src/test/java/com/zfoo/net/base/netty/subscribe/SubscribeReq.java deleted file mode 100644 index 680991d8..00000000 --- a/net/src/test/java/com/zfoo/net/base/netty/subscribe/SubscribeReq.java +++ /dev/null @@ -1,82 +0,0 @@ -package com.zfoo.net.base.netty.subscribe; - -import java.io.Serializable; - -/** - * @author godotg - * @version 1.0 - * @since 2017 05.23 16:51 - */ -public class SubscribeReq implements Serializable { - - private static final long serialVersionUID = 1L; - - private int reqID; - - private String userName; - - private String productName; - - private String phoneNumber; - - private String address; - - public SubscribeReq(int reqID, String userName, String productName, String phoneNumber, String address) { - this.reqID = reqID; - this.userName = userName; - this.productName = productName; - this.phoneNumber = phoneNumber; - this.address = address; - } - - public int getReqID() { - return reqID; - } - - public void setReqID(int reqID) { - this.reqID = reqID; - } - - public String getUserName() { - return userName; - } - - public void setUserName(String userName) { - this.userName = userName; - } - - public String getProductName() { - return productName; - } - - public void setProductName(String productName) { - this.productName = productName; - } - - public String getPhoneNumber() { - return phoneNumber; - } - - public void setPhoneNumber(String phoneNumber) { - this.phoneNumber = phoneNumber; - } - - public String getAddress() { - return address; - } - - public void setAddress(String address) { - this.address = address; - } - - @Override - public String toString() { - return "SubscribeReq{" + - "reqID=" + reqID + - ", userName='" + userName + '\'' + - ", productName='" + productName + '\'' + - ", phoneNumber='" + phoneNumber + '\'' + - ", address='" + address + '\'' + - '}'; - } -} diff --git a/net/src/test/java/com/zfoo/net/base/netty/subscribe/SubscribeResp.java b/net/src/test/java/com/zfoo/net/base/netty/subscribe/SubscribeResp.java deleted file mode 100644 index c2bce7b2..00000000 --- a/net/src/test/java/com/zfoo/net/base/netty/subscribe/SubscribeResp.java +++ /dev/null @@ -1,58 +0,0 @@ -package com.zfoo.net.base.netty.subscribe; - -import java.io.Serializable; - -/** - * @author godotg - * @version 1.0 - * @since 2017 05.23 16:52 - */ -public class SubscribeResp implements Serializable { - - private static final long serialVersionUID = 1L; - - private int respID; - - private int respCode; - - private String desc; - - public SubscribeResp(int respID, int respCode, String desc) { - this.respID = respID; - this.respCode = respCode; - this.desc = desc; - } - - public int getRespID() { - return respID; - } - - public void setRespID(int respID) { - this.respID = respID; - } - - public int getRespCode() { - return respCode; - } - - public void setRespCode(int respCode) { - this.respCode = respCode; - } - - public String getDesc() { - return desc; - } - - public void setDesc(String desc) { - this.desc = desc; - } - - @Override - public String toString() { - return "SubscribeResp{" + - "respID=" + respID + - ", respCode=" + respCode + - ", desc='" + desc + '\'' + - '}'; - } -} diff --git a/net/src/test/java/com/zfoo/net/base/netty/subscribe/SubscribeServerTest.java b/net/src/test/java/com/zfoo/net/base/netty/subscribe/SubscribeServerTest.java deleted file mode 100644 index 0ac15260..00000000 --- a/net/src/test/java/com/zfoo/net/base/netty/subscribe/SubscribeServerTest.java +++ /dev/null @@ -1,71 +0,0 @@ -package com.zfoo.net.base.netty.subscribe; - -import com.zfoo.protocol.util.ThreadUtils; -import io.netty.bootstrap.ServerBootstrap; -import io.netty.channel.ChannelFuture; -import io.netty.channel.ChannelInitializer; -import io.netty.channel.ChannelOption; -import io.netty.channel.EventLoopGroup; -import io.netty.channel.nio.NioEventLoopGroup; -import io.netty.channel.socket.SocketChannel; -import io.netty.channel.socket.nio.NioServerSocketChannel; -import io.netty.handler.codec.serialization.ClassResolvers; -import io.netty.handler.codec.serialization.ObjectDecoder; -import io.netty.handler.codec.serialization.ObjectEncoder; -import io.netty.handler.logging.LogLevel; -import io.netty.handler.logging.LoggingHandler; -import org.junit.Ignore; -import org.junit.Test; - -/** - * 对象读写序列化通信 - * - * @author godotg - * @version 1.0 - * @since 2017 05.23 16:51 - */ -@Ignore -public class SubscribeServerTest { - - @Test - public void serverTest() { - var server = new SubscribeServerTest(); - server.init(); - ThreadUtils.sleep(Long.MAX_VALUE); - } - - public void init() { - //配置服务端nio线程组 - EventLoopGroup bossGroup = new NioEventLoopGroup();//服务端接受客户端连接 - EventLoopGroup workerGroup = new NioEventLoopGroup();//SocketChannel的网络读写 - try { - ServerBootstrap bootstrap = new ServerBootstrap(); - bootstrap.group(bossGroup, workerGroup).channel(NioServerSocketChannel.class) - .option(ChannelOption.SO_BACKLOG, 100).handler(new LoggingHandler(LogLevel.INFO)) - .childHandler(new ChildChannelHandler()); - - //绑定端口,同步等待成功 - ChannelFuture future = bootstrap.bind(9999).sync(); - //等待服务端监听端口关闭 - future.channel().closeFuture().sync(); - } catch (InterruptedException e) { - e.printStackTrace(); - } finally { - //优雅的退出,释放线程池资源 - bossGroup.shutdownGracefully(); - workerGroup.shutdownGracefully(); - } - } - - private class ChildChannelHandler extends ChannelInitializer { - @Override - protected void initChannel(SocketChannel channel) throws Exception { - channel.pipeline().addLast(new ObjectDecoder(1024 * 1024 - , ClassResolvers.weakCachingConcurrentResolver(this.getClass().getClassLoader()))); - //*****一定要实现Serializable,不实现不会抛异常,很难查找bug - channel.pipeline().addLast(new ObjectEncoder()); - channel.pipeline().addLast(new Serverhandler()); - } - } - -} diff --git a/net/src/test/java/com/zfoo/net/zookeeper/recipes/nodecache/NodeChildrenListenerTest.java b/net/src/test/java/com/zfoo/net/zookeeper/recipes/nodecache/NodeChildrenListenerTest.java deleted file mode 100644 index 672936e6..00000000 --- a/net/src/test/java/com/zfoo/net/zookeeper/recipes/nodecache/NodeChildrenListenerTest.java +++ /dev/null @@ -1,58 +0,0 @@ -package com.zfoo.net.zookeeper.recipes.nodecache; - - -import org.apache.curator.RetryPolicy; -import org.apache.curator.framework.CuratorFramework; -import org.apache.curator.framework.CuratorFrameworkFactory; -import org.apache.curator.framework.recipes.cache.PathChildrenCache; -import org.apache.curator.framework.recipes.cache.PathChildrenCacheEvent; -import org.apache.curator.framework.recipes.cache.PathChildrenCacheListener; -import org.apache.curator.retry.RetryUntilElapsed; -import org.junit.Ignore; -import org.junit.Test; - -@Ignore -public class NodeChildrenListenerTest { - - @Test - public void test() throws Exception { - RetryPolicy retryPolicy = new RetryUntilElapsed(5000, 1000); - - CuratorFramework client = CuratorFrameworkFactory - .builder() - .connectString("localhost:2181") - .sessionTimeoutMs(5000) - .connectionTimeoutMs(5000) - .retryPolicy(retryPolicy) - .build(); - - client.start(); - - // PathChildrenCache用于监听指定Zookeeper数据节点的子节点变换的情况。 - // 注意:Zookeeper无法对二级子节点监听 - final PathChildrenCache cache = new PathChildrenCache(client, "/node_test", true); - cache.start(); - cache.getListenable().addListener(new PathChildrenCacheListener() { - - @Override - public void childEvent(CuratorFramework client, PathChildrenCacheEvent event) throws Exception { - switch (event.getType()) { - case CHILD_ADDED: - System.out.println("CHILD_ADDED:" + event.getData()); - break; - case CHILD_UPDATED: - System.out.println("CHILD_UPDATED:" + event.getData()); - break; - case CHILD_REMOVED: - System.out.println("CHILD_REMOVED:" + event.getData()); - break; - default: - break; - } - } - }); - - Thread.sleep(Integer.MAX_VALUE); - } - -} diff --git a/net/src/test/java/com/zfoo/net/zookeeper/recipes/nodecache/NodeListenerTest.java b/net/src/test/java/com/zfoo/net/zookeeper/recipes/nodecache/NodeListenerTest.java deleted file mode 100644 index 183d5366..00000000 --- a/net/src/test/java/com/zfoo/net/zookeeper/recipes/nodecache/NodeListenerTest.java +++ /dev/null @@ -1,55 +0,0 @@ -package com.zfoo.net.zookeeper.recipes.nodecache; - - -import org.apache.curator.RetryPolicy; -import org.apache.curator.framework.CuratorFramework; -import org.apache.curator.framework.CuratorFrameworkFactory; -import org.apache.curator.framework.recipes.cache.NodeCache; -import org.apache.curator.framework.recipes.cache.NodeCacheListener; -import org.apache.curator.retry.RetryUntilElapsed; -import org.junit.Ignore; -import org.junit.Test; - - -/** - * recipe [rec·i·pe || 'resɪpɪ] n. 食谱; 处方; 烹饪法; 制作法 - *

- * 数据发布/订阅(Publish/Subscribe)系统,一般有两种方式推(Push)模式;客户端定时轮询,发现有改变就拉(Pull)。 - *

- * Zookeeper采用两者相结合,客户端订阅,节点有变化服务器就推送 - */ -@Ignore -public class NodeListenerTest { - - @Test - public void test() throws Exception { - RetryPolicy retryPolicy = new RetryUntilElapsed(5000, 1000); - - CuratorFramework curator = CuratorFrameworkFactory - .builder() - .connectString("localhost:2181") - .sessionTimeoutMs(5000) - .connectionTimeoutMs(5000) - .retryPolicy(retryPolicy) - .build(); - - curator.start(); - - // Cache是Curator中对事件监听的包装,其对事件的监听其实可以看做本地缓存视图和远程Zookeeper视图的对比过程。 - // 同时Recipes能够反复注册监听,从而大大简化了原生API开发的繁琐程度。 - final NodeCache cache = new NodeCache(curator, "/node_test"); - cache.start(); - - // 节点创建,数据改变都能检测到;如果节点被删除则不能检测。 - cache.getListenable().addListener(new NodeCacheListener() { - @Override - public void nodeChanged() throws Exception { - byte[] ret = cache.getCurrentData().getData(); - System.out.println("new data:" + new String(ret)); - } - }); - - Thread.sleep(Integer.MAX_VALUE); - } - -}