【问题标题】:Akka streams: attach new publishers/subscribers to FlowAkka 流:将新的发布者/订阅者附加到 Flow
【发布时间】: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


【解决方案1】:

我在我的一个项目中使用以下解决方案来处理来自多个请求处理器的来自 websocket 的请求,这些请求可以产生响应流或提供无限订阅。

// requests coming from websocket, it could be any source, it's doesn't matter
val requests: Source[Request, NotUsed] = ... 

// the request processing here can provide endles stream of responses
val requestProcessing: Flow[Request, Response, NotUsed] = ...

val (outSink, outSource) =
  MergeHub
    .source[Result](perProducerBufferSize = 4)
    .toMat(BroadcastHub.sink(bufferSize = 32))(Keep.both)
    .run()

Source.tick(Duration.Zero, KeepAliveInterval, ConnectionKeepAlive)
  .to(outSink)
  .run()

requests.fold {
  case State(state, AuthRequest(r)) if checkAuth(r) => 
    Source.single(AuthenticationAck).to(outSink)
    state.copy(isAuthenticated = true)
 case State(state, AuthRequest(r)) => 
    Source.single(AuthenticationFailedError).to(outSink)
    state

 case State(state, request) if s.isAuthenticated =>
    // here the most of busines
    Source.single(request).via(requestProcessing).to(outSink)
    state

 case State(state, _) => 
    Source.single(NonAuthorizedError).to(outSink)
    state
}.toMat(Sink.ignore)(Keep.right)

outSource.runForeach { response =>
  // here we get the stream of responses mixed from all requests
}

outSource.runForeach { response =>
  // of course, we could have as many subscribers as we need
}

希望对你有帮助:)

【讨论】:

    【解决方案2】:

    对于这个特殊问题,我不会使用akka-stream。您描述的多播发布订阅类型更适合原始Actor 消息传递和EventStream。

    在某些情况下,我是 akka-stream 的超级粉丝,但在这种情况下,我认为您正在尝试将方形钉子穿过圆孔。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-12-14
      • 2019-07-19
      • 2021-07-28
      • 1970-01-01
      相关资源
      最近更新 更多