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
相关产品推荐
相关产品推荐

