【发布时间】: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 它)。 Sink 和 source2 按预期工作,但 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?
【问题讨论】: