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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 00:06:11