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

Kotlin中多线程间共享ServerSocket连接及线程通信方案问询

Kotlin服务器线程通信管理方案

核心设计思路

采用线程安全注册表+阻塞队列通信的模式,实现管理端、用户线程、UDP线程的解耦通信,同时保证所有IO操作互不阻塞,不影响服务器整体待命能力。

关键组件

  • 线程注册表:用ConcurrentHashMap存储线程标识与对应通信通道、关联资源(如用户Socket、UDP Socket),确保多线程下的安全访问。
  • 命令消息模型:定义标准化指令结构,统一管理端发送的命令类型(如推送JSON、创建UDP转发)。
  • 阻塞队列:每个用户线程、UDP线程持有独立的BlockingQueue,用于接收外部命令,避免线程在IO等待时无法响应指令。

具体实现步骤

1. 定义命令消息模型

// 命令类型枚举
enum class CommandType {
    PUSH_JSON_TO_USER, // 向用户终端推送JSON
    CREATE_UDP_FORWARD // 创建UDP转发线程
}

// 通用命令消息
data class ThreadCommand(
    val targetThreadId: String, // 目标线程标识
    val commandType: CommandType,
    val payload: Any? = null // 命令携带的数据,如JSON字符串、远程UDP地址
)

2. 实现全局线程注册表

object ThreadRegistry {
    // key: 线程唯一标识,value: 该线程的命令队列
    private val threadCommandQueues = ConcurrentHashMap<String, BlockingQueue<ThreadCommand>>()
    // 存储用户线程关联的Socket,用于后续数据推送
    private val userSockets = ConcurrentHashMap<String, Socket>()

    // 注册用户线程
    fun registerUserThread(threadId: String, queue: BlockingQueue<ThreadCommand>, socket: Socket) {
        threadCommandQueues[threadId] = queue
        userSockets[threadId] = socket
    }

    // 注册UDP线程
    fun registerUdpThread(threadId: String, queue: BlockingQueue<ThreadCommand>) {
        threadCommandQueues[threadId] = queue
    }

    // 发送命令到目标线程
    fun sendCommand(command: ThreadCommand): Boolean {
        return threadCommandQueues[command.targetThreadId]?.offer(command) ?: false
    }

    // 获取用户Socket
    fun getUserSocket(threadId: String): Socket? {
        return userSockets[threadId]
    }

    // 移除线程注册(连接关闭时调用)
    fun unregisterThread(threadId: String) {
        threadCommandQueues.remove(threadId)
        userSockets.remove(threadId)
    }
}

3. 改造用户线程逻辑

每个用户线程启动时生成唯一ID,注册到注册表,同时循环处理客户端TCP数据和队列中的命令:

val socketListener_User = ServerSocket(10001)
socketListener_User.use {
    while (true) {
        val socket_User = socketListener_User.accept()
        // 生成唯一线程标识(结合Socket地址+端口或UUID)
        val threadId = "USER-${socket_User.inetAddress.hostAddress}-${socket_User.port}"
        val commandQueue = LinkedBlockingQueue<ThreadCommand>()
        
        ThreadRegistry.registerUserThread(threadId, commandQueue, socket_User)
        
        thread(start = true, name = threadId) {
            try {
                val inputStream = socket_User.getInputStream()
                val outputStream = socket_User.getOutputStream()
                val reader = BufferedReader(InputStreamReader(inputStream))
                val writer = BufferedWriter(OutputStreamWriter(outputStream))

                // 单独线程处理客户端输入,避免阻塞命令逻辑
                val ioExecutor = Executors.newSingleThreadExecutor()
                ioExecutor.submit {
                    var line: String?
                    while (socket_User.isConnected && reader.readLine().also { line = it } != null) {
                        // 处理用户TCP数据,比如触发UDP转发逻辑
                        // ...
                    }
                }

                // 循环处理管理端命令
                while (socket_User.isConnected) {
                    val command = commandQueue.take() // 阻塞等待命令,不占用CPU
                    when (command.commandType) {
                        CommandType.PUSH_JSON_TO_USER -> {
                            val jsonData = command.payload as String
                            writer.write(jsonData)
                            writer.newLine()
                            writer.flush()
                        }
                        CommandType.CREATE_UDP_FORWARD -> {
                            val remoteUdpAddress = command.payload as Pair<String, Int>
                            // 创建UDP转发线程
                            startUdpForwardThread(threadId, remoteUdpAddress, socket_User)
                        }
                    }
                }
            } catch (e: IOException) {
                e.printStackTrace()
            } finally {
                ThreadRegistry.unregisterThread(threadId)
                socket_User.close()
            }
        }
    }
}

4. 实现管理端线程逻辑

监听本地10002端口,仅允许本地连接,解析命令后发送到目标线程:

fun startManagementServer() {
    val managementSocket = ServerSocket(10002, 5, InetAddress.getByName("127.0.0.1"))
    managementSocket.use {
        while (true) {
            val socket = it.accept()
            thread(start = true) {
                val reader = BufferedReader(InputStreamReader(socket.inputStream))
                var commandLine: String?
                while (socket.isConnected && reader.readLine().also { commandLine = it } != null) {
                    // 解析命令格式:COMMAND_TYPE TARGET_THREAD_ID PAYLOAD
                    val parts = commandLine!!.split(" ", limit = 3)
                    if (parts.size < 2) continue
                    
                    val commandType = CommandType.valueOf(parts[0])
                    val targetThreadId = parts[1]
                    val payload = parts.getOrNull(2)
                    
                    val command = when(commandType) {
                        CommandType.PUSH_JSON_TO_USER -> ThreadCommand(targetThreadId, commandType, payload)
                        CommandType.CREATE_UDP_FORWARD -> {
                            val (host, port) = payload!!.split(":")
                            ThreadCommand(targetThreadId, commandType, host to port.toInt())
                        }
                    }
                    
                    // 发送命令到目标线程,返回执行结果
                    val success = ThreadRegistry.sendCommand(command)
                    socket.getOutputStream().write(
                        if (success) "SUCCESS\n" else "FAIL: Thread not found\n".toByteArray()
                    )
                }
                socket.close()
            }
        }
    }
}

5. UDP转发线程实现

由用户线程或管理端触发创建,实现TCP与UDP的双向数据转换:

fun startUdpForwardThread(userThreadId: String, remoteUdpAddr: Pair<String, Int>, userSocket: Socket) {
    val udpThreadId = "UDP-${userThreadId}"
    val commandQueue = LinkedBlockingQueue<ThreadCommand>()
    ThreadRegistry.registerUdpThread(udpThreadId, commandQueue)
    
    thread(start = true, name = udpThreadId) {
        DatagramSocket().use { udpSocket ->
            val remoteAddress = InetAddress.getByName(remoteUdpAddr.first)
            val remotePort = remoteUdpAddr.second
            
            // 转发用户TCP数据到UDP服务器
            thread(start = true) {
                val inputStream = userSocket.getInputStream()
                val buffer = ByteArray(1024)
                var len: Int
                while (userSocket.isConnected && inputStream.read(buffer).also { len = it } != -1) {
                    val packet = DatagramPacket(buffer, len, remoteAddress, remotePort)
                    udpSocket.send(packet)
                }
            }
            
            // 转发UDP服务器数据到用户TCP连接
            val buffer = ByteArray(1024)
            while (userSocket.isConnected) {
                val packet = DatagramPacket(buffer, buffer.size)
                udpSocket.receive(packet)
                userSocket.getOutputStream().write(packet.data, 0, packet.length)
                userSocket.getOutputStream().flush()
            }
        }
        ThreadRegistry.unregisterThread(udpThreadId)
    }
}

关键注意事项

  • 线程安全:所有注册表操作依赖ConcurrentHashMap,避免并发访问冲突。
  • 非阻塞IO:用户线程的TCP数据处理单独用线程池执行,不阻塞命令响应逻辑。
  • 资源清理:线程关闭时必须从注册表移除,防止内存泄漏。
  • 本地访问限制:管理端ServerSocket绑定127.0.0.1,仅允许本地应用连接,提升安全性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 06:50:27