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

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)
    }
}

关键说明

  1. 非阻塞模式:必须将SocketChannel设置为非阻塞,否则Selector无法正常工作。
  2. 事件监听:
    • OP_CONNECT:处理Socket连接完成的逻辑。
    • OP_READ:服务端推送数据时触发,此时读取数据即可。
    • OP_WRITE:仅当有待发送数据时才注册该事件,避免无意义的触发。
  3. 线程安全队列:toSend使用ConcurrentLinkedQueue,保证多线程添加数据时的安全性。
  4. 唤醒Selector:当添加新的待发送数据时,调用selector.wakeup()让Selector立即处理写事件,无需等待超时。

替代方案:多线程模式

如果你觉得NIO的Selector上手复杂,也可以用简单的多线程方案:

  • 一个线程专门负责读取服务端数据(阻塞读,但不影响发送线程)。
  • 另一个线程负责发送数据(和你现有逻辑类似)。
    这种方案代码更直观,但对于大量连接场景,Selector的性能更优。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 20:21:43