You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Java NIO多线程Socket Server异步请求处理问题求助

问题分析

你的NIO服务器当前是单线程Reactor模型,所有IO事件的监听、接收以及业务逻辑处理都在主线程中执行。当handleRead方法中执行Thread.sleep(30000)时,主线程被阻塞,整个Selector事件循环暂停,无法处理后续的客户端请求,这就是第三到第五个请求必须等待第二个请求sleep完成才能被处理的核心原因。

解决方案

要实现真正的异步处理,需要将IO事件处理和耗时业务逻辑处理分离,让主线程专注于监听和处理IO事件,耗时操作交给独立的线程池执行。以下是两种可行的实现方式:

方式一:线程池异步处理耗时业务

直接在handleRead中把耗时的业务逻辑提交到线程池,主线程完成IO读取后立即返回,继续处理其他事件。

修改后的代码示例:

import java.io.IOException;
import java.net.InetSocketAddress;
import java.net.ServerSocket;
import java.nio.ByteBuffer;
import java.nio.channels.SelectionKey;
import java.nio.channels.Selector;
import java.nio.channels.ServerSocketChannel;
import java.nio.channels.SocketChannel;
import java.util.Date;
import java.util.Iterator;
import java.util.Set;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class NIOServer {
    private static Selector selector = null;
    // 根据业务需求调整线程池大小
    private static final ExecutorService businessThreadPool = Executors.newFixedThreadPool(10);

    public static void main(String[] args) {
        try {
            selector = Selector.open();
            ServerSocketChannel serverSocketChannel = ServerSocketChannel.open();
            ServerSocket serverSocket = serverSocketChannel.socket();
            serverSocket.bind(new InetSocketAddress("localhost", 8089));
            serverSocketChannel.configureBlocking(false);
            int ops = serverSocketChannel.validOps();
            serverSocketChannel.register(selector, ops, null);
            
            while (true) {
                selector.select();
                Set<SelectionKey> selectedKeys = selector.selectedKeys();
                Iterator<SelectionKey> i = selectedKeys.iterator();

                while (i.hasNext()) {
                    SelectionKey key = i.next();
                    if (key.isAcceptable()) {
                        handleAccept(serverSocketChannel, key);
                    } else if (key.isReadable()) {
                        handleRead(key);
                    }
                    i.remove();
                }
            }
        } catch (IOException e) {
            e.printStackTrace();
        } finally {
            businessThreadPool.shutdown();
        }
    }

    private static void handleAccept(ServerSocketChannel mySocket, SelectionKey key) throws IOException {
        System.out.println("Connection Accepted..");
        SocketChannel client = mySocket.accept();
        client.configureBlocking(false);
        client.register(selector, SelectionKey.OP_READ);
    }

    private static void handleRead(SelectionKey key) throws IOException {
        SocketChannel client = (SocketChannel) key.channel();
        ByteBuffer buffer = ByteBuffer.allocate(1024);
        int readBytes = client.read(buffer);
        
        if (readBytes > 0) {
            buffer.flip(); // 切换为读模式,避免读取到空字节
            String data = new String(buffer.array(), 0, readBytes).trim();
            System.out.println("Received message: " + data + " at " + new Date());

            // 耗时逻辑提交到线程池
            businessThreadPool.submit(() -> {
                if (data.equalsIgnoreCase("Hello, server! Request 2")) {
                    try {
                        Thread.sleep(30000);
                        System.out.println("Sleep is executed for request: " + data + " at " + new Date());
                    } catch (InterruptedException e) {
                        Thread.currentThread().interrupt();
                        throw new RuntimeException(e);
                    }
                }

                if (data.equalsIgnoreCase("Testing5")) {
                    try {
                        client.close();
                        System.out.println("Connection closed...");
                    } catch (IOException e) {
                        e.printStackTrace();
                    }
                }
            });
        } else if (readBytes == -1) {
            // 客户端主动关闭连接
            client.close();
            key.cancel();
        }
    }
}

方式二:多线程Reactor模型(进阶)

如果需要更高的并发性能,可以采用主从Reactor模型:

  • 主线程(Main Reactor)仅负责监听客户端连接请求,接收连接后将SocketChannel注册到从Reactor线程(Sub Reactor)的Selector上
  • 从Reactor线程负责处理对应连接的后续IO事件,业务逻辑再交给线程池执行
    这种模型能更好地利用多核CPU,适合高并发场景。
关键注意点
  • 线程池大小需根据业务场景调整,避免线程过多导致上下文切换开销过大
  • SocketChannel本身不是线程安全的,异步处理时建议一个连接的所有IO操作由同一个线程负责(比如主从Reactor模型的方式)
  • 如果需要给客户端回写数据,处理完业务后要重新注册OP_WRITE事件,避免阻塞主线程

内容的提问来源于stack exchange,提问作者mc ser

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.19 03:22:37