Kotlin如何检测Socket是否有可读数据?求类似Unix select()的方案
Kotlin中类似Unix select()的实现方案
在Kotlin(基于JVM)中,你可以使用Java NIO的Selector来实现类似Unix select()或Python select.select()的功能,它能监听Socket的可读/可写事件,无需阻塞在单个IO操作上,完美适配你同时收发数据的需求。
核心思路
Selector允许你将非阻塞的SocketChannel注册进去,监听OP_READ(服务端推送数据时触发)和OP_WRITE(Socket可写时触发,用于发送用户请求数据)事件,通过轮询处理这些事件,就能在单个线程中同时处理收发逻辑,避免阻塞。
改造后的实现示例
结合你现有的发送逻辑,改造后的代码大致如下:
import java.nio.ByteBuffer import java.nio.channels.SelectionKey import java.nio.channels.Selector import java.nio.channels.SocketChannel import java.nio.charset.Charset import java.util.concurrent.ConcurrentLinkedQueue // 线程安全的待发送队列 private val toSend = ConcurrentLinkedQueue<String>() private var selector: Selector? = null private var socketChannel: SocketChannel? = null fun startSocketClient(host: String, port: Int) { try { // 初始化SocketChannel并设置为非阻塞模式 socketChannel = SocketChannel.open() socketChannel?.configureBlocking(false) socketChannel?.connect(java.net.InetSocketAddress(host, port)) // 创建Selector selector = Selector.open() // 注册连接事件,连接成功后再注册读/写事件 socketChannel?.register(selector, SelectionKey.OP_CONNECT) // 启动IO处理线程 Thread { while (selector?.isOpen == true) { // 阻塞等待事件触发(可设置超时时间) selector?.select() // 处理所有触发的事件 val selectedKeys = selector?.selectedKeys()?.iterator() while (selectedKeys?.hasNext() == true) { val key = selectedKeys.next() selectedKeys.remove() if (!key.isValid) continue when { key.isConnectable -> { // 完成连接 val channel = key.channel() as SocketChannel if (channel.finishConnect()) { // 连接成功后,注册读事件,写事件按需注册 key.interestOps(SelectionKey.OP_READ) } } key.isReadable -> { // 处理服务端推送的数据 val channel = key.channel() as SocketChannel val buffer = ByteBuffer.allocate(1024) val bytesRead = channel.read(buffer) if (bytesRead > 0) { buffer.flip() val receivedData = Charset.forName("UTF-8").decode(buffer).toString() // 这里将数据展示给用户,比如通过Handler发送到UI线程 // runOnUiThread { updateUI(receivedData) } println("收到服务端数据: $receivedData") } else if (bytesRead == -1) { // 服务端关闭连接 key.cancel() channel.close() } } key.isWritable -> { // 发送待处理的数据 val channel = key.channel() as SocketChannel while (toSend.isNotEmpty()) { val msg = toSend.poll() ?: break // 按照你的协议发送:先发送长度+2,再发送内容 val lengthBytes = (msg.length + 2).toString().toByteArray(Charsets.UTF_8) channel.write(ByteBuffer.wrap(lengthBytes)) val contentBytes = msg.toByteArray(Charsets.UTF_8) channel.write(ByteBuffer.wrap(contentBytes)) } // 待发送队列空了,取消写事件监听,避免频繁触发 if (toSend.isEmpty()) { key.interestOps(key.interestOps() and SelectionKey.OP_WRITE.inv()) } } } } } }.start() } catch (e: Exception) { e.printStackTrace() } } // 调用此方法添加待发送数据,并触发写事件 fun sendMessage(msg: String) { toSend.offer(msg) // 唤醒Selector,立即处理写事件 selector?.wakeup() // 注册写事件(如果之前没注册的话) socketChannel?.keyFor(selector)?.let { key -> if (!key.isValid) return@let key.interestOps(key.interestOps() or SelectionKey.OP_WRITE) } }
关键说明
- 非阻塞模式:必须将
SocketChannel设置为非阻塞,否则Selector无法正常工作。 - 事件监听:
OP_CONNECT:处理Socket连接完成的逻辑。OP_READ:服务端推送数据时触发,此时读取数据即可。OP_WRITE:仅当有待发送数据时才注册该事件,避免无意义的触发。
- 线程安全队列:
toSend使用ConcurrentLinkedQueue,保证多线程添加数据时的安全性。 - 唤醒Selector:当添加新的待发送数据时,调用
selector.wakeup()让Selector立即处理写事件,无需等待超时。
替代方案:多线程模式
如果你觉得NIO的Selector上手复杂,也可以用简单的多线程方案:
- 一个线程专门负责读取服务端数据(阻塞读,但不影响发送线程)。
- 另一个线程负责发送数据(和你现有逻辑类似)。
这种方案代码更直观,但对于大量连接场景,Selector的性能更优。
内容的提问来源于stack exchange,提问作者Atom
相关产品推荐
相关产品推荐

