【问题标题】:Simple server-push broadcast flow with Akka使用 Akka 的简单服务器推送广播流程
【发布时间】:2019-02-14 20:49:56
【问题描述】:

我正在努力实现一个相当简单的 Akka 流程。 这是我认为我需要的:

我有一个服务器和 n 个客户端,并且希望能够通过向客户端广播消息 (JSON) 来响应外部事件。客户可以随时注册/注销。

例如:

  • 1 个客户已注册
  • 服务器抛出一个事件(“Hello World!”)
  • 服务器广播“Hello World!”给所有客户(一个客户)
  • 一个新的客户端打开一个 websocket 连接
  • 服务器引发另一个事件(“Hello Akka!”)
  • 服务器广播“Hello Akka!”给所有客户(两个客户)

这是我目前所拥有的:

def route: Route = {
   val register = path("register") {
     // registration point for the clients
     handleWebSocketMessages(serverPushFlow)
   }
}

// ...

def broadcast(msg: String): Unit = {
  // use the previously created flow to send messages to all clients
  // ???
}

// my broadcast sink to send messages to the clients
val broadcastSink: Sink[String, Source[String, NotUsed]] = BroadcastHub.sink[String]

// a source that emmits simple strings
val simpleMsgSource = Source(Nil: List[String])

def serverPushFlow = {
  Flow[Message].mapAsync(1) {
    case TextMessage.Strict(text) =>       Future.successful(text)
    case streamed: TextMessage.Streamed => streamed.textStream.runFold("")(_ ++ _)
  }
  .via(Flow.fromSinkAndSource(broadcastSink, simpleMsgSource))
  .map[Message](string => TextMessage(string))
}

【问题讨论】:

    标签: scala akka akka-stream akka-http


    【解决方案1】:

    为了能够使用广播集线器,您必须定义两个流。一个运行你的 websocket TextMessagebroadcastHub。您必须运行它,它会生成一个连接到每个客户端的 Source。

    这是简单可运行应用程序中描述的这个概念。

    import akka.NotUsed
    import akka.actor.ActorSystem
    import akka.stream.ActorMaterializer
    import akka.stream.scaladsl.{BroadcastHub, Sink, Source}
    import org.slf4j.LoggerFactory
    
    import scala.concurrent.duration._
    
    object BroadcastSink extends App {
    
      private val logger = LoggerFactory.getLogger("logger")
    
      implicit val actorSystem = ActorSystem()
      implicit val actorMaterializer = ActorMaterializer()
    
      val broadcastSink: Sink[String, Source[String, NotUsed]] =
        BroadcastHub.sink[String]
    
      val simpleMsgSource = Source.tick(500.milli, 500.milli, "Single Message")
    
      val sourceForClients: Source[String, NotUsed] = simpleMsgSource.runWith(broadcastSink)
    
      sourceForClients.to(Sink.foreach(t => logger.info(s"Client 1: $t"))).run()
      Thread.sleep(1000)
    
      sourceForClients.to(Sink.foreach(t => logger.info(s"Client 2: $t"))).run()
      Thread.sleep(1000)
    
      sourceForClients.to(Sink.foreach(t => logger.info(s"Client 3: $t"))).run()
      Thread.sleep(1000)
    
      sourceForClients.to(Sink.foreach(t => logger.info(s"Client 4: $t"))).run()
      Thread.sleep(1000)
    
      actorSystem.terminate()
    }
    

    打印

    10:52:01.774 Client 1: Single Message
    10:52:02.273 Client 1: Single Message
    10:52:02.273 Client 2: Single Message
    10:52:02.773 Client 2: Single Message
    10:52:02.773 Client 1: Single Message
    10:52:03.272 Client 3: Single Message
    10:52:03.272 Client 2: Single Message
    10:52:03.272 Client 1: Single Message
    10:52:03.772 Client 1: Single Message
    10:52:03.772 Client 3: Single Message
    10:52:03.773 Client 2: Single Message
    10:52:04.272 Client 2: Single Message
    10:52:04.272 Client 4: Single Message
    10:52:04.272 Client 1: Single Message
    10:52:04.273 Client 3: Single Message
    10:52:04.772 Client 1: Single Message
    10:52:04.772 Client 2: Single Message
    10:52:04.772 Client 3: Single Message
    10:52:04.772 Client 4: Single Message
    10:52:05.271 Client 4: Single Message
    10:52:05.271 Client 1: Single Message
    10:52:05.271 Client 3: Single Message
    10:52:05.272 Client 2: Single Message
    

    如果事先知道客户,则不需要BrodacastHub,可以使用alsoTo方法:

      def webSocketHandler(clients: List[Sink[Message, NotUsed]]): Flow[Message, Message, Any] = {
        val flow = Flow[Message]
        clients.foldLeft(flow) {case (fl, client) =>
          fl.alsoTo(client)
        }
      }
    

    【讨论】:

    • 对不起,我不太明白这里的两件事: 1. 除了“单一消息”之外,您如何发送其他内容? 2. 你如何通过路由连接客户端?
    • 更新了答案。 1.只是一个例子,消息来源来自akka http。 2. 查看更新,在这种情况下不需要 BrodacastHub
    • "这只是一个例子。消息的来源来自 akka http" 但是在我的用例中,消息不是来自 akka-http。如何获得触发广播的事件(服务器推送)?
    • 我更新了我的问题,因为我认为它引起了一些误解。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-04-11
    • 2017-06-12
    • 1970-01-01
    • 2017-02-11
    • 2015-06-22
    相关资源
    最近更新 更多