【问题标题】:How to create a Source in Akka with Actor to serve WebSocket connection如何在 Akka 中使用 Actor 创建一个 Source 来服务 WebSocket 连接
【发布时间】:2020-05-14 10:21:10
【问题描述】:

我尝试使用 Akka Actors、Streams 和 WebSockets 创建一个简单的聊天。我想为serve WebSocket 连接创建单独的 Sink 和 Source。

我为roomId创建了一个聊天室:

path("ws" / "room" / IntNumber) { roomId => {
  println(s"Connecting to room $roomId")
  parameter("userName") { userName =>
    extractUpgradeToWebSocket { upgrade =>
      val chatRoom = ChatRooms.findOrCreate(roomId)
      val (sink, source) = chatRoom.getSinkAndSource(userName)
      complete(upgrade.handleMessagesWithSinkSource(sink, source))
    }
  }
}

聊天创建并传递一个生成输出的Source[Message, _] 和一个接收输入的Sink[Message, _]handleMessagesWithSinkSource 方法。

我在创建一个工作 source 时遇到问题,Messages 由我的 Actor 填充(Source.actorRefWithBackpressure 应该是 allow 它)。 Sinksource2 按预期工作,但 source 没有:

 def getSinkAndSource(name: String) = {
   val source = Source.actorRefWithBackpressure[Message](AckMessage, {
     case _: Success => CompletionStrategy.draining
   }, PartialFunction.empty)

   val wsActorRef = source.to(Sink.ignore).run()

   val receiver = actorSystem.actorOf(Props(classOf[ChatParticipantActor], name, ChatRoomActor, wsActorRef))
   val sink = Sink.actorRefWithBackpressure(receiver, InitMessage, AckMessage, OnCompleteMessage, onErrorMessage)

   val source2 = Source.tick(FiniteDuration(1, TimeUnit.SECONDS), FiniteDuration(1, TimeUnit.SECONDS), {
     implicit val writer = shared.Protocol.chatMessageRW
     TextMessage(write(ChatMessage(sender = "Bob", message = "Hi")))
   })
   (sink, source2) // This works
   (sing, source) // This does not
 }

我怎样才能制作这样一个可以与Akka Actor 集成的Source

【问题讨论】:

    标签: scala websocket akka


    【解决方案1】:

    Source.actorRefSource.actorRefWithBackpressure 都创建了作为物化值提供的参与者。但是,服务器端 WS API 并不能轻松访问该具体化的值。

    获取已创建的actorRef最简单的方法是使用mapMaterializedValue

       val source = Source.actorRefWithBackpressure[Message](AckMessage, {
         case _: Success => CompletionStrategy.draining
       }, PartialFunction.empty).mapMaterializedValue { actorRef =>
         // Do something with the ActorRef here. Messages you want to send to this client will have to be sent to this ActorRef.
         // e.g.: chatRoom ! NewClient(actorRef)
       }
    

    我的聊天室示例的先前版本仍然基于 Actor,并在完整示例中显示了这一点:

    https://github.com/jrudolph/akka-http-scala-js-websocket-chat/blob/b01b234376c4984dce19effcaf001a9ffb4c6981/backend/src/main/scala/example/akkawschat/Chat.scala#L64

    【讨论】:

      猜你喜欢
      • 2018-07-17
      • 2017-04-15
      • 1970-01-01
      • 2016-08-16
      • 1970-01-01
      • 2020-04-15
      • 2015-06-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多