使用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会自动处理这个消息:
- 清理内部的订阅者状态,将当前订阅者置为
None - 后续新的
ActorPublisher实例订阅时,就能正常建立新的订阅关系
如果需要在Actor中添加自定义的清理逻辑,可以重写onCancel方法:
class MyActor(...) extends Actor with ActorPublisher[A] { override def onCancel(): Unit = { // 在这里添加自定义清理逻辑,比如关闭资源、重置状态等 super.onCancel() // 必须调用父类方法,确保内部订阅状态被正确清理 } // 你的业务逻辑... }
内容的提问来源于stack exchange,提问作者Miguel
相关产品推荐
相关产品推荐

