zoukankan      html  css  js  c++  java
  • netty服务端实战(一)

      先分享一下自己的经历。

      去年7月进入新公司没多久,部门领导就给我分配了一个任务:给公司的一个户外设备写一个采集数据程序,将数据入库,然后做一个web端。因为领导是做.NET的,当时在来之前有和领导沟通过,领导的意思是希望来一个会网络编程和多线程,部门急需一个可以来做采集程序的java,我当时有点心虚,是这样回复领导:自己也只是有1年多java的工作,不会网络编程,自己搭一个简单的web项目架构还是可以应付的,我只能尽自己的最大的可能做好这件事。

      后来进入公司后,做采集程序是一头雾水,自己之前只是做web的,对web项目的架构以及业务还是比较熟悉,但是关于和硬件之间的通讯还真的是从来都没有接触过。It is the first step that is troublesome。问过一些朋友,也百度过看过一些技术博客,有用socket 连接,多线程后加while死循环来接受做采集,后来慢慢发现用netty这种NIO网络应用框架比较合适。

      接受就是熟悉一下具体的业务,户外设备会主动连接服务器IP,定时向服务器的指定端口发送报文。所以当前的采集程序结合netty的情况就是需要写一个netty的server端,来监听服务器本机的该指定端口,接受户外设备发送过来的报文即可,而不需要再写一个netty的client端。

      首先说明一下,本项目应用的框架:springboot+netty+rabbitMq+mybatis+lombok+logback,数据库选用的是MYSQL。

      开发环境:idea2018+jdk1.8+mysql5.6.35+maven3.5.3

      接下来就是搭建项目:

      1.idea快速创建springboot项目,在pom.xml文件中配置依赖包

    <dependencies>
            <!--web模块starter依赖-->
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-web</artifactId>
            </dependency>
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-test</artifactId>
                <scope>test</scope>
            </dependency>
    
            <!--  netty依赖 springboot2.0自动导入版本-->
            <dependency>
                <groupId>io.netty</groupId>
                <artifactId>netty-all</artifactId>
            </dependency>
            <!--rabbitmq依赖-->
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-amqp</artifactId>
            </dependency>
            <!--mybatis依赖-->
            <dependency>
                <groupId>org.mybatis.spring.boot</groupId>
                <artifactId>mybatis-spring-boot-starter</artifactId>
                <version>1.3.2</version>
            </dependency>
            <!--mysql依赖-->
            <dependency>
                <groupId>mysql</groupId>
                <artifactId>mysql-connector-java</artifactId>
                <scope>runtime</scope>
            </dependency>
    
            <!--pagehelper 分页插件-->
            <dependency>
                <groupId>com.github.pagehelper</groupId>
                <artifactId>pagehelper-spring-boot-starter</artifactId>
                <version>1.2.5</version>
            </dependency>
            <!--lombok 插件-->
            <dependency>
                <groupId>org.projectlombok</groupId>
                <artifactId>lombok</artifactId>
                <version>1.16.6</version>
            </dependency>
            <dependency>
                <groupId>org.junit.jupiter</groupId>
                <artifactId>junit-jupiter-api</artifactId>
                <version>RELEASE</version>
                <scope>compile</scope>
            </dependency>
            <!--编写更少量的代码:使用apache commons工具类库:
            https://www.cnblogs.com/ITtangtang/p/3966955.html-->
            <!--apache.commons.lang3-->
            <dependency>
                <groupId>org.apache.commons</groupId>
                <artifactId>commons-lang3</artifactId>
            </dependency>
            <!--你可以把这个工具看成是java.util的扩展-->
            <dependency>
                <groupId>org.apache.commons</groupId>
                <artifactId>commons-collections4</artifactId>
                <version>${commons-collections4.version}</version>
            </dependency>
            <!--apache.codec:编码方法的工具类包
            https://blog.csdn.net/u012881904/article/details/52767853-->
            <dependency>
                <groupId>commons-codec</groupId>
                <artifactId>commons-codec</artifactId>
            </dependency>
        </dependencies>

    2.自定义netty的server端BootNettyServer类

    /**
    * netty的server类
    */
    @Component
    @Slf4j
    public class BootNettyServer {
    @Autowired
    BootNettyInitializer bootNettyInitializer;

    public void run(int port) throws Exception {
    /**
    * 配置服务端的NIO线程组
    * NioEventLoopGroup 是用来处理I/O操作的Reactor线程组
    * bossGroup:用来接收进来的连接,workerGroup:用来处理已经被接收的连接,进行socketChannel的网络读写,
    * bossGroup接收到连接后就会把连接信息注册到workerGroup
    * workerGroup的EventLoopGroup默认的线程数是CPU核数的二倍
    */
    EventLoopGroup bossGroup = new NioEventLoopGroup(1);
    EventLoopGroup workerGroup = new NioEventLoopGroup();
    try {
    ServerBootstrap serverBootstrap = new ServerBootstrap();//ServerBootstrap 是一个启动NIO服务的辅助启动类
    serverBootstrap
    .group(bossGroup, workerGroup)//设置group,将bossGroup, workerGroup线程组传递到ServerBootstrap
    .channel(NioServerSocketChannel.class)//ServerSocketChannel是以NIO的selector为基础进行实现的,用来接收新的连接,这里告诉Channel通过NioServerSocketChannel获取新的连接
    /**
    * option是设置 bossGroup,childOption是设置workerGroup
    * netty 默认数据包传输大小为1024字节, 设置它可以自动调整下一次缓冲区建立时分配的空间大小,避免内存的浪费 最小 初始化 最大 (根据生产环境实际情况来定)
    * 使用对象池,重用缓冲区
    */
    .option(ChannelOption.RCVBUF_ALLOCATOR, new AdaptiveRecvByteBufAllocator(64, 10496, 1048576))
    .childOption(ChannelOption.RCVBUF_ALLOCATOR, new AdaptiveRecvByteBufAllocator(64, 10496, 1048576))
    .childHandler(bootNettyInitializer);//设置 I/O处理类,主要用于网络I/O事件,记录日志,编码、解码消息
    System.out.println("正在监听8234端口中....");
    ChannelFuture f = serverBootstrap.bind(port).sync();//绑定端口,同步等待成功
    f.channel().closeFuture().sync();//等待服务器监听端口关闭
    } catch (InterruptedException e) {

    } finally {
    /**
    * 退出,释放线程池资源
    */
    bossGroup.shutdownGracefully().sync();
    workerGroup.shutdownGracefully().sync();
    }
    }
    }

    3.自定义初始化类BootNettyInitializer类

      继承ChannelInitializer,实现initChannel方法,初始化信道时可以配置心跳、编码器、解码器、以及业务处理ChannelHandler类。

    /**
     * Channel初始化器
     * @param <SocketChannel>
     */
    @Component
    public class BootNettyInitializer<SocketChannel> extends ChannelInitializer<Channel> {
    
        @Autowired
        BootNettyHandler bootNettyHandler;
        @Autowired
        MyDecoder myDecoder;
        @Override
        protected void initChannel(Channel ch) throws Exception {
            // ChannelOutboundHandler,依照逆序执行
            ch.pipeline().addLast("decoder", myDecoder);
          
            ch.pipeline().addLast(bootNettyHandler);//添加业务处理handler
            System.out.println("信道初始化中....");
        }
    }

    4.自定义业务处理类BootNettyHandler

    import io.netty.buffer.ByteBuf;
    import io.netty.buffer.Unpooled;
    import io.netty.channel.*;
    import io.netty.channel.socket.SocketChannel;
    import io.netty.handler.timeout.IdleStateEvent;
    import lombok.extern.slf4j.Slf4j;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.stereotype.Component;
    import java.io.IOException;
    import java.net.InetSocketAddress;
    import java.text.SimpleDateFormat;
    import java.time.LocalDateTime;
    import java.time.format.DateTimeFormatter;
    import java.util.Date;/**
     * I/O数据读写处理类
     */
    @Component
    @ChannelHandler.Sharable
    @Slf4j
    public class BootNettyHandler extends ChannelInboundHandlerAdapter {
        //  将当前客户端连接存入map,实现长连接,控制设备下发指令
        public  static Map<String, Channel> ctxMap = new LinkedHashMap<String, Channel>();
        /**
         * 从客户端收到新的数据时,这个方法会在收到消息时被调用
         *
         * @param ctx
         * @param msg
         */
        @Override
        public void channelRead(ChannelHandlerContext ctx, Object msg) {
            SocketChannel channel = (SocketChannel) ctx.channel();//将字符串转成每两个字符加空格形式的字符串 7D 7D 21 3C C5 3B 0D 0D 7D 7D 25 3C
            String regex = "(.{2})";
            String input = msg.toString().replaceAll(regex, "$1 ");
         //将报文信息记录到日志中 log.info(channel.remoteAddress().getHostString()
    + ": " + input); System.out.println("服务端接受信息为: " + new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").format(new Date()) + " 接收到消息:" + input); //具体的处理业务逻辑,解析报文帧结构,将数据入库 ...... } /** * 从客户端收到新的数据、读取完成时调用 * @param ctx */ @Override public void channelReadComplete(ChannelHandlerContext ctx) { ctx.flush(); } /** * 当出现 Throwable 对象才会被调用,即当 Netty 由于 IO 错误或者处理器在处理事件时抛出的异常时 * @param ctx * @param cause */ @Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { System.out.println("程序异常,断开客户端连接"); //获取客户端的请求地址 取到的值为客户端的 ip+端口号 InetSocketAddress insocket = (InetSocketAddress) ctx.channel().remoteAddress(); String clientIp = insocket.getAddress().getHostAddress();//设备IP地址(个人将设备的IP地址当作 map 的key,channel对象做value) if(ctxMap.get(clientIp)!=null){//如果不为空就剔除 ctxMap.remove(clientIp, ctx.channel()); }else{//否则就将当前的设备ip+端口存进map 当做下发设备的标识的key } cause.printStackTrace(); ctx.close();//抛出异常,断开与客户端的连接 log.info(clientIp+"连接断开"); } /** * 客户端与服务端第一次建立连接时执行 * @param ctx * @throws Exception */ @Override public void channelActive(ChannelHandlerContext ctx) { InetSocketAddress insocket = (InetSocketAddress) ctx.channel().remoteAddress(); String clientIp = insocket.getAddress().getHostAddress(); if(ctxMap.get(clientIp)!=null){//如果不为空就不存 }else{//否则就将当前的设备ip+端口存进map 当做下发设备的标识的key ctxMap.put(clientIp, ctx.channel()); } log.info(clientIp+"连接"); } /** * 客户端与服务端断连时执行 * @param ctx * @throws Exception */ @Override public void channelInactive(ChannelHandlerContext ctx) throws Exception { super.channelInactive(ctx); InetSocketAddress insocket = (InetSocketAddress) ctx.channel().remoteAddress(); String clientIp = insocket.getAddress().getHostAddress(); if(ctxMap.get(clientIp)!=null){//如果不为空就剔除 ctxMap.remove(clientIp, ctx.channel()); }else{//否则就将当前的设备ip+端口存进map 当做下发设备的标识的key } ctx.close(); //断开连接时,必须关闭,否则造成资源浪费,并发量很大情况下可能造成宕机 log.info(clientIp+"断开"); } /** * 服务端当read超时, 会调用这个方法 * @param ctx * @param evt * @throws Exception */ @Override public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception, IOException { if (evt instanceof IdleStateEvent) { IdleStateEvent event = (IdleStateEvent) evt; String eventType = null; switch (event.state()) { case READER_IDLE: eventType = "读空闲"; break; case WRITER_IDLE: eventType = "些空闲"; break; case ALL_IDLE: eventType = "读写空闲"; break; default: } String date = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")); log.info(date + " " + ctx.channel().remoteAddress() + " " + eventType); } } }

    回到开始位置,自定义netty的server端BootNettyServer类该如何启动?

      我的方案是在application启动类中来启动。

    application类实现CommandLineRunner,来启动nettyserver服务

    import org.mybatis.spring.annotation.MapperScan;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.beans.factory.annotation.Value;
    import org.springframework.boot.CommandLineRunner;
    import org.springframework.boot.SpringApplication;
    import org.springframework.boot.autoconfigure.SpringBootApplication;
    
    @MapperScan("com.rtstjkx.jkxlistener.repository")
    @SpringBootApplication
    public class KjxApplication implements CommandLineRunner {
    
        @Value("${netty.port}")
        private Integer  port;
        @Autowired
        BootNettyServer bootNettyServer;
        public static void main(String[] args) throws Exception {
            SpringApplication.run(KjxApplication.class, args);
        }
    
        @Override
        public void run(String... args) throws Exception {
            /**
             * 启动netty服务端服务
             */
            bootNettyServer.run(port);
        }
    }

    至此,采集程序完成。

      此程序仅个人学习理解后完成,如有不足的地方,还请进一步多多给出指导,thank you !!!

    代码git:git@github.com:white66/jkxLinstener.git

      

    时光静好,与君语;细水长流,与君同;繁华落尽,与君老!
  • 相关阅读:
    【java虚拟机】垃圾回收机制详解
    【java虚拟机】分代垃圾回收策略的基础概念
    【java虚拟机】内存分配与回收策略
    【java虚拟机】jvm内存模型
    【转】新说Mysql事务隔离级别
    【转】互联网项目中mysql应该选什么事务隔离级别
    有关PHP的字符串知识
    php的查询数据
    php练习题:投票
    php的数据访问
  • 原文地址:https://www.cnblogs.com/lyzj/p/13278414.html
Copyright © 2011-2022 走看看