mirror of
https://github.com/tiennm99/zfoo.git
synced 2026-10-11 03:13:51 +00:00
fix[zfoo]: not found javax.annotation.meta.When.MAYBE
This commit is contained in:
1 parent
2f8098aa10
commit
94e7e39e9b
12 files changed
+16
-455
No files matched your search
@@ -72,6 +72,14 @@
|
||||
<version>${zfoo.version}</version>
|
||||
</dependency>
|
||||
|
||||
<!--fix-->
|
||||
<!-- Warning:java: 未知的枚举常量 javax.annotation.meta.When.MAYBE 原因: 找不到-->
|
||||
<dependency>
|
||||
<groupId>com.google.code.findbugs</groupId>
|
||||
<artifactId>jsr305</artifactId>
|
||||
<version>3.0.2</version>
|
||||
</dependency>
|
||||
|
||||
<!-- 依赖的通信类库 -->
|
||||
<!-- <dependency>-->
|
||||
<!-- <groupId>io.netty</groupId>-->
|
||||
|
||||
@@ -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之后调用
|
||||
|
||||
@@ -83,6 +83,7 @@ public class JProtobufTcpCodecHandler extends ByteToMessageCodec<EncodedPacketIn
|
||||
|
||||
public void write(ByteBuf buffer, Object packet, Object attachment) throws IOException {
|
||||
// 写入protobuf协议
|
||||
@SuppressWarnings("unchecked")
|
||||
var protobufCodec = (Codec<Object>) ProtobufProxy.create(packet.getClass());
|
||||
byte[] bytes = protobufCodec.encode(packet);
|
||||
// header(4byte) + protocolId(2byte)
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
}
|
||||
@@ -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();
|
||||
}
|
||||
}
|
||||
@@ -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<SocketChannel> {
|
||||
@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());
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
@@ -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 + '\'' +
|
||||
'}';
|
||||
}
|
||||
}
|
||||
@@ -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 + '\'' +
|
||||
'}';
|
||||
}
|
||||
}
|
||||
@@ -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<SocketChannel> {
|
||||
@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());
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
-58
@@ -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);
|
||||
}
|
||||
|
||||
}
|
||||
@@ -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. 食谱; 处方; 烹饪法; 制作法
|
||||
* <p>
|
||||
* 数据发布/订阅(Publish/Subscribe)系统,一般有两种方式推(Push)模式;客户端定时轮询,发现有改变就拉(Pull)。
|
||||
* <p>
|
||||
* 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);
|
||||
}
|
||||
|
||||
}
|
||||
Reference in new issue
Block a user