非阻塞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; } }
问题排查与解答
主问题:读客户端无法收到消息的原因及修复
- 服务器转发逻辑错误:遍历
keyMap时包含了ServerSocketChannel的SelectionKey,强转SocketChannel会抛出异常,导致后续转发逻辑中断。需过滤掉ServerSocketChannel类型的通道:if (chan instanceof SocketChannel && (chan.validOps() & SelectionKey.OP_WRITE) > 0) { // 执行写操作 } - 消息读取与写入的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)且未消费队列中的消息,建议转发后移除已发送的消息,或维护正确的消息分发逻辑。
- 服务器读取客户端消息时,未根据实际读取的字节数构建字符串,导致包含缓冲区空字符,建议修改为:
- 读客户端读取逻辑问题:同样需根据实际读取字节数构建字符串,避免空字符干扰:
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()); }
技术疑问解答
ServerSocketChannelImpl的来源:
这是ServerSocketChannel.open()返回的底层实现类,属于sun.nio.ch包下的NIO具体实现。服务器初始化时,serverSocketChannel.register(selector, SelectionKey.OP_ACCEPT)会将它的SelectionKey注册到选择器,处理事件时这个Key会被加入keyMap,因此会出现该类型。为何两个客户端的chan.validOps() & OP_WRITE都返回true:
validOps()返回的是通道本身支持的操作集合,SocketChannel本身就支持读和写操作,所以它的validOps()必然包含OP_READ和OP_WRITE。客户端自己注册的事件(写客户端注册OP_WRITE、读客户端注册OP_READ)是客户端本地选择器的监听配置,和服务器端通道的validOps()无关,因此两者的validOps() & OP_WRITE都会大于0。isWritable()始终为true的原因:
isWritable()返回true表示通道的发送缓冲区有空闲空间,可执行写操作而不阻塞。刚建立的连接发送缓冲区通常处于空闲状态,因此会一直返回true;若注册了OP_WRITE事件,选择器会持续触发该事件,直到缓冲区被占满。如果不需要持续监听写事件,写完后应修改SelectionKey的监听事件,取消OP_WRITE注册。
内容的提问来源于stack exchange,提问作者Johnyb

