Scala+fs2/cats+Http4s开发WebSocket:类实例发消息难题
在Http4s中让Worker实例主动推送WebSocket消息
我帮你搞定这个问题!在Http4s里实现Worker主动给WebSocket连接发消息,核心是用fs2.concurrent.Queue做消息中转层,而不是直接持有WebSocket连接引用——毕竟Http4s的WebSocket是基于fs2流设计的,直接抓连接会绕开框架的资源管理逻辑,反而容易出问题。下面是具体的实现方案和代码示例:
核心思路
- 为每个WebSocket连接创建一个线程安全的消息队列,Worker只需要往队列里丢消息,不用管连接细节
- 把队列的输出流绑定到WebSocket的发送端,框架会自动把队列里的消息推给前端
- 利用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
相关产品推荐
相关产品推荐

