Selector.select()跨线程注册可读通道后仍无法解除阻塞问题
问题描述
我正在开发一个支持多连接的服务器,创建了两个Selector:一个用于ServerSocketChannel执行accept操作,另一个处理连接的数据读取。负责accept的Selector能正常通过阻塞的select()接收新连接,但将接收的SocketChannel注册到处于select()阻塞状态的读Selector后,即使发送数据,读Selector的select()仍无法解除阻塞。
服务器代码
import java.io.IOException; import java.net.InetSocketAddress; 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.Iterator; import java.util.Set; import java.util.TreeSet; public class Server implements Runnable { public Thread threadAccept; public Thread threadRead; protected ServerSocketChannel serverSocket; protected Selector selectorAccept; protected Selector selectorIO; protected Set<Connection> connections; private ByteBuffer buffer; public Server() { try { selectorAccept = Selector.open(); selectorIO = Selector.open(); serverSocket = ServerSocketChannel.open(); serverSocket.configureBlocking(false); InetSocketAddress hostAddress = new InetSocketAddress("localhost",4444); serverSocket.bind(hostAddress); int ops = serverSocket.validOps(); SelectionKey selectKey = serverSocket.register(selectorAccept, ops); threadAccept = new Thread(this); threadAccept.start(); threadRead = new Thread(this); threadRead.start(); } catch(Exception e) { System.out.println("Error "+e); } } public void run() { if(Thread.currentThread() == threadAccept) { acceptNewConnections(); } else if(Thread.currentThread() == threadRead) { readData(); } } private void acceptNewConnections() { int numberOfKeys = 0; while(true) { try { numberOfKeys = selectorAccept.select(); Set<SelectionKey> keys = selectorAccept.selectedKeys(); Iterator<SelectionKey> itr = keys.iterator(); while(itr.hasNext()) { SelectionKey key = itr.next(); if(key.isAcceptable()) { SocketChannel client = serverSocket.accept(); client.configureBlocking(false); client.register(selectorIO, SelectionKey.OP_READ); System.out.println("New connection"); } } } catch(Exception e) { System.out.println(e); } } } public void readData() { int numberOfKeys = 0; ByteBuffer buffer = ByteBuffer.allocate(256); while(true) { try { System.out.println("About to block on IO selector"); numberOfKeys = selectorIO.select(); System.out.println("I NEVER GET HERE"); Set<SelectionKey> keys = selectorIO.selectedKeys(); Iterator<SelectionKey> itr = keys.iterator(); while(itr.hasNext()) { SelectionKey key = itr.next(); if(key.isReadable()) { SocketChannel channel = (SocketChannel)key.channel(); channel.read(buffer); String s = buffer.toString(); System.out.println(s); } } } catch(Exception e) { System.out.println(e); } } } }
主类代码(启动服务器并创建客户端)
import java.io.IOException; import java.net.Socket; import java.net.UnknownHostException; public class Main { public static void main(String[] args) { Server s = new Server(); Thread t = new Thread(new Runnable() { public void run() { try { Socket s = new Socket("localhost", 4444); byte[] data = "hello".getBytes(); s.getOutputStream().write(data); s.getOutputStream().flush(); } catch (UnknownHostException e) { e.printStackTrace(); } catch (IOException e) { e.printStackTrace(); } } }); t.start(); } }
解决方案
问题出在三个关键细节上,以下是修复方案:
必须移除已处理的SelectionKey
Selector的selectedKeys集合是累积的,每次select()会将就绪的Key添加进去,如果不手动移除已处理的Key,下一次select()会重复处理这些Key,导致Selector无法正确响应新事件。处理每个Key后需调用itr.remove()。跨线程注册Channel后要唤醒Selector
当在accept线程向处于阻塞状态的读Selector注册新Channel时,Selector不会自动感知新注册的Key,需调用selectorIO.wakeup()唤醒阻塞的select(),使其重新检查就绪状态。ByteBuffer读取后需切换模式并重置
读取数据到ByteBuffer后,直接调用toString()会读取整个缓冲区(包括空字节),需先调用buffer.flip()切换到读模式,再转换为字符串,最后用buffer.clear()重置缓冲区以便下次读取。
修复后的核心代码
修复后的acceptNewConnections()方法
private void acceptNewConnections() { int numberOfKeys = 0; while(true) { try { numberOfKeys = selectorAccept.select(); Set<SelectionKey> keys = selectorAccept.selectedKeys(); Iterator<SelectionKey> itr = keys.iterator(); while(itr.hasNext()) { SelectionKey key = itr.next(); itr.remove(); // 移除已处理的Key if(key.isAcceptable()) { SocketChannel client = serverSocket.accept(); client.configureBlocking(false); client.register(selectorIO, SelectionKey.OP_READ); selectorIO.wakeup(); // 唤醒读Selector System.out.println("New connection"); } } } catch(Exception e) { System.out.println(e); } } }
修复后的readData()方法
public void readData() { int numberOfKeys = 0; ByteBuffer buffer = ByteBuffer.allocate(256); while(true) { try { System.out.println("About to block on IO selector"); numberOfKeys = selectorIO.select(); System.out.println("Selector woke up, ready keys: " + numberOfKeys); Set<SelectionKey> keys = selectorIO.selectedKeys(); Iterator<SelectionKey> itr = keys.iterator(); while(itr.hasNext()) { SelectionKey key = itr.next(); itr.remove(); // 移除已处理的Key if(key.isReadable()) { SocketChannel channel = (SocketChannel)key.channel(); int bytesRead = channel.read(buffer); if(bytesRead > 0) { buffer.flip(); // 切换到读模式 String s = new String(buffer.array(), 0, buffer.limit()); System.out.println("Received data: " + s); buffer.clear(); // 重置缓冲区 } else if(bytesRead == -1) { // 客户端关闭连接 key.cancel(); channel.close(); System.out.println("Connection closed"); } } } } catch(Exception e) { System.out.println(e); } } }
内容的提问来源于stack exchange,提问作者Matthew
相关产品推荐
相关产品推荐

