【发布时间】: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