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

使用ActorPublisher源与重连WebSocketClientFlow时出现多订阅者错误

解决ActorPublisher重连WebSocket时的多订阅者异常问题

这个问题的核心在于:ActorPublisher遵循Reactive Streams规范,只允许存在一个订阅者。当WebSocket连接断开后,旧流的订阅关系并没有被主动取消,导致新流尝试订阅同一个Actor时触发了IllegalStateException。我们不需要停止并重建Actor,只需要在旧流终止时主动取消它对Actor的订阅即可。

方案1:跟踪当前订阅者,断开时主动取消

我们可以维护当前活跃的ActorPublisher实例,在WebSocket连接断开时调用其cancel()方法,或者直接向Actor发送Cancel消息,让Actor清理旧的订阅状态。

调整后的代码示例:

import akka.http.scaladsl.Http
import akka.http.scaladsl.model._
import akka.http.scaladsl.model.ws._
import akka.stream.ActorMaterializer
import akka.stream.actor.ActorPublisher
import akka.stream.scaladsl._
import akka.stream.actor.ActorPublisherMessage.Cancel

val actor = ...
// 用Option跟踪当前活跃的Publisher实例
var currentPublisher: Option[ActorPublisher[A]] = None

def connect(): Unit = {
  val publisher = ActorPublisher[A](actor)
  currentPublisher = Some(publisher)
  val source = Source.fromPublisher(publisher)
  openWebSocket(source)
}

def openWebSocket(source: Source[A, NotUsed]): Unit = {
  val flow = Http().webSocketClientFlow(WebSocketRequest(URL))
  val (response, closed) = source
    .map { instanceOfA => 
      TextMessage(instanceOfA.asJson)
    }
    .viaMat(flow)(Keep.right)
    .toMat(Sink.foreach { 
      case message: TextMessage.Strict => println(message.text) 
    })(Keep.both)
    .run()

  // 连接断开时,先取消旧订阅再重连
  closed.foreach { _ =>
    currentPublisher.foreach { pub =>
      // 方式1:调用Publisher的cancel方法(内部会向Actor发送Cancel消息)
      pub.cancel()
      // 方式2:直接向Actor发送Cancel消息,效果和上面一致
      // actor ! Cancel
    }
    currentPublisher = None
    connect()
  }
}

方案2:利用流的生命周期钩子(更优雅)

不需要维护可变状态,我们可以通过watchTermination操作监听流的终止事件,自动触发订阅取消逻辑,更符合Akka Stream的声明式风格:

def connect(): Unit = {
  val publisher = ActorPublisher[A](actor)
  val source = Source.fromPublisher(publisher)
    // 监听流的终止事件,无论正常结束还是异常终止都取消订阅
    .watchTermination() { (_, terminationFuture) =>
      terminationFuture.onComplete { _ =>
        publisher.cancel()
      }(ActorMaterializer.mat.system.dispatcher)
      NotUsed
    }
  openWebSocket(source)
}

def openWebSocket(source: Source[A, NotUsed]): Unit = {
  val flow = Http().webSocketClientFlow(WebSocketRequest(URL))
  val (response, closed) = source
    .map { instanceOfA => 
      TextMessage(instanceOfA.asJson)
    }
    .viaMat(flow)(Keep.right)
    .toMat(Sink.foreach { 
      case message: TextMessage.Strict => println(message.text) 
    })(Keep.both)
    .run()

  closed.foreach { _ => connect() }
}

关键原理说明

当调用publisher.cancel()或向Actor发送Cancel消息时,继承自ActorPublisher的Actor会自动处理这个消息:

  1. 清理内部的订阅者状态,将当前订阅者置为None
  2. 后续新的ActorPublisher实例订阅时,就能正常建立新的订阅关系

如果需要在Actor中添加自定义的清理逻辑,可以重写onCancel方法:

class MyActor(...) extends Actor with ActorPublisher[A] {
  override def onCancel(): Unit = {
    // 在这里添加自定义清理逻辑,比如关闭资源、重置状态等
    super.onCancel() // 必须调用父类方法,确保内部订阅状态被正确清理
  }

  // 你的业务逻辑...
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:39:06