如何让Akka HTTP WebSocket客户端连接持久保持打开不自动关闭
WebSocket连接自动断开及后续问题解决方案
连接自动断开的根本原因
你最初的代码中outgoing使用Source.single()实现,该源在发送完唯一的一条消息后就会主动结束,而Akka HTTP的WebSocket流只要上游发送源或者下游接收源任意一端完成,整个连接流就会终止,因此连接会自动断开。
要保持连接永久打开直到主动关停应用,核心逻辑就是让发送源不主动完成,你后续修改中用到的Source.maybe就是为此设计的,它默认永远不会发出元素也不会主动结束,拼接在业务消息后可以保证连接不被主动关闭。
拼接Source.maybe后服务端无推送的排查点
你修改后的拼接逻辑本身没有问题,不会主动断开连接,遇到服务端无推送、程序挂起的问题可以从以下方向排查:
- 检查服务端鉴权逻辑是否正常,确认你发送的auth请求格式、参数完全符合服务端要求,多数服务端在鉴权失败时不会主动断开连接,也不会返回任何响应,自然也不会推送订阅数据
- 你当前的
incomingSink仅处理了TextMessage.Strict类型的消息,如果服务端返回的是流式文本消息(TextMessage.Streamed)、二进制消息、Ping/Pong帧都会被直接忽略,建议先加全类型日志确认是否有响应被漏掉 - 确认服务端是否要求客户端定时发送心跳帧维持连接,很多WebSocket服务如果超过一定时间没有收到客户端消息,会主动断开连接或者停止推送数据
完整长连接实现示例
package docs.http.scaladsl import akka.actor.ActorSystem import akka.Done import akka.http.scaladsl.Http import akka.stream.scaladsl._ import akka.http.scaladsl.model._ import akka.http.scaladsl.model.ws._ import scala.concurrent.Future import scala.concurrent.duration._ object WebSocketLongLiveClient { def main(args: Array[String]): Unit = { implicit val system = ActorSystem() import system.dispatcher // 全类型入站消息处理,避免漏掉响应 val incoming: Sink[Message, Future[Done]] = Sink.foreach[Message] { case TextMessage.Strict(text) => println(s"收到严格文本消息: $text") case TextMessage.Streamed(stream) => stream.runFold("")(_ + _).onComplete { text => println(s"收到流式文本消息: ${text.get}") } case BinaryMessage.Strict(data) => println(s"收到严格二进制消息,长度: ${data.length}") case BinaryMessage.Streamed(stream) => stream.runFold(0)(_ + _.length).onComplete { len => println(s"收到流式二进制消息,总长度: ${len.get}") } case PingMessage(_) => println(s"收到Ping,连接正常") case PongMessage(_) => println(s"收到Pong,心跳响应正常") } // 业务消息+定时心跳+永不结束源组合,保证连接长开 val authMsg = TextMessage("""{"action":"auth","params":"APIKEY_123"}""") val subMsg = TextMessage("""{"action":"subscribe","params":"topic123"}""") // 每30秒发送一次心跳,适配服务端存活检测规则 val heartbeatSource = Source.tick(30.seconds, 30.seconds, TextMessage("""{"action":"heartbeat"}""")) val outgoing = Source.combine( Source.single(authMsg), Source.single(subMsg), heartbeatSource, Source.maybe[Message] )(Concat(_)) val webSocketFlow = Http().webSocketClientFlow(WebSocketRequest("ws://你的服务端地址")) val (upgradeResponse, closed) = outgoing .viaMat(webSocketFlow)(Keep.right) .toMat(incoming)(Keep.both) .run() val connected = upgradeResponse.flatMap { upgrade => if (upgrade.response.status == StatusCodes.SwitchingProtocols) { Future.successful(Done) } else { throw new RuntimeException(s"连接失败: ${upgrade.response.status}") } } connected.onComplete(println) closed.foreach(_ => println("连接已关闭")) } }
内容的提问来源于stack exchange,提问作者Blankman
相关产品推荐
相关产品推荐

