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
相关产品推荐
相关产品推荐

