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

Play Framework如何将POST接收的数据推送给WebSocket客户端?

嘿,我之前刚好碰到过一模一样的需求!Play官方文档确实在主动推送这块讲得比较少,主要聚焦在响应WebSocket接收的消息上,给你分享个我亲测可行的方案吧

解决方案:Play Framework实现POST触发WebSocket主动推送

核心思路

要实现“POST接收数据后推送给WebSocket客户端”,关键是要有一个全局的连接管理机制,把所有活跃的WebSocket连接保存起来,当POST接口收到数据时,遍历这些连接发送消息。下面给你两种不同场景下的实现方案:


方案一:手动管理WebSocket连接(适合需要精细控制的场景)

如果需要针对特定客户端推送、或者要做自定义的连接状态管理,这个方案更灵活。

1. 创建连接管理器

先写一个单例类,用来维护所有活跃的WebSocket连接,保证线程安全:

import akka.stream.scaladsl.SourceQueueWithComplete
import java.util.concurrent.{ConcurrentHashMap, AtomicInteger}
import scala.concurrent.ExecutionContext.global

object WebSocketConnectionManager {
  // 用ConcurrentHashMap存储连接,键是唯一ID,值是消息发送队列
  private val connections = new ConcurrentHashMap[String, SourceQueueWithComplete[String]]()
  private val idGenerator = new AtomicInteger(0)

  // 添加新连接,返回唯一ID
  def addConnection(queue: SourceQueueWithComplete[String]): String = {
    val connectionId = idGenerator.incrementAndGet().toString
    connections.put(connectionId, queue)
    connectionId
  }

  // 移除指定连接
  def removeConnection(queue: SourceQueueWithComplete[String]): Unit = {
    connections.values().removeIf(_ == queue)
  }

  // 向所有连接推送消息
  def pushToAllClients(message: String): Unit = {
    val iterator = connections.values().iterator()
    while (iterator.hasNext) {
      val queue = iterator.next()
      // 尝试发送消息,发送失败则移除无效连接
      queue.offer(message).onComplete {
        case scala.util.Failure(_) => iterator.remove()
        case _ => // 发送成功,无需操作
      }(global)
    }
  }
}

2. 实现WebSocket路由

在控制器里写WebSocket的处理逻辑,当客户端建立连接时加入管理器,断开时移除:

import akka.stream.scaladsl.{Flow, Sink, Source}
import play.api.mvc.{WebSocket, Action, ControllerComponents, JsValue}
import play.api.libs.json.Json

class MyController(cc: ControllerComponents) extends AbstractController(cc) {
  def ws: WebSocket = WebSocket.accept[String, String] { request =>
    // 创建一个消息队列,用来向客户端发送消息
    val (messageQueue, clientSource) = Source.queue[String](10, akka.stream.OverflowStrategy.dropNew).preMaterialize()
    
    // 将连接加入管理器
    WebSocketConnectionManager.addConnection(messageQueue)

    // 处理客户端发送的消息(如果不需要处理客户端消息,用Sink.ignore即可)
    val clientSink = Sink.foreach[String] { clientMsg =>
      // 可选:处理客户端发来的消息,比如心跳检测
      println(s"Received from client: $clientMsg")
    }

    // 监听连接关闭事件,移除无效连接
    val closeHandler = Sink.onComplete[akka.Done] { _ =>
      WebSocketConnectionManager.removeConnection(messageQueue)
    }

    // 合并Sink,同时处理客户端消息和连接关闭事件
    val combinedSink = Sink.combine(clientSink, closeHandler)(akka.stream.scaladsl.Broadcast[String])

    // 返回Flow,连接建立后就可以通过messageQueue给客户端发消息了
    Flow.fromSinkAndSource(combinedSink, clientSource)
  }
}

3. 实现POST接口触发推送

在同一个控制器里写POST接口,收到客户端1的数据后调用推送方法:

def receiveClientData: Action[JsValue] = Action(parse.json) { request =>
  // 解析客户端1发送的数据(根据实际需求调整)
  val data = Json.stringify(request.body)
  
  // 可选:对数据进行存储、校验等处理
  // saveDataToDatabase(data)

  // 推送给所有活跃的WebSocket客户端
  WebSocketConnectionManager.pushToAllClients(data)

  Ok("Data received and pushed to WebSocket clients")
}

方案二:用Akka Stream的广播流(更简洁,适合广播场景)

如果你的需求是把消息广播给所有WebSocket客户端,不需要针对单个客户端推送,那可以用Akka Stream的MergeHub和BroadcastHub来实现,不需要手动管理连接:

1. 创建广播流单例

import akka.stream.scaladsl.{BroadcastHub, MergeHub, Sink, Source}

object WebSocketBroadcaster {
  // 创建广播流:MergeHub接收消息,BroadcastHub分发给所有客户端
  private val (broadcastSink, clientSource) = MergeHub.source[String](perProducerBufferSize = 16)
    .toMat(BroadcastHub.sink(bufferSize = 256))(Keep.both)
    .run()

  // 对外暴露的Sink,用来发送广播消息
  def sendBroadcastMessage: Sink[String, akka.NotUsed] = broadcastSink

  // 对外暴露的Source,每个WebSocket客户端订阅这个Source接收消息
  def clientMessageSource: Source[String, akka.NotUsed] = clientSource
}

2. 实现WebSocket路由

def ws: WebSocket = WebSocket.accept[String, String] { _ =>
  // 客户端只需要订阅广播Source,Sink用ignore忽略客户端消息(如果需要处理可以调整)
  Flow.fromSinkAndSource(Sink.ignore, WebSocketBroadcaster.clientMessageSource)
}

3. POST接口触发广播

def receiveClientData: Action[JsValue] = Action(parse.json) { request =>
  val data = Json.stringify(request.body)
  
  // 将数据发送到广播Sink,自动推送给所有客户端
  Source.single(data).runWith(WebSocketBroadcaster.sendBroadcastMessage)

  Ok("Data broadcast to WebSocket clients")
}

注意事项

  • 连接清理:不管用哪种方案,都要确保无效连接被及时移除,避免内存泄漏。方案一通过监听连接关闭事件和发送失败的情况清理,方案二Akka Stream会自动处理。
  • 消息格式:示例里用了字符串,实际项目中建议用JSON格式,方便客户端解析,Play的Json工具类可以轻松处理。
  • 线程安全:方案一用了ConcurrentHashMap保证多线程安全;方案二基于Akka Stream,本身就是线程安全的。
  • 错误处理:可以根据需求添加异常捕获,比如推送消息失败时记录日志。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 09:33:57