Netty实现简易的应用层协议

前一阵在工作中用到了RabbitMQ,因此对几种常见的消息队列产生了兴趣。首先从GitHub上下载了RocketMQ的源码打算一探究竟。在阅读remoting这个模块时遇到了很大的障碍。RocketMQ的网络编程完全基于Netty,而本人对Netty的理解还只停留在了知道这是一款封装了NIO的优秀框架上。于是正好就借此机会先揭开Netty的面纱。
阅读完《Netty IN ACTION》后有些手痒,就用Netty实现了一个简易的应用层协议以及一个同步调用的方法。
github:https://github.com/ztglcy/netty-protocol

整体结构

程序结构

图里demo中是client和server的简易demo;handler中则是自定义协议的编码器和解码器;message中是与传输的消息相关的类;processor是服务端的业务处理类;service中则是client和server的启动类。

传输消息格式及编解码

传输的消息格式

length是一个表示消息大小的int型数字,自定义长度解码器解决TCP黏包问题。headerLength则是表示消息头大小的int型数字,用以将传输的消息头与消息体分开进行序列化。header和content分别存储消息的消息头以及消息体。

public class MessageHeader{

    private int messageId;
    private int clientId;
    private int serverId;
    private int code;

    private MessageHeader(){}

    public MessageHeader(int code) {
        this.code = code;
    }

    public int getCode() {
        return code;
    }
    public int getMessageId() {
        return messageId;
    }

    public void setMessageId(int messageId) {
        this.messageId = messageId;
    }

    public int getClientId() {
        return clientId;
    }

    public void setClientId(int clientId) {
        this.clientId = clientId;
    }

    public int getServerId() {
        return serverId;
    }

    public void setServerId(int serverId) {
        this.serverId = serverId;
    }

    public void setCode(int code) {
        this.code = code;
    }

}

消息头包含messageId,clientId,serverId,code四个参数,分别用以表征Message的ID,客户端ID,服务端ID,以及消息体的格式code(维护在一个常量中)。公有的构造方法中必传消息体格式code,私有的构造方法用于fastjson的反序列化。

public class Message {

    private MessageHeader messageHeader;

    private byte[] content;

    private Message(){
    }

    public static Message createMessage(MessageHeader messageHeader){
        Message message = new Message();
        message.messageHeader = messageHeader;
        return message;
    }

    public byte[] getContent() {
        return content;
    }

    public void setContent(byte[] content) {
        this.content = content;
    }

    public MessageHeader getMessageHeader() {
        return messageHeader;
    }
     //还有编码与解码的方法
    ......
}

Message的结构比较直白,包含消息头和消息体以及编码与解码的方法。解码与编码的方法在解码器与编码器中进行调用。先来介绍一下编码的过程。

public class ProtocolEncoder extends MessageToByteEncoder<Message> {

    @Override
    protected void encode(ChannelHandlerContext channelHandlerContext, Message message, ByteBuf byteBuf) throws Exception {

        ByteBuffer byteBuffer = message.encode();
        byteBuf.writeBytes(byteBuffer);

    }
}

编码器继承了MessageToByteEncoder并实现了其encode()方法进行编码。编码器直接调用传进来的Message自己的编码方法,将编码后的ByteBuffer写入ByteBuf中。再来看一下Message怎么实现这个方法的。

 public ByteBuffer encode(){
        int length = 4;
        byte[] bytes = SerializableHelper.encode(messageHeader);
        if(bytes != null){
            length += bytes.length;
        }
        if(content!=null){
            length += content.length;
        }

        ByteBuffer byteBuffer = ByteBuffer.allocate(length + 4);
        byteBuffer.putInt(length);
        if(bytes != null){
            byteBuffer.putInt(bytes.length);
            byteBuffer.put(bytes);
        }else{
            byteBuffer.putInt(0);
        }
        if(content!=null){
            byteBuffer.put(content);
        }
        byteBuffer.flip();

        return byteBuffer;
    }

length用以表示整个消息大小,计算方式为表示消息头大小的int+序列化后的消息头大小+消息体大小。计算完成后申请一块length+4大小的ByteBuffer(因为length本身存储也要4个字节)。将消息内容按照上面给出的格式依次写入Buffer中。

public class ProtocolDecoder extends LengthFieldBasedFrameDecoder {

    public ProtocolDecoder() {
        super(16777216, 0, 4,0,4);
    }

    @Override
    public Object decode(ChannelHandlerContext ctx, ByteBuf in) throws Exception {
        ByteBuf byteBuf = (ByteBuf) super.decode(ctx, in);
        if(byteBuf == null){
            return null;
        }
        ByteBuffer byteBuffer = byteBuf.nioBuffer();
        return Message.decode(byteBuffer);
    }

}

解码器继承了LengthFieldBasedFrameDecoder用以解决粘包的问题。构造函数中的参数分别表示包的最大值、长度字段的偏移量、长度字段占的字节数、添加到长度字段的补偿值以及从解码帧中第一次去除的字节数。因为消息的头部存储了4个字节的表示消息大小的Int型,所以后4个参数为0、4、0、4。经过处理后的消息已经剥离掉了最消息头部的Int型。再调用Message自身的decode()方法进行解码。

public static Message decode(ByteBuffer byteBuffer){
        int length = byteBuffer.limit();
        int headerLength = byteBuffer.getInt();

        byte[] headerData = new byte[headerLength];
        byteBuffer.get(headerData);
        MessageHeader messageHeader = SerializableHelper.decode(headerData,MessageHeader.class);

        byte[] content = new byte[length - headerLength -4];
        byteBuffer.get(content);

        Message message = Message.createMessage(messageHeader);
        message.setContent(content);
        return message;
    }

解码时首先将消息头的长度从ByteBuffer中取出,然后读取该长度的字节作为消息头进行反序列化,其他部分则作为消息体,重新组装成新的Message。

服务端与客户端的引导

public interface ProtocolService {

    void start();

    void shutdown();
    
}

服务端与客户端都继承了ProtocolService接口,实现了start()和shutdown()两个方法。

public class NettyProtocolServer implements ProtocolService {

    private EventLoopGroup bossGroup = new NioEventLoopGroup(1);
    private EventLoopGroup workerGroup = new NioEventLoopGroup();
    private Map<Integer,ProtocolProcessor> processorMap = new HashMap<>();

    @Override
    public void start(){
        try{
            ServerBootstrap bootstrap = new ServerBootstrap();
            bootstrap.group(bossGroup,workerGroup)
                    .channel(NioServerSocketChannel.class)
                    .localAddress(8888)
                    .childHandler(new ChannelInitializer<SocketChannel>() {
                        @Override
                        public void initChannel(SocketChannel ch) throws Exception {
                            ch.pipeline()
                                    .addLast(new ProtocolEncoder()
                                            ,new ProtocolDecoder()
                                            ,new ProtocolServerProcessor()
                                    );
                        }
                    });
            ChannelFuture cf = bootstrap.bind().sync();
        } catch (InterruptedException e) {

        }
    }

    @Override
    public void shutdown() {
        bossGroup.shutdownGracefully();
        workerGroup.shutdownGracefully();
    }

    public void registerProcessor(Integer code,ProtocolProcessor protocolProcessor){
        processorMap.put(code,protocolProcessor);
    }

    public class ProtocolServerProcessor extends SimpleChannelInboundHandler<Message> {

        @Override
        protected void channelRead0(ChannelHandlerContext channelHandlerContext, Message message) throws Exception {
            Integer code = message.getMessageHeader().getCode();
            ProtocolProcessor processor = processorMap.get(code);
            if(processor != null){
                Message response = processor.process(message);
                channelHandlerContext.writeAndFlush(response);
            }
        }
    }
}

服务端的start()方法中完成了服务端的初始化,很常见的netty写法,将编码器、解码器以及业务处理器ProtocolServerProcessor加入了Worker的Pipeline中。这里业务处理器也可以放在线程池里执行,防止业务处理时间太长造成堵塞。shutdown()方法则将两个EventLoopGroup进行关闭,防止资源泄露。registerProcessor()方法则是将业务处理器以KV的形式注册到服务端中,ProtocolServerProcessor会根据消息头中的code在map中查找对应的业务处理器进行业务的处理。

public class DemoProcessor implements ProtocolProcessor{

    @Override
    public Message process(Message message) {
        byte[] bodyDate = message.getContent();
        DemoMessageBody messageBody = SerializableHelper.decode(bodyDate,DemoMessageBody.class);
        System.out.println(messageBody.getDemo());

        MessageHeader messageHeader = new MessageHeader(1);
        messageHeader.setMessageId(message.getMessageHeader().getMessageId());
        DemoMessageBody responseBody = new DemoMessageBody();
        responseBody.setDemo("I received!");

        Message response = Message.createMessage(messageHeader);
        response.setContent(SerializableHelper.encode(responseBody));
        return response;

    }
}

DemoProcessor是一个示例的业务处理器,将传来的消息体解码后返回一个I received!的回复,这里注意的是messageId要与请求的消息一致,用以表征这是哪个请求的返回。

public class NettyProtocolClient implements ProtocolService {

    private EventLoopGroup eventLoopGroup = new NioEventLoopGroup();
    private Bootstrap bootstrap = new Bootstrap();
    private ConcurrentHashMap<Integer,Response> responseMap = new ConcurrentHashMap<>();

    @Override
    public void start() {
        bootstrap.group(eventLoopGroup)
                .channel(NioSocketChannel.class)
                .handler(new ChannelInitializer<SocketChannel>() {
                    @Override
                    protected void initChannel(SocketChannel ch) throws Exception {
                        ch.pipeline()
                                .addLast(new ProtocolDecoder(),
                                        new ProtocolEncoder(),
                                        new ProtocolClientProcessor());
                    }
                });
    }

    @Override
    public void shutdown() {
        eventLoopGroup.shutdownGracefully();
    }

    public Message send(String address, Message message){
        try {
            Response response = new Response();
            responseMap.put(message.getMessageHeader().getMessageId(),response);
            Channel channel = bootstrap.connect(address,8888).sync().channel();

            channel.writeAndFlush(message);
            Message responseMessage = response.waitResponse();
            responseMap.remove(message.getMessageHeader().getMessageId());
            channel.close();
            return responseMessage;
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        return null;
    }

    public class ProtocolClientProcessor extends SimpleChannelInboundHandler<Message> {

        @Override
        protected void channelRead0(ChannelHandlerContext channelHandlerContext, Message message) throws Exception {
            Response response = responseMap.get(message.getMessageHeader().getMessageId());
            if (response != null){
                response.putResponse(message);
            }
        }
    }
}

客户端与服务端的start()方法和shutdown()方法类似。客户端提供了一个send()方法用以消息的同步调用,send()方法在发送信息后以消息的messageId为key生成一个Response的实例缓存在responseMap中,调用Response中的countDownLatch.await()方法堵塞住等待返回(这里应该加一个时间限制以防止线程无限期地堵塞住)。ProtocolClientProcessor会处理返回的消息,将其存入对应的Response中,并调用countDownLatch.countDown()。这样客户端线程就可以收到结果同步返回。还可以改进的一点在于保持客户端与服务端的长连接,将其缓存在客户端中,每次发送消息都用已缓存的连接,减少开销。

DEMO

最后分别编写一个客户端与服务端的demo用以测试我们的协议。

public class ServerDemo {

    public static void main(String[] args) {
        NettyProtocolServer nettyProtocolServer = new NettyProtocolServer();
        nettyProtocolServer.registerProcessor(1,new DemoProcessor());
        nettyProtocolServer.start();
        try {
            Thread.sleep(10000);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        nettyProtocolServer.shutdown();
    }
}

服务端的demo首先构建一个NettyProtocolServer的实例,将DemoProcessor注册到服务端中之后挂起主线程等待客户端的消息,最后shutdown掉NettyProtocolServer。

public class ClientDemo {
    public static void main(String[] args) {
        NettyProtocolClient client = new NettyProtocolClient();
        client.start();
        Message message = demoMessage();
        Message messageResponse = client.send("localhost",message);
        System.out.println(SerializableHelper.decode(messageResponse.getContent(),DemoMessageBody.class).getDemo());
        client.shutdown();
    }

    private static Message demoMessage(){
        MessageHeader messageHeader = new MessageHeader(1);
        messageHeader.setMessageId(1);
        messageHeader.setClientId(1);
        messageHeader.setServerId(1);
        Message message =Message.createMessage(messageHeader);
        DemoMessageBody responseBody = new DemoMessageBody();
        responseBody.setDemo("Hello World!");
        message.setContent(SerializableHelper.encode(responseBody));
        return message;
    }
}

客户端的demo也很简单。构建一个NettyProtocolClient的实例,拼装一个消息,调用send()方法,再对返回的消息稍加处理就OK啦(客户端拼装和处理消息可以再抽出一个中间层)。

最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 205,033评论 6 478
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 87,725评论 2 381
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 151,473评论 0 338
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 54,846评论 1 277
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 63,848评论 5 368
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 48,691评论 1 282
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 38,053评论 3 399
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 36,700评论 0 258
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 42,856评论 1 300
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 35,676评论 2 323
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 37,787评论 1 333
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 33,430评论 4 321
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 39,034评论 3 307
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 29,990评论 0 19
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 31,218评论 1 260
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 45,174评论 2 352
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 42,526评论 2 343