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

如何用Source.asSubscriber包装响应式监听器?遇异常求解决方案

解决AsyncHttpClient WebSocket转Akka Stream Source的背压异常问题

首先,你遇到的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

关键细节说明

  1. 背压处理:queue.offer返回一个Future[QueueOfferResult],当下游有需求时,这个Future会成功返回QueueOfferResult.Enqueued;如果队列满了且用了backpressure策略,Future会一直等待直到有空间,完全符合Reactive Streams规范。
  2. 流生命周期管理:WebSocket的onClose和onError会正确触发队列的complete()或fail(),保证流的生命周期和WebSocket连接一致。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 07:04:11