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

WebSocket连接断开后BroadcastHub实例未被清理问题求助

问题:BroadcastHub实例无法被GC回收导致内存泄漏

我有一个仅在内网可用的Web应用,通过WebSocket向前端异步推送数据。运行一段时间后,生成的多个BroadcastHub实例无法被全部清理,部分连接对应的BroadcastHub缓存被填满,推测由客户端断开连接导致。长此以往会耗尽Java VM内存。我已使用PoisonPill,日志显示其调用正常,但BroadcastHub及其缓存仍未被GC回收,恳请技术帮助。

原始处理WebSocket数据的Flow创建代码

val (actor, source) =
      // producer is materialized ActorRef
      Source.actorRef[String](10, akka.stream.OverflowStrategy.dropTail)
      .named(WebServer.socketWelcomeFlowName)
      // attach a BroadcastHub Sink to the producer
      .toMat(BroadcastHub.sink[String])(Keep.both)
      // By running/materializing the producer, we get back a Source, which
      // gives us access to the elements published by the producer
      .run()

connectedClients += 1
log.info(s"New client session with path: ${actor.path}.")
communication.welcomeNewClient(actor, log)(10 seconds, IO.contextShift(blockingIO)).unsafeRunAsync(_ => {})
communication.tellTopUserActor(log)(GenericMessageGalliumAdapterActor.name, GenericMessageGalliumAdapterActor.NewClientConnected).unsafeRunSync()

Flow[ws.Message]
  .watchTermination() { (m, f) =>
    f.onComplete(_ => {
      log.info(s"Client session '${actor.path}' left.")
      actor ! PoisonPill
      connectedClients -= 1
      }
    )
    m
  }
  .merge(source)
  .map {
    case TextMessage.Strict(tm) => handleMessage(actor, tm).map(Future.successful)
    case PassThroughMessage(tm) => Some(Future.successful(tm))
    case _ => Some(Future.successful(TextMessage.Strict("Gallium can only handle TextMessages for now!")))
  }
  .filter(_.isDefined)
  .map(_.get)
  .mapAsync(parallelism = 3)(identity)

更新后的代码

Source.actorRef[String]({
      case akka.actor.Status.Success(s: CompletionStrategy) => s
      case akka.actor.Status.Success(_)                     => CompletionStrategy.immediately
      case akka.actor.Status.Success                        => CompletionStrategy.immediately
      case akka.actor.Status.Failure                        => CompletionStrategy.immediately
    }, { case akka.actor.Status.Failure(cause)              => cause }
      , 10, OverflowStrategy.fail)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 12:40:25