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

基于Akka Http+WebSocket+BlockingQueue实现类GraphQL订阅遇阻求助

Akka Http WebSocket订阅功能问题排查与优化指导

方向判断

你的整体思路是可行的:用Akka Http承载WebSocket连接、通过消息中转机制处理数据库更新推送、独立处理阻塞队列事务数据,符合类GraphQL订阅的架构逻辑,但细节实现存在关键问题,导致Sink收不到消息、消息推送逻辑不符合订阅的“按需分发”要求。

Sink收不到消息的核心原因

  1. 消息类型匹配不全:你只处理了BinaryMessage.Strict,但客户端可能发送BinaryMessage.Streamed或TextMessage(比如调试阶段的文本格式查询),导致匹配不到对应case,代码无法执行。
  2. WebSocket连接可能未正确建立:CORS配置可能拦截了WebSocket握手请求,或者客户端连接逻辑有误,需先确认连接是否成功(可在路由中添加连接日志验证)。

代码修复与优化建议

1. 修复Sink的消息处理逻辑

覆盖所有Message类型,并将流式消息转为Strict消息统一处理:

val incomingMessages: Sink[Message, Future[Done]] = {
  Sink.foreachAsync(1) {
    case binaryMsg: BinaryMessage =>
      // 将流式二进制消息转为Strict,确保完整接收
      binaryMsg.toStrict(3.seconds).map { strictMsg =>
        println("Received binary query")
        val query = deserialize(strictMsg.asByteBuffer)
        // 后续订阅逻辑,需绑定当前WebSocket连接的推送通道
        subscribe(query, getCallbackForCurrentConnection)
      }
    case textMsg: TextMessage =>
      // 兼容文本格式的查询(调试或多场景适配)
      textMsg.toStrict(3.seconds).map { strictMsg =>
        println("Received text query")
        val query = strictMsg.text
        subscribe(query, getCallbackForCurrentConnection)
      }
    case other =>
      Future.failed(new Exception(s"Unsupported message type: ${other.getClass.getName}"))
  }
}

2. 重构消息推送逻辑(替换全局事件总线)

当前全局事件总线会导致所有WebSocket连接收到所有推送消息,不符合订阅的“按需推送”逻辑。应为每个WebSocket连接创建专属的订阅Actor,仅推送匹配该连接查询的结果:

定义订阅管理Actor

class SubscriptionActor(query: String, dbQueue: BlockingQueue[Data], wsSender: ActorRef) 
  extends Actor with ActorLogging {

  // 使用Akka调度器+阻塞线程池轮询队列,避免手动创建线程
  private val queuePoller = context.system.scheduler.scheduleAtFixedRate(
    initialDelay = 0.seconds,
    interval = 100.millis,
    receiver = self,
    message = PollQueue
  )(context.system.dispatchers.lookup("blocking-dispatcher"))

  override def receive: Receive = {
    case PollQueue =>
      // 用poll而非take,避免阻塞Actor线程
      Option(dbQueue.poll()).foreach { txData =>
        matchResult(query, txData).foreach { matchedResult =>
          val serialized = serialize(matchedResult)
          wsSender ! BinaryMessage(ByteString(serialized))
        }
      }
    case Terminated(_) =>
      // 连接关闭时停止调度器,释放资源
      queuePoller.cancel()
      context.stop(self)
  }
}

// 定义调度器消息
case object PollQueue

配置阻塞线程池(避免阻塞Akka默认线程池)

在application.conf中添加专门处理阻塞操作的调度器,防止阻塞Akka的非阻塞处理能力:

blocking-dispatcher {
  type = Dispatcher
  executor = "thread-pool-executor"
  thread-pool-executor {
    core-pool-size-min = 4
    core-pool-size-max = 4
  }
  throughput = 100
}

关联WebSocket流与订阅Actor

lazy val route: Route = cors() {
  path("ws") {
    handleWebSocketMessages {
      // 使用Sink和Source的Mat版本,关联Actor与WebSocket输出流
      Flow.fromSinkAndSourceMat(
        Sink.foreach[Message] {
          case binaryMsg: BinaryMessage =>
            binaryMsg.toStrict(3.seconds).map { strictMsg =>
              val query = deserialize(strictMsg.asByteBuffer)
              // 创建专属订阅Actor,传入当前WebSocket连接的输出ActorRef
              val subscriptionActor = context.system.actorOf(
                Props(new SubscriptionActor(query, javaBlockingQueue, sender()))
              )
              context.watch(subscriptionActor)
            }
          // 其他消息类型处理...
        },
        Source.actorRef[BinaryMessage](
          completionMatcher = PartialFunction.empty,
          failureMatcher = PartialFunction.empty,
          bufferSize = 1000,
          overflowStrategy = OverflowStrategy.backpressure
        )
      )((_, _) => Done)
    }
  }
}

3. 移除手动创建线程的逻辑

不要用Executors.newSingleThreadExecutor,改用Akka调度器+阻塞线程池处理队列轮询,避免线程泄漏和破坏Akka的线程模型。

关键注意事项

  • 每个订阅Actor仅负责单个连接的查询匹配,确保消息推送的准确性
  • 所有阻塞操作必须放到专门的线程池,避免影响Akka的非阻塞处理能力
  • 连接关闭时要及时停止调度器和Actor,防止资源泄漏

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 21:15:01