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

如何让Akka HTTP WebSocket客户端连接持久保持打开不自动关闭

WebSocket连接自动断开及后续问题解决方案

连接自动断开的根本原因

你最初的代码中outgoing使用Source.single()实现,该源在发送完唯一的一条消息后就会主动结束,而Akka HTTP的WebSocket流只要上游发送源或者下游接收源任意一端完成,整个连接流就会终止,因此连接会自动断开。

要保持连接永久打开直到主动关停应用,核心逻辑就是让发送源不主动完成,你后续修改中用到的Source.maybe就是为此设计的,它默认永远不会发出元素也不会主动结束,拼接在业务消息后可以保证连接不被主动关闭。

拼接Source.maybe后服务端无推送的排查点

你修改后的拼接逻辑本身没有问题,不会主动断开连接,遇到服务端无推送、程序挂起的问题可以从以下方向排查:

  • 检查服务端鉴权逻辑是否正常,确认你发送的auth请求格式、参数完全符合服务端要求,多数服务端在鉴权失败时不会主动断开连接,也不会返回任何响应,自然也不会推送订阅数据
  • 你当前的incoming Sink仅处理了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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 08:54:03