• Java I/O(4):AIO和NIO中的Selector


    您好,我是湘王,这是我的CSDN博客,欢迎您来,欢迎您再来~

    在Java NIO的三大核心中,除了Channel和Buffer,剩下的就是Selector了。有的地方叫它选择器,也有叫多路复用器的(比如Netty)。

    之前提过,数据总是从Channel读取到Buffer,或者从Buffer写入到Channel,单个线程可以监听多个Channel——Selector就是这个线程背后的实现机制(所以得名Selector)。

    Selector通过控制单个线程处理多个Channel,如果应用打开了多个Channel,但每次传输的流量都很低,使用Selector就会很方便(至于为什么,具体到Netty中再分析)。所以使用Selector的好处就显而易见:用最少的资源实现最多的操作,避免了线程切换带来的开销。

    还是以代码为例来演示Selector的作用。新建一个类,在main()方法中输入下面的代码:

    1. public static void main(String args[]) throws IOException {
    2. // 创建ServerSocketChannel
    3. ServerSocketChannel channel1 = ServerSocketChannel.open();
    4. channel1.socket().bind(new InetSocketAddress("127.0.0.1", 8080));
    5. channel1.configureBlocking(false);
    6. ServerSocketChannel channel2 = ServerSocketChannel.open();
    7. channel2.socket().bind(new InetSocketAddress("127.0.0.1", 9090));
    8. channel2.configureBlocking(false);
    9. // 创建一个Selector对象
    10. Selector selector = Selector.open();
    11. // 按照字面意思理解,应该是这样的:selector.register(channel, event);
    12. // 但其实是这样的:channel.register(selector, SelectionKey.OP_READ);
    13. // 四种监听事件:
    14. // OP_CONNECT(连接就绪)
    15. // OP_ACCEPT(接收就绪)
    16. // OP_READ(读就绪)
    17. // OP_WRITE(写就绪)
    18. // 注册Channel到Selector,事件一旦被触发,监听随之结束
    19. SelectionKey key1 = channel1.register(selector, SelectionKey.OP_ACCEPT);
    20. SelectionKey key2 = channel2.register(selector, SelectionKey.OP_ACCEPT);
    21. // 模板代码:在编写程序时,大多数时间都是在模板代码中添加相应的业务代码
    22. while(true) {
    23. int readyNum = selector.select();
    24. if (readyNum == 0) {
    25. continue;
    26. }
    27. Set selectedKeys = selector.selectedKeys();
    28. // 轮询
    29. for (SelectionKey key : selectedKeys) {
    30. Channel channel = key.channel();
    31. if (key.isConnectable()) {
    32. if (channel == channel1) {
    33. System.out.println("channel1连接就绪");
    34. } else {
    35. System.out.println("channel2连接就绪");
    36. }
    37. } else if (key.isAcceptable()) {
    38. if (channel == channel1) {
    39. System.out.println("channel1接收就绪");
    40. } else {
    41. System.out.println("channel2接收就绪");
    42. }
    43. }
    44. // 触发后删除,这里不删
    45. // it.remove();
    46. }
    47. }
    48. }

    代码写好后启动ServerSocketChannel服务,可以看到我这里已经启动成功:

     

    然后在网上下载一个叫做SocketTest.jar的工具(在一些工具网站下载的时候当心中毒,如果不放心,可以私信我,给你地址),双击打开,并按下图方式执行:

     

    点击「Connect」可以看到变化:

     

    然后点击「Disconnect」,再输入「9090」后,再点击「Connect」试试:

     

    可以看到结果显示结果变了:

     

    两次连接,打印了三条信息:说明selector的轮询在起作用(因为Set中包含了所有处于监听的SelectionKey)。但是「接收就绪」监听事件仅执行了一次就再不响应。如果感兴趣的话你可以把OP_READ、OP_WRITE这些事件也执行一下试试看。

    因为Selector是单线程轮询监听多个Channel,那么如果Selector(线程)之间需要传递数据,怎么办呢?——Pipe登场了。Pipe就是一种用于Selector之间数据传递的「管道」。

    先来看个图:

    可以清楚地看到它的工作方式。

    还是用代码来解释。

    1. public static void main(String args[]) throws IOException {
    2. // 打开管道
    3. Pipe pipe = Pipe.open();
    4. // 将Buffer数据写入到管道
    5. Pipe.SinkChannel sinkChannel = pipe.sink();
    6. ByteBuffer buffer = ByteBuffer.allocate(32);
    7. buffer.put("ByteBuffer".getBytes());
    8. // 切换到写模式
    9. buffer.flip();
    10. sinkChannel.write(buffer);
    11. // 从管道读取数据
    12. Pipe.SourceChannel sourceChannel = pipe.source();
    13. buffer = ByteBuffer.allocate(32);
    14. sourceChannel.read(buffer);
    15. System.out.println(new String(buffer.array()));
    16. // 关闭管道
    17. sinkChannel.close();
    18. sourceChannel.close();
    19. }

    之前说过,同步指的按顺序一次完成一个任务,直到前一个任务完成并有了结果以后,才能再执行后面的任务。而异步指的是前一个任务结束后,并不等待任务结果,而是继续执行后一个任务,在所有任务都「执行」完后,通过任务的回调函数去获得结果。所以异步使得应用性能有了极大的提高。为了更加生动地说明什么是异步,可以来做个实验:

    通过调用CompletableFuture.supplyAsync()方法可以很明显地观察到,处于位置2的「这一步先执行」会最先显示,然后才执行位置1的代码。而这就是异步的具体实现。

    NIO为了支持异步,升级到了NIO2,也就是AIO。而AIO引入了新的异步Channel的概念,并提供了异步FileChannel和异步SocketChannel的实现。AIO的异步SocketChannel是真正的异步非阻塞I/O。通过代码可以更好地说明:

    1. /**
    2. * AIO客户端
    3. *
    4. * @author xiangwang
    5. */
    6. public class AioClient {
    7. public void start() throws IOException, InterruptedException {
    8. AsynchronousSocketChannel channel = AsynchronousSocketChannel.open();
    9. if (channel.isOpen()) {
    10. // socket接收缓冲区recbuf大小
    11. channel.setOption(StandardSocketOptions.SO_RCVBUF, 128 * 1024);
    12. // socket发送缓冲区recbuf大小
    13. channel.setOption(StandardSocketOptions.SO_SNDBUF, 128 * 1024);
    14. // 保持长连接状态
    15. channel.setOption(StandardSocketOptions.SO_KEEPALIVE, true);
    16. // 连接到服务端
    17. channel.connect(new InetSocketAddress(8080), null,
    18. new AioClientHandler(channel));
    19. // 阻塞主进程
    20. for(;;) {
    21. TimeUnit.SECONDS.sleep(1);
    22. }
    23. } else {
    24. throw new RuntimeException("Channel not opened!");
    25. }
    26. }
    27. public static void main(String[] args) throws IOException, InterruptedException {
    28. new AioClient().start();
    29. }
    30. }

     

    1. /**
    2. * AIO客户端CompletionHandler
    3. *
    4. * @author xiangwang
    5. */
    6. public class AioClientHandler implements CompletionHandler {
    7. private final AsynchronousSocketChannel channel;
    8. private final CharsetDecoder decoder = Charset.defaultCharset().newDecoder();
    9. private final BufferedReader input = new BufferedReader(new InputStreamReader(System.in));
    10. public AioClientHandler(AsynchronousSocketChannel channel) {
    11. this.channel = channel;
    12. }
    13. @Override
    14. public void failed(Throwable exc, AioClient attachment) {
    15. throw new RuntimeException("channel not opened!");
    16. }
    17. @Override
    18. public void completed(Void result, AioClient attachment) {
    19. System.out.println("send message to server: ");
    20. try {
    21. // 将输入内容写到buffer
    22. String line = input.readLine();
    23. channel.write(ByteBuffer.wrap(line.getBytes()));
    24. // 在操作系统中的Java本地方法native已经把数据写到了buffer中
    25. // 这里只需要一个缓冲区能接收就行了
    26. ByteBuffer buffer = ByteBuffer.allocate(1024);
    27. while (channel.read(buffer).get() != -1) {
    28. buffer.flip();
    29. System.out.println("from server: " + decoder.decode(buffer).toString());
    30. if (buffer.hasRemaining()) {
    31. buffer.compact();
    32. } else {
    33. buffer.clear();
    34. }
    35. // 将输入内容写到buffer
    36. line = input.readLine();
    37. channel.write(ByteBuffer.wrap(line.getBytes()));
    38. }
    39. } catch (IOException | InterruptedException | ExecutionException e) {
    40. e.printStackTrace();
    41. }
    42. }
    43. }

     

    1. /**
    2. * AIO服务端
    3. *
    4. * @author xiangwang
    5. */
    6. public class AioServer {
    7. public void start() throws InterruptedException, IOException {
    8. AsynchronousServerSocketChannel channel = AsynchronousServerSocketChannel.open();
    9. if (channel.isOpen()) {
    10. // socket接受缓冲区recbuf大小
    11. channel.setOption(StandardSocketOptions.SO_RCVBUF, 4 * 1024);
    12. // 端口重用,防止进程意外终止,未释放端口,重启时失败
    13. // 因为直接杀进程,没有显式关闭套接字来释放端口,会等待一段时间后才可以重新use这个关口
    14. // 解决办法就是用SO_REUSEADDR
    15. channel.setOption(StandardSocketOptions.SO_REUSEADDR, true);
    16. channel.bind(new InetSocketAddress(8080));
    17. } else {
    18. throw new RuntimeException("channel not opened!");
    19. }
    20. // 处理client连接
    21. channel.accept(null, new AioServerHandler(channel));
    22. System.out.println("server started");
    23. // 阻塞主进程
    24. for(;;) {
    25. TimeUnit.SECONDS.sleep(1);
    26. }
    27. }
    28. public static void main(String[] args) throws IOException, InterruptedException {
    29. AioServer server = new AioServer();
    30. server.start();
    31. }
    32. }

    1. /**
    2. * AIO服务端CompletionHandler
    3. *
    4. * @author xiangwang
    5. */
    6. public class AioServerHandler implements CompletionHandler {
    7. private final AsynchronousServerSocketChannel serverChannel;
    8. private final CharsetDecoder decoder = Charset.defaultCharset().newDecoder();
    9. private final BufferedReader input = new BufferedReader(new InputStreamReader(System.in));
    10. public AioServerHandler(AsynchronousServerSocketChannel serverChannel) {
    11. this.serverChannel = serverChannel;
    12. }
    13. @Override
    14. public void failed(Throwable exc, Void attachment) {
    15. // 处理下一次的client连接
    16. serverChannel.accept(null, this);
    17. }
    18. @Override
    19. public void completed(AsynchronousSocketChannel result, Void attachment) {
    20. // 处理下一次的client连接,类似链式调用
    21. serverChannel.accept(null, this);
    22. try {
    23. // 将输入内容写到buffer
    24. String line = input.readLine();
    25. result.write(ByteBuffer.wrap(line.getBytes()));
    26. // 在操作系统中的Java本地方法native已经把数据写到了buffer中
    27. // 这里只需要一个缓冲区能接收就行了
    28. ByteBuffer buffer = ByteBuffer.allocate(1024);
    29. while (result.read(buffer).get() != -1) {
    30. buffer.flip();
    31. System.out.println("from client: " + decoder.decode(buffer).toString());
    32. if (buffer.hasRemaining()) {
    33. buffer.compact();
    34. } else {
    35. buffer.clear();
    36. }
    37. // 将输入内容写到buffer
    38. line = input.readLine();
    39. result.write(ByteBuffer.wrap(line.getBytes()));
    40. }
    41. } catch (InterruptedException | ExecutionException | IOException e) {
    42. e.printStackTrace();
    43. }
    44. }
    45. }

    执行测试后显示,不管是在客户端还是在服务端,读写完全是异步的。


    感谢您的大驾光临!咨询技术、产品、运营和管理相关问题,请关注后留言。欢迎骚扰,不胜荣幸~

  • 相关阅读:
    【软件测试笔试题】阿里巴巴(中国)网络技术有限公司
    (完全解决)pycharm运行或者调试项目的时候报错:test setup failed
    3.4 常用操作
    ROS2踩坑记录
    大学生餐饮主题网页制作 美食网页设计模板 学生静态网页作业成品 dreamweaver美食甜品蛋糕HTML网站制作
    DevEco Hvigor高效编译,构建过程新秘籍
    华为云 云证书(SSL)管理服务
    C. Joyboard
    数组元素全排列Java
    BertTokenizer 使用方法
  • 原文地址:https://blog.csdn.net/lostrex/article/details/127436894