JAVA NIO non

Java自1.4以后,加入了新IO特性,NIO. 号称new IO. NIO带来了non-blocking特性. 这篇文章主要讲的是如何使用NIO的网络新特性,来构建高性能非阻塞并发服务器.

文章基于个人理解,我也来搞搞NIO.,求指正.

在NIO之前

服务器还是在使用阻塞式的java socket. 以Tomcat最新版本没有开启NIO模式的源码为例, tomcat会accept出来一个socket连接,然后调用processSocket方法来处理socket.

1while(true) { 2.... 3 Socket socket = null; 4 try { 5 // Accept the next incoming connection from the server 6 // socket 7 socket = serverSocketFactory.acceptSocket(serverSocket); 8 } 9... 10... 11 // Configure the socket 12 if (running && !paused && setSocketOptions(socket)) { 13 // Hand this socket off to an appropriate processor 14 if (!processSocket(socket)) { 15 countDownConnection(); 16 // Close socket right away(socket); 17 closeSocket(socket); 18 } 19 } 20.... 21}

使用ServerSocket.accept()方法来创建一个连接. accept方法是阻塞方法,在下一个connection进来之前,accept会阻塞.

在一个socket进来之后,Tomcat会在thread pool里面拿出一个thread来处理连接的socket. 然后自己快速的脱身去接受下一个socket连接. 代码如下:

1protected boolean processSocket(Socket socket) { 2 // Process the request from this socket 3 try { 4 SocketWrapper<Socket> wrapper = new SocketWrapper<Socket>(socket); 5 wrapper.setKeepAliveLeft(getMaxKeepAliveRequests()); 6 // During shutdown, executor may be null - avoid NPE 7 if (!running) { 8 return false; 9 } 10 getExecutor().execute(new SocketProcessor(wrapper)); 11 } catch (RejectedExecutionException x) { 12 log.warn("Socket processing request was rejected for:"+socket,x); 13 return false; 14 } catch (Throwable t) { 15 ExceptionUtils.handleThrowable(t); 16 // This means we got an OOM or similar creating a thread, or that 17 // the pool and its queue are full 18 log.error(sm.getString("endpoint.process.fail"), t); 19 return false; 20 } 21 return true; 22}

而每个处理socket的线程,也总是会阻塞在while(true) sockek.getInputStream().read() 方法上. 

总结就是, 一个socket必须使用一个线程来处理. 致使服务器需要维护比较多的线程. 线程本身就是一个消耗资源的东西,并且每个处理socket的线程都会阻塞在read方法上,使得系统大量资源被浪费.

以上这种socket的服务方式适用于HTTP服务器,每个http请求都是短期的,无状态的,并且http后台的业务逻辑也一般比较复杂. 使用多线程和阻塞方式是合适的.

倘若是做游戏服务器,尤其是CS架构的游戏.这种传统模式服务器毫无胜算.游戏有以下几个特点是传统服务器不能胜任的:
1, 持久TCP连接. 每一个client和server之间都存在一个持久的连接.当CCU(并发用户数量)上升,阻塞式服务器无法为每一个连接运行一个线程.
2, 自己开发的二进制流传输协议. 游戏服务器讲究响应快.那网络传输也要节省时间. HTTP协议的冗余内容太多,一个好的游戏服务器传输协议,可以使得message压缩到3-6倍甚至以上.这就使得游戏服务器要开发自己的协议解析器.
3, 传输双向,且消息传输频率高.假设一个游戏服务器instance连接了2000个client,每个client平均每秒钟传输1-10个message,一个message大约几百字节或者几千字节.而server也需要向client广播其他玩家的当前信息.这使得服务器需要有高速处理消息的能力.
4, CS架构的游戏服务器端的逻辑并不像APP服务器端的逻辑那么复杂. 网络游戏在client端处理了大部分逻辑,server端负责简单逻辑,甚至只是传递消息.

在Java NIO出现以后

出现了使用NIO写的非阻塞网络引擎,比如Apache Mina, JBoss Netty, Smartfoxserver BitSwarm. 比较起来, Mina的性能不如后两者.Tomcat也存在NIO模式,不过需要人工开启.

首先要说明一下, 与App Server的servlet开发模式不一样, 在Mina, Netty和BitSwarm上开发应用程序都是Event Driven的设计模式.Server端会收到Client端的event,Client也会收到Server端的event,Server端与Client端的都要注册各种event的EventHandler来handle event.

用大白话来解释NIO:
1, Buffers, 网络传输字节存放的地方.无论是从channel中取,还是向channel中写,都必须以Buffers作为中间存贮格式.
2, Socket Channels. Channel是网络连接和buffer之间的数据通道.每个连接一个channel.就像之前的socket的stream一样.
3, Selector. 像一个巡警,在一个片区里面不停的巡逻. 一旦发现事件发生,立刻将事件select出来.不过这些事件必须是提前注册在selector上的. select出来的事件打包成SelectionKey.里面包含了事件的发生事件,地点,人物. 如果警察不巡逻,每个街道(socket)分配一个警察(thread),那么一个片区有几条街道,就需要几个警察.但现在警察巡逻了,一个巡警(selector)可以管理所有的片区里面的街道(socketchannel).

以上把警察比作线程,街道比作socket或socketchannel,街道上发生的一切比作stream.把巡警比作selector,引起巡警注意的事件比作selectionKey.

从上可以看出,使用NIO可以使用一个线程,就能维护多个持久TCP连接.

NIO实例

下面给出NIO编写的EchoServer和Client. Client连接server以后,将发送一条消息给server. Server会原封不懂的把消息发送回来.Client再把消息发送回去.Server再发回来.用不休止. 在性能的允许下,Client可以启动任意多.

以下Code涵盖了NIO里面最常用的方法和连接断开诊断.注释也全.

首先是Server的实现. Server端启动了2个线程,connectionBell线程用于巡逻新的连接事件. readBell线程用于读取所有channel的数据. 注解**: Mina采取了同样的做法,只是readBell线程启动的个数等于处理器个数+1.** 由此可见,NIO只需要少量的几个线程就可以维持非常多的并发持久连接.

每当事件发生,会调用dispatch方法去处理event. 一般情况,会使用一个ThreadPool来处理event. ThreadPool的大小可以自定义.但不是越大越好.如果处理event的逻辑比较复杂,比如需要额外网络连接或者复杂数据库查询,那ThreadPool就需要稍微大些.**(猜测)**Smartfoxserver处理上万的并发,也只用到了3-4个线程来dispatch event.

EchoServer

1public class EchoServer { 2 public static SelectorLoop connectionBell; 3 public static SelectorLoop readBell; 4 public boolean isReadBellRunning=false; 5 6 public static void main(String[] args) throws IOException { 7 new EchoServer().startServer(); 8 } 9 10 // 启动服务器 11 public void startServer() throws IOException { 12 // 准备好一个闹钟.当有链接进来的时候响. 13 connectionBell = new SelectorLoop(); 14 15 // 准备好一个闹装,当有read事件进来的时候响. 16 readBell = new SelectorLoop(); 17 18 // 开启一个server channel来监听 19 ServerSocketChannel ssc = ServerSocketChannel.open(); 20 // 开启非阻塞模式 21 ssc.configureBlocking(false); 22 23 ServerSocket socket = ssc.socket(); 24 socket.bind(new InetSocketAddress("localhost",7878)); 25 26 // 给闹钟规定好要监听报告的事件,这个闹钟只监听新连接事件. 27 ssc.register(connectionBell.getSelector(), SelectionKey.OP_ACCEPT); 28 new Thread(connectionBell).start(); 29 } 30 31 // Selector轮询线程类 32 public class SelectorLoop implements Runnable { 33 private Selector selector; 34 private ByteBuffer temp = ByteBuffer.allocate(1024); 35 36 public SelectorLoop() throws IOException { 37 this.selector = Selector.open(); 38 } 39 40 public Selector getSelector() { 41 return this.selector; 42 } 43 44 @Override 45 public void run() { 46 while(true) { 47 try { 48 // 阻塞,只有当至少一个注册的事件发生的时候才会继续. 49 this.selector.select(); 50 51 Set<SelectionKey> selectKeys = this.selector.selectedKeys(); 52 Iterator<SelectionKey> it = selectKeys.iterator(); 53 while (it.hasNext()) { 54 SelectionKey key = it.next(); 55 it.remove(); 56 // 处理事件. 可以用多线程来处理. 57 this.dispatch(key); 58 } 59 } catch (IOException e) { 60 e.printStackTrace(); 61 } catch (InterruptedException e) { 62 e.printStackTrace(); 63 } 64 } 65 } 66 67 public void dispatch(SelectionKey key) throws IOException, InterruptedException { 68 if (key.isAcceptable()) { 69 // 这是一个connection accept事件, 并且这个事件是注册在serversocketchannel上的. 70 ServerSocketChannel ssc = (ServerSocketChannel) key.channel(); 71 // 接受一个连接. 72 SocketChannel sc = ssc.accept(); 73 74 // 对新的连接的channel注册read事件. 使用readBell闹钟. 75 sc.configureBlocking(false); 76 sc.register(readBell.getSelector(), SelectionKey.OP_READ); 77 78 // 如果读取线程还没有启动,那就启动一个读取线程. 79 synchronized(EchoServer.this) { 80 if (!EchoServer.this.isReadBellRunning) { 81 EchoServer.this.isReadBellRunning = true; 82 new Thread(readBell).start(); 83 } 84 } 85 86 } else if (key.isReadable()) { 87 // 这是一个read事件,并且这个事件是注册在socketchannel上的. 88 SocketChannel sc = (SocketChannel) key.channel(); 89 // 写数据到buffer 90 int count = sc.read(temp); 91 if (count < 0) { 92 // 客户端已经断开连接. 93 key.cancel(); 94 sc.close(); 95 return; 96 } 97 // 切换buffer到读状态,内部指针归位. 98 temp.flip(); 99 String msg = Charset.forName("UTF-8").decode(temp).toString(); 100 System.out.println("Server received ["+msg+"] from client address:" + sc.getRemoteAddress()); 101 102 Thread.sleep(1000); 103 // echo back. 104 sc.write(ByteBuffer.wrap(msg.getBytes(Charset.forName("UTF-8")))); 105 106 // 清空buffer 107 temp.clear(); 108 } 109 } 110 111 } 112 113}

接下来就是Client的实现.Client可以用传统IO,也可以使用NIO.这个例子使用的NIO,单线程.

1public class Client implements Runnable { 2 // 空闲计数器,如果空闲超过10次,将检测server是否中断连接. 3 private static int idleCounter = 0; 4 private Selector selector; 5 private SocketChannel socketChannel; 6 private ByteBuffer temp = ByteBuffer.allocate(1024); 7 8 public static void main(String[] args) throws IOException { 9 Client client= new Client(); 10 new Thread(client).start(); 11 //client.sendFirstMsg(); 12 } 13 14 public Client() throws IOException { 15 // 同样的,注册闹钟. 16 this.selector = Selector.open(); 17 18 // 连接远程server 19 socketChannel = SocketChannel.open(); 20 // 如果快速的建立了连接,返回true.如果没有建立,则返回false,并在连接后出发Connect事件. 21 Boolean isConnected = socketChannel.connect(new InetSocketAddress("localhost", 7878)); 22 socketChannel.configureBlocking(false); 23 SelectionKey key = socketChannel.register(selector, SelectionKey.OP_READ); 24 25 if (isConnected) { 26 this.sendFirstMsg(); 27 } else { 28 // 如果连接还在尝试中,则注册connect事件的监听. connect成功以后会出发connect事件. 29 key.interestOps(SelectionKey.OP_CONNECT); 30 } 31 } 32 33 public void sendFirstMsg() throws IOException { 34 String msg = "Hello NIO."; 35 socketChannel.write(ByteBuffer.wrap(msg.getBytes(Charset.forName("UTF-8")))); 36 } 37 38 @Override 39 public void run() { 40 while (true) { 41 try { 42 // 阻塞,等待事件发生,或者1秒超时. num为发生事件的数量. 43 int num = this.selector.select(1000); 44 if (num ==0) { 45 idleCounter ++; 46 if(idleCounter >10) { 47 // 如果server断开了连接,发送消息将失败. 48 try { 49 this.sendFirstMsg(); 50 } catch(ClosedChannelException e) { 51 e.printStackTrace(); 52 this.socketChannel.close(); 53 return; 54 } 55 } 56 continue; 57 } else { 58 idleCounter = 0; 59 } 60 Set<SelectionKey> keys = this.selector.selectedKeys(); 61 Iterator<SelectionKey> it = keys.iterator(); 62 while (it.hasNext()) { 63 SelectionKey key = it.next(); 64 it.remove(); 65 if (key.isConnectable()) { 66 // socket connected 67 SocketChannel sc = (SocketChannel)key.channel(); 68 if (sc.isConnectionPending()) { 69 sc.finishConnect(); 70 } 71 // send first message; 72 this.sendFirstMsg(); 73 } 74 if (key.isReadable()) { 75 // msg received. 76 SocketChannel sc = (SocketChannel)key.channel(); 77 this.temp = ByteBuffer.allocate(1024); 78 int count = sc.read(temp); 79 if (count<0) { 80 sc.close(); 81 continue; 82 } 83 // 切换buffer到读状态,内部指针归位. 84 temp.flip(); 85 String msg = Charset.forName("UTF-8").decode(temp).toString(); 86 System.out.println("Client received ["+msg+"] from server address:" + sc.getRemoteAddress()); 87 88 Thread.sleep(1000); 89 // echo back. 90 sc.write(ByteBuffer.wrap(msg.getBytes(Charset.forName("UTF-8")))); 91 92 // 清空buffer 93 temp.clear(); 94 } 95 } 96 } catch (IOException e) { 97 e.printStackTrace(); 98 } catch (InterruptedException e) { 99 e.printStackTrace(); 100 } 101 } 102 } 103 104}

下载以后黏贴到eclipse中, 先运行EchoServer,然后可以运行任意多的Client. 停止Server和client的方式就是直接terminate server.

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )

JAVA NIO non - HelloWorld