Play Framework(Scala)实现POST请求数据转发至WebSocket输出
Play Framework(Scala):将POST请求数据转发到WebSocket客户端
问题描述
需要实现以下流程:
- 接收外部API发来的POST请求消息
- 将消息传递至WebSocket组件并推送给Angular客户端
- 客户端实时接收WebSocket消息
当前仅实现了基础WebSocket功能(回复客户端"I received your message"),需要在getMessage POST方法中完成消息转发逻辑。
现有代码
package controllers import org.apache.pekko.actor._ import org.apache.pekko.stream.Materializer import play.api.libs.streams.ActorFlow import javax.inject._ import play.api.mvc._ @Singleton class HomeController @Inject()(val controllerComponents: ControllerComponents) (implicit system: ActorSystem, mat: Materializer) extends BaseController { def index(): Action[AnyContent] = Action { implicit request: Request[AnyContent] => Ok(views.html.index()) } def getMessage: Action[AnyContent] = Action { request: Request[AnyContent] => println(request.body.toString) //connect the websocket here Ok("It works!") } def socket: WebSocket = WebSocket.accept[String, String] { request => ActorFlow.actorRef { out => MyWebSocketActor.props(out) } } } object MyWebSocketActor { def props(out: ActorRef): Props = Props(new MyWebSocketActor(out)) } class MyWebSocketActor(out: ActorRef) extends Actor { def receive: Receive = { case msg: String => out ! ("I received your message: " + msg) } }
解决方案
核心是通过全局Akka广播Actor实现POST接口与WebSocket客户端的消息桥接,具体步骤如下:
1. 创建全局广播Actor
这个Actor负责维护所有在线WebSocket客户端的引用,并处理消息广播逻辑:
package controllers import org.apache.pekko.actor._ // 全局广播Actor的消息协议 object WebSocketBroadcastActor { case class RegisterClient(client: ActorRef) case class UnregisterClient(client: ActorRef) case class BroadcastMessage(message: String) def props: Props = Props(new WebSocketBroadcastActor()) } class WebSocketBroadcastActor extends Actor { // 存储所有在线客户端的Actor引用 private var clients = Set.empty[ActorRef] override def receive: Receive = { // 注册新客户端 case WebSocketBroadcastActor.RegisterClient(client) => clients += client context.watch(client) // 监听客户端Actor终止事件,自动注销 // 注销客户端 case WebSocketBroadcastActor.UnregisterClient(client) => clients -= client // 客户端断开连接时自动注销 case Terminated(client) => clients -= client // 向所有客户端广播消息 case WebSocketBroadcastActor.BroadcastMessage(msg) => clients.foreach(_ ! msg) } }
2. 修改HomeController
注入并初始化广播Actor,在POST接口中发送消息到广播Actor:
@Singleton class HomeController @Inject()(val controllerComponents: ControllerComponents) (implicit system: ActorSystem, mat: Materializer) extends BaseController { // 创建全局唯一的广播Actor实例 private val broadcastActor = system.actorOf(WebSocketBroadcastActor.props, "websocket-broadcast") def index(): Action[AnyContent] = Action { implicit request: Request[AnyContent] => Ok(views.html.index()) } def getMessage: Action[AnyContent] = Action { request: Request[AnyContent] => // 提取POST请求中的文本消息(根据实际请求格式调整,如JSON/表单) val message = request.body.asText.getOrElse("No content received") // 将消息发送给广播Actor,由它转发给所有WebSocket客户端 broadcastActor ! WebSocketBroadcastActor.BroadcastMessage(message) Ok("Message forwarded to clients!") } def socket: WebSocket = WebSocket.accept[String, String] { request => // 将广播Actor引用传入WebSocket Actor ActorFlow.actorRef { out => MyWebSocketActor.props(out, broadcastActor) } } }
3. 修改WebSocket Actor
让WebSocket Actor在启动时注册到广播Actor,终止时注销,并处理广播消息:
object MyWebSocketActor { // 更新props方法,传入广播Actor引用 def props(out: ActorRef, broadcastActor: ActorRef): Props = Props(new MyWebSocketActor(out, broadcastActor)) } class MyWebSocketActor(out: ActorRef, broadcastActor: ActorRef) extends Actor { // 启动时向广播Actor注册自己 override def preStart(): Unit = { broadcastActor ! WebSocketBroadcastActor.RegisterClient(self) } // 终止时向广播Actor注销自己 override def postStop(): Unit = { broadcastActor ! WebSocketBroadcastActor.UnregisterClient(self) } override def receive: Receive = { // 保留原有处理客户端消息的逻辑 case clientMsg: String => out ! ("I received your message: " + clientMsg) // 处理来自广播Actor的消息,转发给WebSocket客户端 case broadcastMsg: String => out ! broadcastMsg } }
关键说明
- 客户端生命周期管理:通过
context.watch监听客户端Actor的终止事件,自动移除无效连接,避免发送消息失败 - 消息格式适配:如果POST请求是JSON格式,需修改
getMessage中的消息提取逻辑,例如使用request.body.asJson解析后转换为字符串 - 单例广播Actor:确保广播Actor全局唯一,所有WebSocket连接和POST请求都与同一个实例交互
内容的提问来源于stack exchange,提问作者DanCrts
相关产品推荐
相关产品推荐

