如何用Source.asSubscriber包装响应式监听器?遇异常求解决方案
首先,你遇到的IllegalStateException是典型的Reactive Streams规范违反:当上游还没收到下游的需求信号(demand)时,就调用了onNext。Source.asSubscriber的设计并不是用来让你主动调用onNext推送消息的,这是你用错工具了。
为什么Source.asSubscriber不适合你的场景?
Source.asSubscriber的作用是把Akka Stream的Source转换成一个Reactive Streams Subscriber,用来订阅其他符合规范的Publisher。它的核心是被动接收Publisher的消息,而不是让你主动调用onNext来推送数据——这种主动推送的行为完全违背了Reactive Streams的背压机制,所以才会触发需求检查的异常。
正确的实现方式:用Source.queue适配背压
对于WebSocket这种主动推送消息的场景,你需要一个能协调外部推送和Akka Stream背压的中间层,Source.queue就是干这个的。它允许你从外部(比如WebSocketListener)发送元素,同时遵循Akka Stream的背压规则:当下游没有需求时,队列会缓存或者阻塞(根据你配置的溢出策略),不会直接触发违规的onNext。
改写后的代码示例
import akka.stream.scaladsl.Source import akka.stream.QueueOfferResult import org.asynchttpclient.{WebSocket, WebSocketListener, WebSocketUpgradeHandler} def createWsSource(url: String, onStartAction: Option[WebSocket => Unit]): Source[WsMessage, KillSwitch] = { // 创建支持背压的队列Source,bufferSize和溢出策略可根据业务调整 val queueSource = Source.queue[WsMessage](bufferSize = 16, OverflowStrategy.backpressure) queueSource.mapMaterializedValue { queue => val listener = new WebSocketListener { private var activeWs: Option[WebSocket] = None override def onOpen(ws: WebSocket): Unit = { activeWs = Some(ws) // 执行启动回调 onStartAction.foreach(_(ws)) // 将WsOpen消息发送到队列,自动处理背压 queue.offer(WsOpen(ws)) .recover { case _ => QueueOfferResult.Failure } // 处理发送失败的情况 } override def onClose(ws: WebSocket, code: Int, reason: String): Unit = { queue.complete() // 关闭队列,结束流 activeWs = None } override def onBinaryFrame(payload: Array[Byte], finalFragment: Boolean, rsv: Int): Unit = { // 处理二进制帧,构建你的WsMessage val binaryMsg = WsBinary(payload, finalFragment, rsv) queue.offer(binaryMsg) .recover { case _ => QueueOfferResult.Failure } } override def onTextFrame(payload: String, finalFragment: Boolean, rsv: Int): Unit = { // 处理文本帧,构建你的WsMessage val textMsg = WsText(payload, finalFragment, rsv) queue.offer(textMsg) .recover { case _ => QueueOfferResult.Failure } } override def onError(t: Throwable): Unit = { queue.fail(t) // 将错误传入流中 activeWs.foreach(_.sendCloseFrame()) } override def onPongFrame(payload: Array[Byte]): Unit = { super.onPongFrame(payload) } } // 初始化WebSocket连接 val websocket = asyncHttpClient .prepareGet(url) .execute(new WebSocketUpgradeHandler.Builder().addWebSocketListener(listener).build) .get() // 实现KillSwitch,用于手动终止流和WebSocket连接 new KillSwitch { override def shutdown(): Unit = { queue.complete() websocket.sendCloseFrame() } override def abort(ex: Throwable): Unit = { queue.fail(ex) websocket.sendCloseFrame() } } } } // 假设你定义了这些WsMessage类型 sealed trait WsMessage case class WsOpen(ws: WebSocket) extends WsMessage case class WsBinary(payload: Array[Byte], finalFragment: Boolean, rsv: Int) extends WsMessage case class WsText(payload: String, finalFragment: Boolean, rsv: Int) extends WsMessage
关键细节说明
- 背压处理:
queue.offer返回一个Future[QueueOfferResult],当下游有需求时,这个Future会成功返回QueueOfferResult.Enqueued;如果队列满了且用了backpressure策略,Future会一直等待直到有空间,完全符合Reactive Streams规范。 - 流生命周期管理:WebSocket的
onClose和onError会正确触发队列的complete()或fail(),保证流的生命周期和WebSocket连接一致。 - KillSwitch集成:手动终止流时,同时关闭队列和WebSocket连接,避免资源泄漏。
再重申Source.asSubscriber的正确用法
如果你需要对接一个现有的Reactive Streams Publisher,把它转换成Akka Source,应该用Source.fromPublisher;而Source.asSubscriber是反过来——把Akka Source变成Subscriber,用来订阅其他Publisher,比如:
import akka.stream.scaladsl.Source import org.reactivestreams.Publisher // 假设有一个第三方的Publisher val externalPublisher: Publisher[String] = ... // 转换成Akka Source val source: Source[String, NotUsed] = Source.fromPublisher(externalPublisher) // 或者用Source.asSubscriber订阅Publisher(很少用,除非你需要自定义Subscriber逻辑) val subscriber = Source.asSubscriber[String].to(Sink.foreach(println)).run() externalPublisher.subscribe(subscriber)
内容的提问来源于stack exchange,提问作者expert

