第04篇:手写JavaRPC框架之搞定网络通信

作者: 西魏陶渊明
博客: https://blog.springlearn.cn/
天下代码一大抄, 抄来抄去有提高, 看你会抄不会抄!

一、前言
- 听说你Sql写的很溜,那么你知道服务端的Sql是如何被传输到SQL服务器执行的吗?
- 听说你10分钟能写两个接口,那你知道数据是如何在两个系统中通讯的吗?
- 听说你微服务玩的很熟练,那你知道微服务的基础是什么吗?
可以这样说,我们写的任何系统都离不开通讯,离不开网络编程,就没有现在我们发达的互联网世界。就没有什么分布式,没有什么微服务。所以由此可见网络编程是非常基础的知识。
但是我们思考下,有多少同学真正使用过Java网络通信的API了呢? 相信百分之80的小伙伴可能都没用过? 为什么呢? 因为我们站在巨人的肩膀上, 底层的代码都被层层的封装起来了,为了使我们能专注于业务的开发。这虽然提高了我们的开发效率。但是呢? 从另一个方面讲,他不利于我们的技术成长,使我们只会用,而不去思考为什么这么用。逐渐沦为CRUD Body。
个人如果想成长,想打破这种现状, 那么网络通信是一定要掌握的,当你掌握了这些,才算掌握了一点核心技术。当你掌握了这些,才能收获一些不一样的东西,看问题的维度又会有所提升。
本系列文章, 我们会一起来写RPC框架,而网络通讯是必要要掌握的知识,如果说以前你不懂,那么没关系跟着小编来Coding。我们一起来从0到1搭建一个网络通信层,然后以此为基础实现一个Java RPC框架吧。
二、目标
2.1 目标介绍

本篇文章是我们的第四篇,内容主要是实现网络通讯。通信层框架主要使用的是Netty进行实现, 说到Netty可能很多同学都没有用过。而要想实现通讯Netty就必须要知道,所以本篇内容篇幅较多。
- 第一个目标,快速学习Netty的架构,掌握Netty 核心的API,最终唯我所用。
- 第二个目标,使用Netty完成我们的通信层。
内容非常的硬核,难度指数比较大,也主要是偏向于实战。请专注学习,内容中如果有差错,欢迎提出。小编会积极回复,并改正。
三、Netty API 学习

前面第三篇我们学习了搞定序列化,在上一篇中我们介绍了这幅图,序列化就是将数据转换成二进制数据,在网络管道中传输。今天开头还是这一张图, 不过今天要说的不在是里面的数据,而是要研究下如何来构建一个通信的管道。本篇文章我们要利用Netty搭建一个网络管道。

这张图是对第一张图的一个细化,可以看到在这个管道的里面有一个双向的数据流【双工】。
客服端向服务端发送数据,服务端也可以同时向客户端发送数据。这个过程叫做全双工。为什么呢? 因为这个管道是TCP管道。我们所知道的Dubbo也是在此基础上实现的。所以说dubbo和http是平级的关系。
- 一个inbound入栈方向,负责将二进制数据转换成Java对象
- 一个outbound出库方向,负责将Java对象转换成二进制对象
3.1 ChannelPipeline 网络管道
上面的那个管道在Netty中就是 ChannelPipeline, ChannelPipeline 是Netty 中一个非常重要的组件,我们说的管道,就可以理解成是这个类,在这个管道中有两个方向的流向。如下
- ChannelInboundHandler 入栈
- ChannelOutboundHandler 出栈
只有管道还不行,要对管道中流动的数据进行处理。怎么来处理呢? 就是 ChannelHandler
3.2 ChannelHandler 管道处理器
ChannelHandle 通道处理器是最顶层接口, ChannelHandler 和 ChannelPipeline 的关系,好比这张图。
ChannelPipeline 相当于是管道,而 ChannelHandler 相当于管道中的每个拦路虎, 通过对数据的拦截,然后进行处理。下面这张图比较形象。

Netty 要想学的好, ChannelHandler 学习不能少,下面是 Netty 中 ChannelHandler 的类关系图。想要处理数据只用继承这其中的一些类就可以了。

请记住这张图,我们下面会利用这些管道处理器来实现我们的网络通信。
3.3 入栈管道解码器
编码器本质上就是一个
ChannelHandler, 所以上图我们也能看出来它是实现了ChannelHandler的。
二进制数据转换成Java对象就要我们来搞一个入栈的解码器,通过上面的图我们知道Netty给我们提供了一个
入栈方向的类。ByteToMessageDecoder。那么我们就先实现他,直接看代码。
1/**
2 * 请求解码器负责将二进制数据转换成能处理的协议
3 * 个人博客:https://java.springlearn.cn
4 * 公众号:西魏陶渊明 {关注获取学习源码}
5 */
6@Slf4j
7public abstract class ChannelDecoder extends ByteToMessageDecoder {
8
9 /**
10 * 解码方法
11 *
12 * @param ctx 通道上下文信息
13 * @param in 网络传过来的信息(注意粘包和拆包问题)
14 * @param out in中的数据转换成对象调用out.add方法
15 * @throws Exception 未知异常
16 */
17 @Override
18 public void decode(ChannelHandlerContext ctx, ByteBuf in, List<Object> out) throws Exception {
19 doDecode(ctx, in, out);
20}
21
22
23 /**
24 * 解码方法
25 *
26 * @param ctx 通道上下文信息
27 * @param in 网络传过来的信息(注意粘包和拆包问题)
28 * @param out in中的数据转换成对象调用out.add方法
29 * @throws Exception 未知异常
30 */
31 protected abstract void doDecode(ChannelHandlerContext ctx, ByteBuf in, List<Object> out) throws Exception;
32}-
doDecode 方法,通过读取ByteBuf中数据,然后经过的规则处理,就是协议处理,然后反序列化成Java对象,然后调用out.add()。
-
这个规则处理,就是协议,我们在本系列文章的第二篇,就说了我们的协议是什么,如下这张图。那么我们就按照这个规则来解析数据吧。

- 实际上这里我们还要面对黏包和拆包的问题。什么是黏包和拆包呢?
3.4 黏包和拆包及解决方案
我们举一个例子,前两天某某国前首相安老三,遇到刺客了。这时候你很悲伤想发一个说说: 只要人没事就好。
拆包:
只要人没事就好 = 只要人没 + 事就好。

就是形容一条完整的数据报文,因为某些原因,数据被分成多段进行传输了,当你读取数据的时候,读到的不是完整的数据,而是一半的数据。此时意思就比较尴尬了。只要人没,事就好 😂。
黏包:

两条报文,连在一起发送了。导致了意思大变样。
就是数据都在网络管道中传输,但是我们服务定位每个报文的长度,导致了读取的数据就是混乱的。
那么拆包和黏包的问题我们都知道了,下面直接说解决方案吧。在Netty中有如下解决方案。

- LineBasedFrameDecoder
遇到了换行符,就当做是一条完整的消息
1 @Test
2 @DisplayName("使用换行符分隔符")
3 void lineBasedFrameDecoder() {
4 int maxLength = 100;
5ByteBuf buffer = ByteBufAllocator.DEFAULT.buffer();
6buffer.writeBytes("hello world
7hello
8world
9".getBytes(StandardCharsets.UTF_8));
10EmbeddedChannel channel = new EmbeddedChannel(new LoggingHandler(LogLevel.DEBUG),
11 new LineBasedFrameDecoder(maxLength),
12 new ChannelInboundHandlerAdapter() {
13 @Override
14 public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
15 ByteBuf buf = ((ByteBuf) msg);
16 String content = buf.toString(StandardCharsets.UTF_8);
17 System.out.println(content);
18 }
19 });
20// hello world
21// hello
22// world
23 channel.writeInbound(ByteBufAllocator.DEFAULT.buffer().writeBytes(buffer));
24}- DelimiterBasedFrameDecoder
遇到了分隔符,就当做是一条完整的消息,分隔符可以自定义。我们可以指定多个分隔符,如下示例。
1 @Test
2 @DisplayName("自定义换行符分隔符")
3 void delimiterBasedFrameDecoder() {
4 int maxLength = 100;
5ByteBuf buffer = ByteBufAllocator.DEFAULT.buffer();
6buffer.writeBytes("hello world
7hello
8world".getBytes(StandardCharsets.UTF_8));
9ByteBuf delimeter1 = Unpooled.buffer().writeBytes("
10".getBytes(StandardCharsets.UTF_8));
11ByteBuf delimeter2 = Unpooled.buffer().writeBytes("
12".getBytes(StandardCharsets.UTF_8));
13ByteBuf delimeter3 = Unpooled.buffer().writeBytes("".getBytes(StandardCharsets.UTF_8));
14EmbeddedChannel channel = new EmbeddedChannel(new LoggingHandler(LogLevel.DEBUG),
15 new DelimiterBasedFrameDecoder(maxLength, delimeter1, delimeter2, delimeter3),
16 new ChannelInboundHandlerAdapter() {
17 @Override
18 public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
19 ByteBuf buf = ((ByteBuf) msg);
20 String content = buf.toString(StandardCharsets.UTF_8);
21 System.out.println(content);
22 }
23 });
24// hello world
25// hello
26// world
27 channel.writeInbound(ByteBufAllocator.DEFAULT.buffer().writeBytes(buffer));
28}- FixedLengthFrameDecoder
固定长度对消息进行拆分
1 @Test
2 @DisplayName("固定长度进行拆解")
3 void fixedLengthFrameDecoder() {
4 //这里每条消息设置的固定长度是5
5 EmbeddedChannel channel = new EmbeddedChannel(new LoggingHandler(LogLevel.DEBUG), new FixedLengthFrameDecoder(5),
6 new ChannelInboundHandlerAdapter() {
7 @Override
8 public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
9 ByteBuf buf = ((ByteBuf) msg);
10 String content = buf.toString(StandardCharsets.UTF_8);
11 System.out.println(content);
12 }
13 });
14// hello
15// world
16// welco
17 channel.writeInbound(ByteBufAllocator.DEFAULT.buffer().writeBytes("helloworldwelcome".getBytes(StandardCharsets.UTF_8)));
18}- LengthFieldBasedFrameDecoder
消息分为两部分,一部分为消息头部,一部分为实际的消息体。其中消息头部是固定长度的,消息体是可变的,且消息头部一般会包含一个Length字段。
1 @Test
2 @DisplayName("动态获取长度报文")
3 void lengthFieldBasedFrameDecoder() {
4 ByteBuf buffer = ByteBufAllocator.DEFAULT.buffer();
5byte[] bytes = "hello world".getBytes(StandardCharsets.UTF_8);
6// 11
7 System.out.println(bytes.length);
8// 4字节
9 buffer.writeInt(bytes.length);
10// 真正的数据
11 buffer.writeBytes(bytes);
12// 最大包长100字节
13 int maxFrameLength = 100;
14// 从0开始,说明开头就是长度
15 int lengthFieldOffset = 0;
16// 0 说明, 报文是有长度+真实数据组成的,没有其他的东西
17 int lengthAdjustment = 0;
18// 跳过长度的字节,因为是int,所以是4字节
19 int initialBytesToStrip = 4;
20EmbeddedChannel channel = new EmbeddedChannel(new LoggingHandler(LogLevel.DEBUG),
21 new LengthFieldBasedFrameDecoder(maxFrameLength, lengthFieldOffset, 4, lengthAdjustment, initialBytesToStrip),
22 new ChannelInboundHandlerAdapter() {
23 @Override
24 public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
25 ByteBuf buf = ((ByteBuf) msg);
26 String content = buf.toString(StandardCharsets.UTF_8);
27 System.out.println(content);
28 }
29 });
30channel.writeInbound(ByteBufAllocator.DEFAULT.buffer().writeBytes(buffer));
31}好了,其实知道了问题产生的原因和已有的解决方案,我们也可以自己重新实现一个拆包和解包的方案,本篇文章我们就会自己来实现一个,具体的思路就是如下。

看起来思路是不是很简单? 底层的API可不简单哦,记得自己看代码,建议拉下来自己走走。
3.5 入栈管道处理器
二进制数据经过前面的解码器,就会转换成Object对象。此时我们下一个处理器就可以直接处理这个Object对象了。此时我们可以来继承 SimpleChannelInboundHandler。自定义一个泛型。如下示例。我们演示下二进制数据转Java对象,并传给我们的业务处理器。
- fillProtocol 首先我们按照我们定义的规则来,生成二进制数据流。
- 然后解析成Java对象,并传给我们的处理器。
1 /**
2 * 请求解码器负责将二进制数据转换成能处理的协议
3 * 虫洞栈:https://java.springlearn.cn
4 * 公众号:西魏陶渊明 {关注获取学习源码}
5 */
6 private ByteBuf fillProtocol() throws Exception {
7 ByteBuf buffer = ByteBufAllocator.DEFAULT.buffer();
8RpcRequest rpcRequest = new RpcRequest();
9//1. 获取协议类型(1个字节)
10 buffer.writeByte(rpcRequest.getProtocolType());
11//2. 获取序列化类型(1个字节)
12 buffer.writeByte(rpcRequest.getSerializationType());
13//3. 根据序列化类型找到数据转换器生成二进制数据
14 Serializer serializer = SerializeEnum.
15 ofByType(rpcRequest.getSerializationType())
16 .getSerialize().newInstance();
17byte[] data = serializer.serialize(rpcRequest);
18//4. 写入报文长度(4个字节)
19 buffer.writeInt(data.length);
20//5. 写入报文内容(数组)
21 buffer.writeBytes(data);
22return buffer;
23}
24
25 @Test
26 @DisplayName("SimpleChannelInboundHandler自动匹配支持的Java对象")
27 public void test() throws Exception {
28 // 根据自定义协议生成二进制数据流
29 ByteBuf byteBuf = fillProtocol();
30EmbeddedChannel channel = new EmbeddedChannel(new LoggingHandler(LogLevel.DEBUG),
31 new ChannelInitializer<EmbeddedChannel>() {
32 @Override
33 protected void initChannel(EmbeddedChannel ch) throws Exception {
34 ChannelPipeline pipeline = ch.pipeline();
35 pipeline.addLast("a handler", new MojitoChannelDecoder("mojito"));
36 pipeline.addLast("b handler",// 自定义一个String类型的
37 new SimpleChannelInboundHandler<RpcRequest>() {
38 @Override
39 protected void channelRead0(ChannelHandlerContext ctx, RpcRequest msg) throws Exception {
40 System.out.println("RpcRequest:" + msg);
41 // 向下传播
42 ctx.fireChannelRead(msg);
43 }
44 });
45 pipeline.addLast("c handler", new SimpleChannelInboundHandler<Integer>() {
46 @Override
47 protected void channelRead0(ChannelHandlerContext ctx, Integer msg) throws Exception {
48 System.out.println("Integer:" + msg);
49 }
50 });
51 }
52 });
53channel.writeInbound(ByteBufAllocator.DEFAULT.buffer().writeBytes(byteBuf));
54}- b handler 的泛型是RpcRequest,二进制数据经过MojitoChannelDecoder将数据转换成RpcRequest对象,此时就会进到b handler。
- 而在b handler中我们继续调用方法向下传播数据。
ctx.fireChannelRead(msg)。会发现c handler并没有执行,为什么呢? 因为SimpleChannelInboundHandler有一个特性,就是只有数据类型为自己定义的泛型的时候才会进入。
如下源码也比较简单,这个特性我们可以抄一下, 可以用到我们需要的地方。
TypeParameterMatcher.find(this, SimpleChannelInboundHandler.class, "I");
读取泛型类型。
1public abstract class SimpleChannelInboundHandler<I> extends ChannelInboundHandlerAdapter {
2
3 private final TypeParameterMatcher matcher;
4
5protected SimpleChannelInboundHandler(boolean autoRelease) {
6 matcher = TypeParameterMatcher.find(this, SimpleChannelInboundHandler.class, "I");
7this.autoRelease = autoRelease;
8}
9 public boolean acceptInboundMessage(Object msg) throws Exception {
10 return matcher.match(msg);
11}
12
13 @Override
14 public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
15 boolean release = true;
16try {
17 if (acceptInboundMessage(msg)) {
18 @SuppressWarnings("unchecked")
19 I imsg = (I) msg;
20channelRead0(ctx, imsg);
21} else {
22 release = false;
23ctx.fireChannelRead(msg);
24}
25 } finally {
26 if (autoRelease && release) {
27 ReferenceCountUtil.release(msg);
28}
29 }
30 }
31
32}好了,到这里我们的入栈流程就说完了。下面我们说出栈的流程。
3.6 出栈管道处理器

通过前面我们对 ChannelHandler 的了解,如果我们要写出栈的处理器,其实就是要继承 ChannelOutboundHandlerAdapter 。以下面这个例子,我们写一个 RpcResponse 对象。
- a handler 是编码器,负责将
RpcResponse对象转成二进制数据 - b handler 是出栈处理器, 而里面的Object类型的msg究竟是二进制数据呢? 还是
RpcResponse呢? 这个就要看出栈执行器的位置了。
这个问题,后面在3.7就能找到答案。
1 @Test
2 @DisplayName("出栈处理器")
3 public void testOutbound() throws Exception {
4 EmbeddedChannel channel = new EmbeddedChannel(new LoggingHandler(LogLevel.DEBUG),
5 new ChannelInitializer<EmbeddedChannel>() {
6 @Override
7 protected void initChannel(EmbeddedChannel ch) throws Exception {
8 ChannelPipeline pipeline = ch.pipeline();
9 pipeline.addLast("a handler", new MojitoChannelEncoder("mojito"));
10 pipeline.addLast("b handler", new ChannelOutboundHandlerAdapter() {
11 @Override
12 public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) throws Exception {
13 System.out.println("Write:" + msg);
14 super.write(ctx, msg, promise);
15 }
16 });
17 }
18 });
19RpcResponse rpcResponse = new RpcResponse();
20channel.writeOutbound(rpcResponse);
21}3.7 处理器顺序

1 ChannelPipeline p = ...;
2p.addLast("1", new InboundHandlerA());
3p.addLast("2", new InboundHandlerB());
4p.addLast("3", new OutboundHandlerA());
5p.addLast("4", new OutboundHandlerB());
6p.addLast("5", new InboundOutboundHandlerX());- 入栈: 1 -> 2 -> 5
- 出栈: 5 -> 4 -> 3
所以由此得出,3.5 Write的地方,会打印 RpcResponse 对象。
好了,到这里Netty的核心API其实就学习完了,了解了这些就算入门了。下面我们看如何来组装这些API吧。
3.8 服务端引导类
通过前面的学习,我们已经知道在Netty 中如何将二进制数据转换成Java对象,并进行业务处理,然后在将业务数据通过出栈处理器器转换成二进制数据返回了。那么我们现在来看下如何用Netty构建一个服务端吧。 直接上代码。
1 /**
2 * @author liuxin
3 * 个人博客:https://java.springlearn.cn
4 * 公众号:西魏陶渊明 {关注获取学习源码}
5 * 2022/8/11 23:12
6 */
7 @Test
8 @DisplayName("构建一个服务端")
9 public void testServer() throws Exception {
10 ServerBootstrap serverBootstrap = new ServerBootstrap();
11// io 线程一个进行轮训即可
12 NioEventLoopGroup bossGroup = new NioEventLoopGroup(1, new NamedThreadFactory("boss"));
13// 业务处理线程组, CPU线程数 + 1 即可: (同一个核心同一时刻只能执行一个任务,所以创建多了也没用,建议给N+1个)
14 NioEventLoopGroup workGroup = new NioEventLoopGroup(
15 Runtime.getRuntime().availableProcessors() + 1, new NamedThreadFactory("work"));
16serverBootstrap.group(bossGroup, workGroup)
17 .childHandler(new ChannelInitializer<SocketChannel>() {
18 @Override
19 protected void initChannel(SocketChannel ch) throws Exception {
20 // 我们的管道信息信息就在这里
21 ChannelPipeline pipeline = ch.pipeline();
22
23 }
24 }).channel(OSinfo.isLinux() ? EpollServerSocketChannel.class : NioServerSocketChannel.class);
25ChannelFuture sync = serverBootstrap.bind(8080).sync();
26sync.addListener((ChannelFutureListener) future -> {
27 if (future.isSuccess()) {
28 System.out.println("端口绑定成功");
29 } else {
30 System.out.println("端口绑定失败:" + future.cause().getCause());
31 }
32 });
33Channel channel = sync.channel();
34// 添加一个关闭时间监听器
35 channel.closeFuture().addListener((ChannelFutureListener) future -> {
36 if (future.isSuccess()) {
37 System.out.println("服务关闭成功");
38 } else {
39 System.out.println("服务关闭失败:" + future.cause().getCause());
40 }
41 });
42channel.close();
43}3.9 客户端引导类
这里我们为了测试,我们先构建一个服务端,然后构建一个客户端然后进行访问。我们看服务端的输出。
- 构建服务端打印请求连接和释放连接事件
1 /**
2 * 请求解码器负责将二进制数据转换成能处理的协议
3 * 虫洞栈:https://java.springlearn.cn
4 * 公众号:西魏陶渊明 {关注获取学习源码}
5 */
6 @Test
7 @DisplayName("构建一个服务端")
8 public void testServer() throws Exception {
9 ServerBootstrap serverBootstrap = new ServerBootstrap();
10// io 线程一个进行轮训即可
11 NioEventLoopGroup bossGroup = new NioEventLoopGroup(1, new NamedThreadFactory("boss"));
12// 业务处理线程组, CPU线程数 + 1 即可: (同一个核心同一时刻只能执行一个任务,所以创建多了也没用,建议给N+1个)
13 NioEventLoopGroup workGroup = new NioEventLoopGroup(
14 Runtime.getRuntime().availableProcessors() + 1, new NamedThreadFactory("work"));
15serverBootstrap.group(bossGroup, workGroup)
16 .childHandler(new ChannelInitializer<SocketChannel>() {
17 @Override
18 protected void initChannel(SocketChannel ch) throws Exception {
19 // 设置我们的管道信息
20 ChannelPipeline pipeline = ch.pipeline();
21 pipeline.addLast(new ChannelInboundHandlerAdapter() {
22 @Override
23 public void channelActive(ChannelHandlerContext ctx) throws Exception {
24 Channel channel = ctx.channel();
25 SocketAddress socketAddress = channel.remoteAddress();
26 System.out.println("收到了一个链接:" + socketAddress);
27 ctx.fireChannelActive();
28 }
29
30 @Override
31 public void channelInactive(ChannelHandlerContext ctx) throws Exception {
32 Channel channel = ctx.channel();
33 SocketAddress socketAddress = channel.remoteAddress();
34 System.out.println(socketAddress + ":关闭连接");
35 }
36 });
37 }
38 }).channel(OSinfo.isLinux() ? EpollServerSocketChannel.class : NioServerSocketChannel.class);
39ChannelFuture sync = serverBootstrap.bind(8080).sync();
40sync.addListener((ChannelFutureListener) future -> {
41 if (future.isSuccess()) {
42 System.out.println("端口绑定成功");
43 } else {
44 System.out.println("端口绑定失败:" + future.cause().getCause());
45 }
46 });
47Channel channel = sync.channel();
48// 添加一个关闭时间监听器
49 channel.closeFuture().addListener((ChannelFutureListener) future -> {
50 if (future.isSuccess()) {
51 System.out.println("服务关闭成功");
52 } else {
53 System.out.println("服务关闭失败:" + future.cause().getCause());
54 }
55 }).sync();
56channel.close();
57}- 构建客户端
1 @Test
2 public void testClient() throws Exception {
3 NioEventLoopGroup workGroup = new NioEventLoopGroup(
4 Runtime.getRuntime().availableProcessors() + 1, new NamedThreadFactory("work"));
5Bootstrap clientBootstrap = new Bootstrap();
6clientBootstrap.group(workGroup);
7clientBootstrap.channel(NioSocketChannel.class);
8clientBootstrap.option(ChannelOption.TCP_NODELAY, false);
9clientBootstrap.handler(new ChannelInitializer<SocketChannel>() {
10 @Override
11 protected void initChannel(SocketChannel ch) throws Exception {
12 // 设置我们的管道
13 ChannelPipeline pipeline = ch.pipeline();
14 // 客户端要将我们发出的Java对象转换成二进制对象输入
15 // 客户端要将服务端发送的二进制对象转换成Java对象
16 }
17 });
18ChannelFuture channelFuture = clientBootstrap.connect("127.0.0.1", 8080).sync();
19channelFuture.channel().write("HelloWord");
20}服务端控制台输出:
- 端口绑定成功
- 收到了一个链接:/127.0.0.1:55732
- /127.0.0.1:55732:关闭连接
3.10 Netty 学习总结
好了,我们一口气吧 Netty 的API都学习了,知识点有点多,大家可以看着图来理解。这里学习Netty是因为我们要用Netty来构建一个网络通道。用于我们开发RPC框架,这点知识已经够用了。但是需要注意的是 Netty 并不仅仅只有这些知识点。Netty 为什么这么快? 吞吐量这么高? 值得我们学习的知识点还有很多。这个后面我们单独再来说,本篇文章我们就了解这么多就够用了。下面我们终于可以开始自己的Coding了。
四、通信层搭建

通过前三篇的学习及上面对Netty的学习,相信上图中关于底层通信的所有知识点都已经了解了吧。那么下面就开始编程了。来一步一步完成我们的通信层。
4.1 工程结构

4.2 架构设计

所谓的架构其实就会对于Netty 管道中的处理逻辑和分层。
- Config API 其实就是我们的Fluent风格的API
- Business 就是我们提供给开发者用来实现业务的接口
- Model 是我们原型开发者自定义自己的数据传输模型,前提是要集成
ProtocolHeader - Exchange 用于屏蔽Netty原生众多的API,通过封装只暴露我们需要感知的API
- Codec 提供自定义的解码器和编码器,同时也能支持HTTP的协议
- Serialize 底层的序列化实现
4.3 服务端

- 首先我们定义出
Server接口,为了尽量让职责单一。我们将配置方法和核心的能力分开,又定义了 提供配置的ConfigurableServer。 这块代码我们就是抄的Spring的ApplicationContext的设计。设计的好处是,接口隔离原则,即一个类与其他类保留最少的关系。这样说可能还不好理解。我们思考下 Server集成了ConfigurableServer。假如说我们把所有的接口定义都放在Server中。当我们要把配置的信息,暴露给外面的时候,只能将Server给提供出去,但是Server中有那么多的非配置的方法,是不是都被外部所感知到了呢? 解决办法就是将接口细化, 给外部只提供ConfigurableServer。这样外部就看不到Server中所有的方法,就不会被困扰。
1 // 给外面提供的接口能力太大了,他不关心的提供出来就是困扰。
2 public void customerConfig(Server server);
3// 提供的都是我想要的,一起都是刚刚好。
4 public void customerConfig(ConfigurableServer confServer);-
目前我们是使用Netty通信框架进行的实现,但是为了以后可以支持其他的通信框架,我们定义了抽象模板类
AbstractServer。在模板类型,只定义统一的创建服务端的流程。而具体的细节,交给了抽象方法。如果说后面我们不使用Netty来,我们的改动也是最小的。 -
NettyServer继承了AbstractServer实现了其定义的抽象方法,具体的负责创建服务。
代码如下,更多细节可以到下载代码学习。
1public interface ConfigurableServer<T extends Server<?>> {
2
3 /**
4 * 给网络通道注册二进制处理协议
5 *
6 * @param protocol 协议
7 */
8 void registryProtocol(Protocol<? extends ProtocolHeader, ? extends ProtocolHeader> protocol);
9
10 /**
11 * 注册钩子程序
12 */
13void registryHooks(Runnable hookTask);
14
15 /**
16 * 协议信息
17 *
18 * @return Protocol
19 */
20Protocol<? extends ProtocolHeader, ? extends ProtocolHeader> getProtocol();
21
22 /**
23 * 这里我们提供一个Server初始化的方法,为什么呢?
24 * 目前我们的服务端是使用NettyServer,我们也支持其他的通信框架。因为可能初始化方法不一样.
25 * 所以我们将具体的实现作为一个泛型。让具体的实现来自己定义自己的初始化方法
26 *
27 * @param initializer 初始化接口
28 */
29void initializer(ServerInitializer<T> initializer);
30}
31
32public interface Server<T extends Server<?>> extends ConfigurableServer<T> {
33
34 /**
35 * 启动服务
36 *
37 * @param port 服务端口号
38 */
39 void start(int port);
40
41 /**
42 * 非阻塞启动
43 *
44 * @param port 端口
45 */
46void startAsync(int port);
47
48 /**
49 * 关闭服务
50 */
51void close();
52
53 /**
54 * 服务器端口
55 *
56 * @return int
57 */
58int getPort();
59
60 /**
61 * 是否运行中
62 *
63 * @return boolean
64 */
65boolean isRun();
66
67}
68
69/**
70 * @author liuxin
71 * 2022/8/10 22:16
72 */
73public abstract class AbstractServer<T extends Server<?>> implements Server<T> {
74
75 private ServerInitializer<T> serverInitializer;
76
77private Integer port;
78
79private final AtomicBoolean runningState = new AtomicBoolean(false);
80
81private Protocol<? extends ProtocolHeader, ? extends ProtocolHeader> protocol;
82
83@Override
84 public void registryProtocol(Protocol<? extends ProtocolHeader, ? extends ProtocolHeader> protocol) {
85 this.protocol = protocol;
86}
87
88 @Override
89 public void registryHooks(Runnable hookTask) {
90 ThreadHookTools.addHook(new Thread(hookTask));
91}
92
93 @Override
94 public Protocol<? extends ProtocolHeader, ? extends ProtocolHeader> getProtocol() {
95 return protocol;
96}
97
98 @Override
99 public void initializer(ServerInitializer<T> initializer) {
100 this.serverInitializer = initializer;
101}
102
103 @Override
104 public ServerInitializer<T> getServerInitializer() {
105 return serverInitializer;
106}
107
108 @Override
109 public void start(int port) {
110 checked();
111activeAndCreateServer(() -> {
112 this.port = port;
113doCreateServer(port, false);
114});
115}
116
117 private void checked() {
118 if (Objects.isNull(protocol)) {
119 throw new RuntimeException("Protocol不能为空,请Server#registryProtocol");
120}
121 if (Objects.isNull(serverInitializer)) {
122 throw new RuntimeException("Protocol不能为空,请Server#initializer");
123}
124 }
125
126 private void activeAndCreateServer(LambdaExecute execute) {
127 if (isRun()) {
128 throw new RuntimeException("运行中");
129}
130 if (runningState.compareAndSet(false, true)) {
131 try {
132 execute.execute();
133} catch (Throwable t) {
134 runningState.compareAndSet(true, false);
135}
136 }
137 }
138
139 private void closeAndDestroyServer(LambdaExecute execute) {
140 if (isRun()) {
141 if (runningState.compareAndSet(true, false)) {
142 execute.execute();
143}
144 }
145 }
146
147 @Override
148 public void startAsync(int port) {
149 activeAndCreateServer(() -> {
150 this.port = port;
151 doCreateServer(port, true);
152 });
153}
154
155 @Override
156 public void close() {
157 closeAndDestroyServer(this::doDestroyServer);
158}
159
160 @Override
161 public int getPort() {
162 return this.port;
163}
164
165 @Override
166 public boolean isRun() {
167 return runningState.get();
168}
169
170 public abstract void doCreateServer(int port, boolean async);
171
172public abstract void doDestroyServer();
173}
174
175@Slf4j
176public class NettyServer extends AbstractServer<NettyServer> {
177
178 private final ServerBootstrap serverBootstrap = new ServerBootstrap();
179
180private Channel serverChannel;
181
182private EventLoopGroup bossGroup;
183
184private EventLoopGroup workerGroup;
185
186private static final int DEFAULT_EVENT_THREADS = Math.min(Runtime.getRuntime().availableProcessors() + 1, 32);
187
188public ServerBootstrap getServerBootstrap() {
189 return this.serverBootstrap;
190}
191
192 @Override
193 @SneakyThrows
194 public void doCreateServer(int port, boolean async) {
195 // 1. io线程数 = cpu * 2
196 bossGroup = new NioEventLoopGroup(1, new NamedThreadFactory("mojito-boss", true));
197// 2. 业务线程数 = cpu + 1
198 workerGroup = new NioEventLoopGroup(DEFAULT_EVENT_THREADS, new NamedThreadFactory("mojito-work", true));
199serverBootstrap.group(bossGroup, workerGroup)
200 .childOption(ChannelOption.TCP_NODELAY, true)
201 .childOption(ChannelOption.SO_KEEPALIVE, true)
202 .option(ChannelOption.SO_BACKLOG, 128)
203 .handler(new LoggingHandler(LogLevel.INFO))
204 .channel(OSinfo.isLinux() ? EpollServerSocketChannel.class : NioServerSocketChannel.class)
205 .localAddress(port).option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 10000);
206getServerInitializer().initializer(this);
207// 3. 阻塞绑定端口
208 ChannelFuture bindFuture = serverBootstrap.bind().addListener((ChannelFutureListener) channelFuture -> {
209 if (channelFuture.isSuccess()) {
210 log.info(Banner.print("麻烦给我的爱人来一杯Mojito,我喜欢阅读她微醺时的眼眸!", Ansi.Color.RED));
211 log.info("Mojito启动成功,端口号:" + port);
212 } else {
213 Throwable cause = channelFuture.cause();
214 throw new RuntimeException(cause);
215 }
216 }).sync();
217serverChannel = bindFuture.channel();
218if (async) {
219 log.info("异步服务启动成功");
220} else {
221 serverChannel.closeFuture().sync();
222log.info("阻塞服务启动成功");
223}
224 }
225
226 @Override
227 public void doDestroyServer() {
228 workerGroup.shutdownGracefully();
229bossGroup.shutdownGracefully();
230serverChannel.close();
231}
232
233}
2344.4 客户端

服务端的设计和客户端的设计是一样的,依赖抽象不依赖细节。架构设计只是定义好流程,具体是什么框架来实现底层的通信,就交给最底层的细节。
- 客户端有几个核心的地方
- 连接服务器
- 发送数据(同步&异步)
- 断线重连(放在优化的时候讲)
- 异步转同步问题
1
2/**
3 * @author liuxin
4 * 2022/8/5 23:12
5 */
6public interface ConfigurableClient<REQ extends ProtocolHeader, RES extends ProtocolHeader, T extends Client<REQ, RES>> {
7
8 /**
9 * 给网络通道注册二进制处理协议
10 *
11 * @param protocol 协议
12 */
13 void registryProtocol(Protocol<REQ, RES> protocol);
14
15 /**
16 * 协议信息
17 *
18 * @return Protocol
19 */
20Protocol<REQ, RES> getProtocol();
21
22 /**
23 * 注册钩子程序
24 */
25void registryHooks(Runnable hookTask);
26
27 /**
28 * 这里我们提供一个Server初始化的方法,为什么呢?
29 * 目前我们的服务端是使用NettyServer,我们也支持其他的通信框架。因为可能初始化方法不一样.
30 * 所以我们将具体的实现作为一个泛型。让具体的实现来自己定义自己的初始化方法
31 *
32 * @param initializer 初始化接口
33 */
34void initializer(ClientInitializer initializer);
35
36 /**
37 * 客户端初始化扩展
38 *
39 * @return ClientInitializer
40 */
41ClientInitializer<Client<REQ, RES>> getClientInitializer();
42
43}
44
45/**
46 * @author liuxin
47 * 个人博客:https://java.springlearn.cn
48 * 公众号:西魏陶渊明 {关注获取学习源码}
49 * 2022/8/5 23:12
50 */
51public interface Client<REQ extends ProtocolHeader, RES extends ProtocolHeader> extends ConfigurableClient<REQ, RES, Client<REQ, RES>> {
52
53
54 /**
55 * 建立连接
56 *
57 * @param host 连接地址
58 * @param port 连接端口
59 */
60 void connect(String host, Integer port);
61
62 /**
63 * 发送消息
64 *
65 * @param req 请求体
66 * @return 异步结果
67 */
68MojitoFuture<RES> send(REQ req);
69
70 /**
71 * 关闭连接
72 */
73void close();
74
75
76 /**
77 * 是否连接
78 *
79 * @return boolean
80 */
81boolean isRun();
82
83 /**
84 * 是否连接
85 *
86 * @return boolean
87 */
88boolean isConnected();
89
90
91}
92/**
93 * @author liuxin
94 * 个人博客:https://java.springlearn.cn
95 * 公众号:西魏陶渊明 {关注获取学习源码}
96 * 2022/8/5 23:12
97 */
98public abstract class AbstractClient<REQ extends ProtocolHeader, RES extends ProtocolHeader> implements Client<REQ, RES> {
99
100
101 /**
102 * 将要连接的远程地址
103 */
104 private String remoteHost;
105
106 /**
107 * 将要连接的远程端口
108 */
109private int remotePort;
110
111private Protocol<REQ, RES> protocol;
112
113private ClientInitializer<Client<REQ, RES>> clientInitializer;
114
115private final AtomicBoolean running = new AtomicBoolean(false);
116
117@Override
118 public void connect(String host, Integer port) {
119 if (!running.get()) {
120 this.remoteHost = host;
121this.remotePort = port;
122}
123 try {
124 if (running.compareAndSet(false, true)) {
125 doConnect();
126}
127 } catch (Throwable t) {
128 t.printStackTrace();
129running.compareAndSet(true, false);
130}
131 }
132
133 @Override
134 public MojitoFuture<RES> send(REQ req) {
135 return doSend(req);
136}
137
138 @Override
139 public void close() {
140 doClose();
141}
142
143 public int getRemotePort() {
144 return remotePort;
145}
146
147 public String getRemoteHost() {
148 return remoteHost;
149}
150
151 @Override
152 public Protocol<REQ, RES> getProtocol() {
153 return this.protocol;
154}
155
156
157 @Override
158 public void registryProtocol(Protocol<REQ, RES> protocol) {
159 this.protocol = protocol;
160}
161
162 @Override
163 public void registryHooks(Runnable hookTask) {
164 ThreadHookTools.addHook(new Thread(hookTask));
165}
166
167 @Override
168 public void initializer(ClientInitializer initializer) {
169 this.clientInitializer = initializer;
170}
171
172 @Override
173 public boolean isRun() {
174 return running.get();
175}
176
177 @Override
178 public ClientInitializer<Client<REQ, RES>> getClientInitializer() {
179 return clientInitializer;
180}
181
182 public abstract void doConnect();
183
184public abstract void doClose();
185
186public abstract MojitoFuture<RES> doSend(REQ req);
187}
188
189/**
190 * @author liuxin
191 * 个人博客:https://java.springlearn.cn
192 * 公众号:西魏陶渊明 {关注获取学习源码}
193 * 2022/8/5 23:12
194 */
195@Slf4j
196public class NettyClient<REQ extends ProtocolHeader, RES extends ProtocolHeader> extends AbstractClient<REQ, RES> {
197
198 private final Bootstrap clientBootstrap = new Bootstrap();
199
200private final EventLoopGroup workerGroup = new NioEventLoopGroup();
201
202private DefaultEnhanceChannel enhanceChannel;
203
204@Override
205 @SneakyThrows
206 public void doConnect() {
207 clientBootstrap.group(workerGroup);
208clientBootstrap.channel(NioSocketChannel.class);
209clientBootstrap.option(ChannelOption.TCP_NODELAY, false);
210getClientInitializer().initializer(this);
211ChannelFuture channelFuture = clientBootstrap.connect(getRemoteHost(), getRemotePort()).sync();
212enhanceChannel = DefaultEnhanceChannel.getOrAddChannel(channelFuture.channel());
213}
214
215 public Bootstrap getClientBootstrap() {
216 return clientBootstrap;
217}
218
219 @Override
220 public void doClose() {
221 workerGroup.shutdownGracefully();
222enhanceChannel.disconnected();
223log.info("Client 关闭成功");
224}
225
226 @Override
227 public MojitoFuture<RES> doSend(REQ req) {
228 // 这里我们也设置,断线重连,后面优化
229 return getProtocol().getClientPromiseHandler().sendAsync(enhanceChannel, req);
230}
231
232 @Override
233 public boolean isConnected() {
234 return Objects.nonNull(enhanceChannel) && enhanceChannel.isConnected();
235}
236
237
238}
2394.5 异步转同步
socket通信发送数据,什么时候回复都可以。甚至可以客户端一直发,而服务端不进行回复。而我们的RPC框架更像一问一答,发送请求后,需要立马就收到回复。每个请求都要对应一个响应,这就需要我们进行特殊的设计来完成,这样的需求。
我们的思路就是实现,Jdk的Future,并给他添加上监听器的功能。因为我们主要是学习,所以不要怕麻烦,不要怕重新造轮子。下面我们开始实现它。
1mojito/mojito-net/src/main/java/cn/lxchinesszz/mojito/future on master [!+?]
2➜ tree
3.
4├── MojitoFuture.java
5├── Promise.java
6└── listener
7 ├── MojitoListener.java
8 └── MojitoListeners.java
9- 定义接口
Promise。承诺,这个接口一定会告诉你成功或者失败。这样的设计Js上运用是最多的。
1public interface Promise<V> {
2
3 /**
4 * 是否成功
5 *
6 * @return boolean
7 */
8 boolean isSuccess();
9
10 /**
11 * 设置成功表示
12 *
13 * @param data 数据
14 */
15void setSuccess(V data);
16
17 /**
18 * 设置失败标识
19 *
20 * @param cause 异常
21 */
22void setFailure(Throwable cause);
23
24 /**
25 * 添加监听器
26 *
27 * @param listener 监听器
28 */
29void addListeners(MojitoListener<V> listener);
30}- 实现核心方法
这里需要注意的就是多线程的可见性和多线程的原子性问题以及如何实现线程等待。如果这部分不太容易理解,建议可以先多了解点关于线程安全的知识点。然后回过头在来看。
1public class MojitoFuture<V> implements Promise<V>, Future<V> {
2
3 /**
4 * volatile 多线程保证可见性
5 */
6 private volatile V result;
7
8 /**
9 * AtomicReferenceFieldUpdater 多线程保证操作的原子性
10 */
11@SuppressWarnings("rawtypes")
12 private static final AtomicReferenceFieldUpdater<MojitoFuture, Object> RESULT_UPDATER =
13 AtomicReferenceFieldUpdater.newUpdater(MojitoFuture.class, Object.class, "result");
14
15 /**
16 * 当前承诺的监听器
17 */
18private final List<MojitoListener<V>> listeners = new ArrayList<>();
19
20 /**
21 * 是否被撤销了
22 */
23private boolean cancelled = false;
24
25 /**
26 * 并发锁
27 */
28private final ReentrantLock lock = new ReentrantLock();
29
30 /**
31 * get 阻塞条件,用于完成时候唤醒get阻塞线程
32 */
33private final Condition condition = lock.newCondition();
34
35 /**
36 * get超时阻塞条件,,用于完成时候唤醒get阻塞线程
37 */
38private final Condition timeoutCondition = lock.newCondition();
39
40 /**
41 * 原子操作
42 *
43 * @param objResult 数据
44 * @return boolean
45 */
46private boolean setValue0(Object objResult) {
47 if (RESULT_UPDATER.compareAndSet(this, null, objResult) ||
48 RESULT_UPDATER.compareAndSet(this, EMPTY, objResult)) {
49 return true;
50}
51 return false;
52}
53
54 @Override
55 public boolean isSuccess() {
56 return isDone() && result != EMPTY;
57}
58
59 @Override
60 public void setSuccess(V data) {
61 boolean updateSuccess;
62try {
63 lock.lock();
64updateSuccess = setValue0(data);
65condition.signalAll();
66timeoutCondition.signalAll();
67} finally {
68 lock.unlock();
69}
70 if (updateSuccess) {
71 for (MojitoListener<V> listener : listeners) {
72 listener.onSuccess(data);
73}
74 }
75 }
76
77 @Override
78 public void setFailure(Throwable cause) {
79 boolean updateSuccess;
80try {
81 lock.lock();
82updateSuccess = setValue0(EMPTY);
83condition.signalAll();
84timeoutCondition.signalAll();
85} finally {
86 lock.unlock();
87}
88 if (updateSuccess) {
89 for (MojitoListener<V> listener : listeners) {
90 listener.onThrowable(cause);
91}
92 }
93 }
94
95
96 @Override
97 public boolean cancel(boolean mayInterruptIfRunning) {
98 this.cancelled = mayInterruptIfRunning;
99return false;
100}
101
102 @Override
103 public boolean isCancelled() {
104 return cancelled;
105}
106
107 @Override
108 public boolean isDone() {
109 // 只要不等于空,说明就是有结果了,不管成功或者失败
110 return result != null;
111}
112
113 @Override
114 public V get() throws InterruptedException, ExecutionException {
115lock.lock();
116try {
117 if (Thread.currentThread().isInterrupted()) {
118 throw new InterruptedException();
119} else {
120 // 如果没有完成并且没有撤销
121 if (!isDone() && !isCancelled()) {
122 // 释放线程,任务在这里阻塞,等待完成时候释放.
123 condition.await();
124}
125 }
126 } finally {
127 lock.unlock();
128}
129 return isSuccess() ? result : null;
130}
131
132 @Override
133 public V get(long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException {
134lock.lock();
135try {
136 if (Thread.currentThread().isInterrupted()) {
137 throw new InterruptedException();
138} else {
139 // 如果没有完成并且没有撤销
140 if (!isDone() && !isCancelled()) {
141 // 释放线程,任务在这里阻塞,等待完成时候释放.
142 boolean await = timeoutCondition.await(timeout, unit);
143if (!await) {
144 throw new TimeoutException();
145}
146 }
147 }
148 } finally {
149 lock.unlock();
150}
151 return isSuccess() ? result : null;
152}
153
154 @Override
155 public void addListeners(MojitoListener<V> listener) {
156 this.listeners.add(listener);
157}
158}4.6 Fluent API
好了,前面我们快速的学习了Netty的API,后面我们使用了Netty来实现了我们的底层通信能力。但是这里还是太复杂了,最后我们要为这些复杂的对象,设计一套简单的使用API。API的风格决定使用 Fluent API
fluent-API 是一种面向对象的 API,其设计主要基于方法链。 这个概念由Eric Evans和Martin Fowler于 2005 年创建,旨在通过创建特定领域语言 ( DSL )来提高代码可读性。
在实践中,创建一个流畅的 API,就是不需要记住接下来的步骤或方法,一切都是那么的自然和连续,下一步的动作,就好像它是一个选项菜单,让我们的选择。
关键词: 自然连续,无需记住
这里面主要是对Java泛型的利用,主要实现在这里,基本都是泛型,所以要好好看。建议获取源码,运行走走。
1mojito/mojito-net/src/main/java/cn/lxchinesszz/mojito/fluent on master [»!?]
2➜ tree
3.
4├── AbstractFactory.java
5├── ConfigurableFactory.java
6├── Factory.java
7└── Mojito.java
8下面我们来展示,来看下这个API,是否合乎你的心意呢?
- 服务端
1class MojitoTest{
2 /**
3 * @author liuxin
4 * 个人博客:https://java.springlearn.cn
5 * 公众号:西魏陶渊明 {关注获取学习源码}
6 * 2022/8/11 23:12
7 */
8 @Test
9 @DisplayName("构建服务端【阻塞方式】")
10 public void server() throws Exception {
11 Server<?> server = Mojito.server(RpcRequest.class, RpcResponse.class)
12 // 业务层,读取请求对象,返回结果
13 .businessHandler((channelContext, request) -> new RpcResponse())
14 .create();
15server.start(6666);
16}
17
18 /**
19 * @author liuxin
20 * 个人博客:https://java.springlearn.cn
21 * 公众号:西魏陶渊明 {关注获取学习源码}
22 * 2022/8/11 23:12
23 */
24 @Test
25 @DisplayName("构建服务端【非阻塞方式】")
26 public void serverAsync() throws Exception {
27 Server<?> server = Mojito.server(RpcRequest.class, RpcResponse.class)
28 // 业务层,读取请求对象,返回结果
29 .businessHandler((channelContext, request) -> new RpcResponse())
30 .create();
31server.startAsync(6666);
32}
33}- 客户端
支持同步和异步两种方式
1 /**
2 * @author liuxin
3 * 个人博客:https://java.springlearn.cn
4 * 公众号:西魏陶渊明 {关注获取学习源码}
5 * 2022/8/11 23:12
6 */
7 @Test
8 @DisplayName("构建客户端【异步方式】")
9 public void clientAsync() throws Exception {
10 // 构建连接
11 Client<RpcRequest, RpcResponse> client = Mojito.client(RpcRequest.class, RpcResponse.class)
12 .connect("127.0.0.1", 6666);
13
14MojitoFuture<RpcResponse> sendFuture = client.sendAsync(new RpcRequest());
15sendFuture.addListeners(new MojitoListener<RpcResponse>() {
16 @Override
17 public void onSuccess(RpcResponse result) {
18 System.out.println("收到结果:" + result);
19 }
20
21 @Override
22 public void onThrowable(Throwable throwable) {
23 System.err.println("处理失败:" + throwable.getMessage());
24 }
25 });
26Thread.currentThread().join();
27}
28
29 /**
30 * @author liuxin
31 * 个人博客:https://java.springlearn.cn
32 * 公众号:西魏陶渊明 {关注获取学习源码}
33 * 2022/8/11 23:12
34 */
35 @Test
36 @DisplayName("构建客户端【同步方式】")
37 public void clientSync() throws Exception{
38 Client<RpcRequest, RpcResponse> client = Mojito.client(RpcRequest.class, RpcResponse.class)
39 .connect("127.0.0.1", 6666);
40System.out.println(client.send(new RpcRequest()));
41}
42- 启动演示
当看到下面的Logo说明服务已经启动成功了。
1 ___ ___ ______ ___ __ ___________ ______
2|" /" | / " |" ||" (" _ ")/ "
3 // | // ____ || ||| |)__/ \__/// ____
4 / /. | / / ) :) |: ||: | \_ / / / ) :)
5|: . |(: (____/ //___| / |. | |. | (: (____/ //
6|. /: | // :|_/ )/ | : | /
7|___|\__/|___| "_____/(_______/(__\_|_) \__| "_____/
8
9 :: Mojito ::
10麻烦给我的爱人来一杯Mojito,我喜欢阅读她微醺时的眼眸!
1122:53:44.652 [mojito-boss-thread-1] INFO cn.lxchinesszz.mojito.server.netty.NettyServer - Mojito启动成功,端口号:6666
1222:53:44.653 [mojito-boss-thread-1] INFO io.netty.handler.logging.LoggingHandler - [id: 0x7fc4d842, L:/0:0:0:0:0:0:0:0:6666] ACTIVEFluent API 风格是自然联系的,仿佛就跟菜单一样,根本不需要去记API。一起都是那么的自然。
五、总结
本篇文章爆肝了11天, 因为只有晚上下班,回来才有时间来思考总结。所以进度有点慢。
文章前部分介绍 Netty API,后半部分介绍基于 Netty 来设计我们的通信层。最终通过
Fluent API 的风格,将复杂的代码,通过简单的API给暴露了出来。
但是做到这一步只能说是完成了需求,后面我们还要做压测和调优。
- 是否可以使用多线程?
- 耗时对象是否可以进行池化?
- 序列化为什么还没有支持
Protostuff? - 各种异常场景是否都捕捉到了,给出清晰的提示?
- 能否提供更多的扩展功能?
mojito-net 只能做RPC吗? 难道不能做一个简单的 web容器? 难道不能实现一个mq 吗?
😊 那么你准备好跟我一起Coding了吗?,如果喜欢麻烦点个关注。
