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

Scala+fs2/cats+Http4s开发WebSocket:类实例发消息难题

在Http4s中让Worker实例主动推送WebSocket消息

我帮你搞定这个问题!在Http4s里实现Worker主动给WebSocket连接发消息,核心是用fs2.concurrent.Queue做消息中转层,而不是直接持有WebSocket连接引用——毕竟Http4s的WebSocket是基于fs2流设计的,直接抓连接会绕开框架的资源管理逻辑,反而容易出问题。下面是具体的实现方案和代码示例:

核心思路

  1. 为每个WebSocket连接创建一个线程安全的消息队列,Worker只需要往队列里丢消息,不用管连接细节
  2. 把队列的输出流绑定到WebSocket的发送端,框架会自动把队列里的消息推给前端
  3. 利用Http4s的资源管理机制,在连接关闭时自动清理Worker任务,避免内存泄漏

具体代码实现

第一步:定义消息类型

先定义我们要传递的批量任务更新结构,方便后续序列化:

import fs2.concurrent.Queue
import cats.effect.{IO, Fiber}
import org.http4s._
import org.http4s.dsl.io._
import org.http4s.websocket.WebSocketFrame
import org.http4s.blaze.server.BlazeServerBuilder
import scala.concurrent.duration._

// 定义任务更新消息,包含任务ID、状态和完成百分比
case class TaskUpdate(taskId: String, status: String, progress: Int)

第二步:实现Worker类

Worker类接收消息队列作为依赖,处理批量任务时直接往队列里推送更新:

class Worker(taskId: String, updateQueue: Queue[IO, TaskUpdate]) {
  // 模拟批量任务处理逻辑
  def startProcessing: IO[Unit] = {
    IO.untilDefinedM {
      // 用当前运行时间模拟任务进度
      IO.realTime.map(_.toSeconds).flatMap { elapsed =>
        val progress = math.min((elapsed * 10).toInt, 100)
        val status = if (progress < 100) "RUNNING" else "COMPLETED"
        
        // 往队列发送更新消息
        updateQueue.enqueue1(TaskUpdate(taskId, status, progress)) >>
        // 进度到100%时结束任务,否则继续等待
        if (progress == 100) IO.pure(Some(())) else IO.sleep(1.second).as(None)
      }
    }
  }
}

第三步:搭建WebSocket路由

在路由中创建队列、实例化Worker,并把队列流绑定到WebSocket发送端:

val wsRoutes = HttpRoutes.of[IO] {
  case GET -> Root / "ws" / taskId =>
    // 为当前连接创建一个无界消息队列
    Queue.unbounded[IO, TaskUpdate].flatMap { updateQueue =>
      // 实例化Worker,传入任务ID和消息队列
      val worker = new Worker(taskId, updateQueue)
      
      // 异步启动Worker的任务处理(不阻塞WebSocket连接建立)
      val workerFiber: IO[Fiber[IO, Throwable, Unit]] = worker.startProcessing.start
      
      // 构建WebSocket连接
      WebSocketBuilder2[IO].build(
        // 处理前端发来的消息(如果不需要可以忽略,这里只是打印日志)
        receive = _.evalMap { frame =>
          IO.println(s"Received message from frontend: $frame") >> IO.unit
        },
        // 发送流:把队列里的TaskUpdate转成JSON格式的Text帧
        send = updateQueue.dequeue.map { update =>
          WebSocketFrame.Text(
            s"""{"taskId":"${update.taskId}","status":"${update.status}","progress":${update.progress}}"""
          )
        }
      ).onFinalize {
        // 连接关闭时,自动取消Worker的任务,释放资源
        workerFiber.cancel >> IO.println(s"WebSocket connection for task $taskId closed")
      }
    }
}

// 启动Blaze服务器
def runServer: IO[Unit] = {
  BlazeServerBuilder[IO]
    .bindHttp(8080, "localhost")
    .withHttpApp(wsRoutes.orNotFound)
    .resource
    .use(_ => IO.never)
}

关键细节说明

  • 解耦Worker与WebSocket连接:Worker不需要知道WebSocket的存在,只负责生成任务更新消息,队列作为中间层实现了关注点分离
  • 资源自动管理:用onFinalize在连接关闭时取消Worker的Fiber,避免任务在后台无意义运行
  • 线程安全:fs2的Queue是线程安全的,支持多线程并发写入,完全适配批量任务的异步处理场景

如果你之前尝试直接持有WebSocket的发送端引用,那肯定会遇到问题——因为Http4s的WebSocket发送流是由框架控制的,外部直接操作会破坏流的生命周期管理。用队列中转才是符合Http4s函数式流设计的正确方式。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 08:36:25