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

Ktor WebSocket未及时关闭问题及内存泄漏解决方案咨询

问题:Ktor WebSocket客户端断网后无法及时关闭连接引发内存泄漏

问题描述

移动端断网约3分钟后,Socket因ping/pong失效,但代码中的finally块迟迟不触发(甚至30分钟后才执行),导致WebSocket连接一直未关闭,进而引发内存泄漏。需要实现客户端断网3分钟后,服务器端自动移除对应WebSocket连接,避免资源浪费。

当前配置与代码

WebSocket安装配置

fun Application.configureWebSocket(){
install(WebSockets) {
    pingPeriod = Duration.ofSeconds(15)
    timeout = Duration.ofSeconds(15)
    maxFrameSize = kotlin.Long.MAX_VALUE
    masking = false
}

WebSocket路由处理代码

routing {
    webSocket("ws") {
        val token = call.request.queryParameters["token"]
        if (token == null) {
            close(CloseReason(CloseReason.Codes.VIOLATED_POLICY, "No Token"))
            return@webSocket
        }

        val decodedJWT = try { JwtFactory.buildverifier().verify(token) }
        catch (e: Exception) {
            close(CloseReason(CloseReason.Codes.VIOLATED_POLICY, "Invalid Token: ${e.message}"))
            return@webSocket
        }

        val userId: UUID = try { UUID.fromString(decodedJWT.getClaim(JwtClaimConstant.claimUserId).asString())  }
        catch (e: Exception) {
            close(CloseReason(CloseReason.Codes.VIOLATED_POLICY, "Invalid Token: ${e.message}"))
            return@webSocket
        }

        val sessionId = decodedJWT.id?.let {
            runCatching { UUID.fromString(it) }.getOrNull()
        } ?: run {
            close(CloseReason(CloseReason.Codes.VIOLATED_POLICY, "Invalid or missing sessionId (jti)"))
            return@webSocket
        }
        logger.info("$userId is connected")

        try {
            println("$userId start")
            incoming.consumeEach {
                when (it) {
                    is Frame.Text -> {
                        val text = it.readText()
                        println("tototot $userId Received: $text")
                    }
                    is Frame.Close -> {
                        println("tototot $userId WebSocket closed by server with reason: ${it.readReason()}")
                    }
                    is Frame.Ping -> {
                        println("tototot $userId ping: $it")
                    }
                    is Frame.Pong -> {
                        println("tototot $userId pong: $it")
                    }  
                    else -> {
                        println("tototot $userId else: $it")
                    }
                }
            }
        } catch (e: Exception) {
            println("$userId error $e")
        } finally {
            println("$userId finally remove")
        }

        println("$userId end")
    }
}

补充信息

iOS端断网10秒后执行ping会触发如下超时错误并关闭连接,但服务器端无任何响应:

Optional(Error Domain=NSPOSIXErrorDomain Code=60 "Operation timed out" UserInfo={NSDescription=Operation timed out})


解决方案

1. 优化超时配置与TCP层检查

当前timeout设置为15秒(即发送Ping后15秒未收到Pong触发超时),但实际可能受系统TCP超时干扰(默认TCP超时可能远大于3分钟)。可以:

  • 确认Ktor的timeout参数生效:该参数是上层WebSocket协议的Pong等待超时,若底层TCP连接未断开,Ktor依赖此参数检测连接失效。
  • 可尝试调整系统TCP KeepAlive配置(如缩短TCP空闲超时),但需注意全局影响。

2. 主动监控Pong响应,手动触发关闭

在连接处理逻辑中添加独立协程,监控Pong响应次数,达到阈值后手动关闭连接:

routing {
    webSocket("ws") {
        // ... 之前的鉴权逻辑省略 ...
        logger.info("$userId is connected")

        var missedPongs = 0
        // 启动监控协程
        val monitorJob = launch {
            while (true) {
                delay(pingPeriod.toMillis())
                if (!isActive) break
                // 每发送一次Ping后,未收到Pong则计数+1
                missedPongs++
                if (missedPongs >= 12) { // 15秒一次Ping,12次即3分钟
                    close(CloseReason(CloseReason.Codes.ABNORMAL_CLOSURE, "No pong for 3 minutes"))
                    break
                }
            }
        }

        try {
            println("$userId start")
            incoming.consumeEach { frame ->
                when (frame) {
                    // ... 其他帧处理逻辑省略 ...
                    is Frame.Pong -> {
                        missedPongs = 0 // 收到Pong重置计数
                        println("tototot $userId pong: $it")
                    }
                }
            }
        } catch (e: IOException) {
            println("$userId connection lost: ${e.message}")
        } catch (e: Exception) {
            println("$userId error: ${e.message}")
        } finally {
            monitorJob.cancel()
            println("$userId finally remove")
            // 若有全局连接集合,在此移除当前连接
            // activeConnections.remove(userId)
        }

        println("$userId end")
    }
}

3. 精准捕获网络异常

将泛型Exception捕获改为优先捕获网络类异常(如IOException),确保连接断开时能及时触发finally块:

catch (e: IOException) {
    println("$userId connection lost due to network issue: ${e.message}")
} catch (e: Exception) {
    println("$userId unexpected error: ${e.message}")
}

4. 规范全局连接管理

若使用全局集合存储活跃WebSocket连接,必须在finally块中移除当前连接,避免内存泄漏:

// 全局连接集合示例
val activeConnections = ConcurrentHashMap<UUID, WebSocketSession>()

// 鉴权成功后添加连接
activeConnections[userId] = this

// finally块中移除
finally {
    activeConnections.remove(userId)
    monitorJob.cancel()
    println("$userId finally remove")
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 21:34:56