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

非阻塞Socket服务器消息无法被客户端接收及相关技术疑问

基于非阻塞Socket的消息队列问题排查与技术解答

我用非阻塞Socket实现了一个简易消息队列,预期逻辑是写客户端发送消息时,服务器将消息转发给读客户端,但读客户端始终无法收到消息,同时有以下三个技术疑问需要解答:


服务器代码

public class NonBlockingMessageQueue implements MessageQueue {

    private final Map<SelectionKey, LocalDateTime> keyMap = new HashMap<>();
    private final List<String> messages = new ArrayList<>();

    private Selector selector;
    private ServerSocketChannel serverSocketChannel;

    @Override
    public void start(int port) {
        try {
            initSelector();
            listenToSockets();
        } catch (IOException e) {
            throw new RuntimeException(e);
        }
    }

    private void initSelector() throws IOException {
        this.selector = Selector.open();
        this.serverSocketChannel = ServerSocketChannel.open();
        serverSocketChannel.configureBlocking(false);

        InetAddress ip = InetAddress.getByName("localhost");
        serverSocketChannel.bind(new InetSocketAddress(ip, 1234));
        serverSocketChannel.register(selector, SelectionKey.OP_ACCEPT);
    }

    private void listenToSockets() throws IOException {
        while (true) {
            if (selector.select() <= 0)
                continue;
            Set<SelectionKey> selectedKeys = selector.selectedKeys();
            Iterator<SelectionKey> iterator = selectedKeys.iterator();
            while (iterator.hasNext()) {
                SelectionKey key = iterator.next();
                iterator.remove();
                keyMap.put(key, LocalDateTime.now());
                if (key.isAcceptable()) {
                    SocketChannel sc = serverSocketChannel.accept();
                    sc.configureBlocking(false);
                    sc.register(selector, SelectionKey.OP_READ | SelectionKey.OP_WRITE);
                    System.out.println("Connection Accepted: " + sc.getLocalAddress() + "n");
                }

                if (key.isReadable()) {
                    boolean exists = keyMap.containsKey(key);
                    SocketChannel sc = (SocketChannel) key.channel();
                    ByteBuffer bb = ByteBuffer.allocate(1024);
                    sc.read(bb);
                    String result = new String(bb.array()).trim();
                    messages.add(result);


                    for( SelectionKey channel : keyMap.keySet()) {
                        try {
                            SelectableChannel chan = channel.channel();
                            if( (chan.validOps() & SelectionKey.OP_WRITE ) > 0) {
                                ByteBuffer bbb = ByteBuffer.wrap(messages.get(0).getBytes());
                                                             ((SocketChannel)chan).write(bbb);
                            }
                        }catch (Throwable e ){
                            System.out.println("Error " + e);
                        }
                    }
                }
            }
        }
    }
}

写客户端代码

public class NonBlockingClientWrite {
    private static BufferedReader input = null;
    public static void main(String[] args) throws Exception {
        InetSocketAddress addr = new InetSocketAddress(
                InetAddress.getByName("localhost"), 1234);
        Selector selector = Selector.open();
        SocketChannel sc = SocketChannel.open();
        sc.configureBlocking(false);
        sc.connect(addr);
        sc.register(selector, SelectionKey.OP_CONNECT | SelectionKey.OP_WRITE);
        input = new BufferedReader(new InputStreamReader(System.in));
        while (true) {
            if (selector.select() > 0) {
                Boolean doneStatus = processReadySet
                        (selector.selectedKeys());
                if (doneStatus) {
                    break;
                }
            }
        }
        sc.close();
    }
    public static Boolean processReadySet(Set readySet)
            throws Exception {
        SelectionKey key = null;
        Iterator iterator = null;
        iterator = readySet.iterator();
        while (iterator.hasNext()) {
            key = (SelectionKey) iterator.next();
            iterator.remove();
        }
        if (key.isConnectable()) {
            Boolean connected = processConnect(key);
            if (!connected) {
                return true;
            }
        }
        if (key.isWritable()) {
            System.out.print("Type a message (type quit to stop): ");
            String msg = input.readLine();
            if (msg.equalsIgnoreCase("quit")) {
                return true;
            }
            SocketChannel sc = (SocketChannel) key.channel();
            ByteBuffer bb = ByteBuffer.wrap(msg.getBytes());
            sc.write(bb);
        }
        return false;
    }
    public static Boolean processConnect(SelectionKey key) {
        SocketChannel sc = (SocketChannel) key.channel();
        try {
            while (sc.isConnectionPending()) {
                sc.finishConnect();
            }
        } catch (IOException e) {
            key.cancel();
            e.printStackTrace();
            return false;
        }
        return true;
    }
}

读客户端代码

public class NonBlockingClientRead {
    private static BufferedReader input = null;
    public static void main(String[] args) throws Exception {
        InetSocketAddress addr = new InetSocketAddress(InetAddress.getByName("localhost"), 1234);
        Selector selector = Selector.open();
        SocketChannel sc = SocketChannel.open();
        sc.configureBlocking(false);
        sc.connect(addr);
        sc.register(selector, SelectionKey.OP_CONNECT | SelectionKey.OP_READ);
        while (true) {
            if (selector.select() > 0) {
                Boolean doneStatus = processReadySet(selector.selectedKeys());
                if (doneStatus) {
                    break;
                }
            }
        }
        sc.close();
    }
    public static Boolean processReadySet(Set readySet) throws Exception {
        SelectionKey key = null;
        Iterator iterator = null;
        iterator = readySet.iterator();
        while (iterator.hasNext()) {
            key = (SelectionKey) iterator.next();
            iterator.remove();
        }
        if (key.isConnectable()) {
            Boolean connected = processConnect(key);
            if (!connected) {
                return true;
            }
        }
        if (key.isReadable()) {
            SocketChannel sc = (SocketChannel) key.channel();
            ByteBuffer bb = ByteBuffer.allocate(1024);
            sc.read(bb);
            String result = new String(bb.array()).trim();
            System.out.println("Message received from Server: " + result + " Message length= "
                    + result.length());
        }
        return false;
    }
    public static Boolean processConnect(SelectionKey key) {
        SocketChannel sc = (SocketChannel) key.channel();
        try {
            while (sc.isConnectionPending()) {
                sc.finishConnect();
            }
        } catch (IOException e) {
            key.cancel();
            e.printStackTrace();
            return false;
        }
        return true;
    }
}

问题排查与解答

主问题:读客户端无法收到消息的原因及修复

  1. 服务器转发逻辑错误:遍历keyMap时包含了ServerSocketChannel的SelectionKey,强转SocketChannel会抛出异常,导致后续转发逻辑中断。需过滤掉ServerSocketChannel类型的通道:
    if (chan instanceof SocketChannel && (chan.validOps() & SelectionKey.OP_WRITE) > 0) {
        // 执行写操作
    }
    
  2. 消息读取与写入的ByteBuffer处理不当:
    • 服务器读取客户端消息时,未根据实际读取的字节数构建字符串,导致包含缓冲区空字符,建议修改为:
      int bytesRead = sc.read(bb);
      if (bytesRead > 0) {
          bb.flip();
          String result = new String(bb.array(), 0, bytesRead).trim();
          messages.add(result);
      }
      
    • 转发消息时,每次只发messages.get(0)且未消费队列中的消息,建议转发后移除已发送的消息,或维护正确的消息分发逻辑。
  3. 读客户端读取逻辑问题:同样需根据实际读取字节数构建字符串,避免空字符干扰:
    int bytesRead = sc.read(bb);
    if (bytesRead > 0) {
        bb.flip();
        String result = new String(bb.array(), 0, bytesRead).trim();
        System.out.println("Message received from Server: " + result + " Message length= " + result.length());
    }
    

技术疑问解答

  1. ServerSocketChannelImpl的来源:
    这是ServerSocketChannel.open()返回的底层实现类,属于sun.nio.ch包下的NIO具体实现。服务器初始化时,serverSocketChannel.register(selector, SelectionKey.OP_ACCEPT)会将它的SelectionKey注册到选择器,处理事件时这个Key会被加入keyMap,因此会出现该类型。

  2. 为何两个客户端的chan.validOps() & OP_WRITE都返回true:
    validOps()返回的是通道本身支持的操作集合,SocketChannel本身就支持读和写操作,所以它的validOps()必然包含OP_READ和OP_WRITE。客户端自己注册的事件(写客户端注册OP_WRITE、读客户端注册OP_READ)是客户端本地选择器的监听配置,和服务器端通道的validOps()无关,因此两者的validOps() & OP_WRITE都会大于0。

  3. isWritable()始终为true的原因:
    isWritable()返回true表示通道的发送缓冲区有空闲空间,可执行写操作而不阻塞。刚建立的连接发送缓冲区通常处于空闲状态,因此会一直返回true;若注册了OP_WRITE事件,选择器会持续触发该事件,直到缓冲区被占满。如果不需要持续监听写事件,写完后应修改SelectionKey的监听事件,取消OP_WRITE注册。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 11:50:54