【发布时间】:2018-08-02 23:07:22
【问题描述】:
我正在构建一个 Akka 应用程序,并希望将某些参与者的 FSM 状态转换公开给外部消费者。 (最终目标是能够将状态转换消息推送到 websocket,以便实时查看。)
根据文档Combining dynamic stages to build a simple Publish-Subscribe service,看来我需要公开表示发布-订阅通道的流,以便消费者和生产者可以使用它。
我遇到问题的部分是将新的 Source(s) 附加到 Flow 中,以便生成的每个新 Actor 都会将其状态转换发布到 Source。另一个问题是向 Flow 添加新的 Sink(最终,这将是 websockets,但出于测试目的,它可以是任何新的 Sink)。
首先我连接一个 MergeHub 和一个 BroadcastHub 以形成一个“通道”,然后从物化的接收器和源创建流:
val orderFlow: Flow[String, String, NotUsed] = {
val (sink, source) = MergeHub.source[String](16)
.toMat(BroadcastHub.sink(256))(Keep.both).run()
Flow.fromSinkAndSource(sink, source)
}
那么问题是我如何动态地将新的生产者和消费者添加到这个流中?有什么想法吗?
【问题讨论】:
-
你有代码吗?您是否尝试过隔离 pub-sub 关注点并编写一个简化版本以及对其进行测试?当问题太宽泛或不清楚时,人们很难帮助你。
-
添加了一个代码 sn-p,或许可以更好地说明问题。
标签: scala websocket akka publish-subscribe akka-stream