【问题标题】:How to provide subscriber with routing capability如何为订户提供路由能力
【发布时间】:2014-11-19 10:15:42
【问题描述】:

我有以下流的路径-

 kafkaStream[message] -> 
 kafkaStream[message] -> mergedKafkaStream[message] -> stream[EnrichedMessage] -> I/O
 kafkaStream[message] -> 

我不确定如何以 akka 流的方式编写它。 我尝试了以下(伪)。

KafkaStream extends ActorPublisher[message] {

}

IOHandler extends ActorSubscriber {

}

k1、k2、k3 是 kafka 流发布者

f = Flow[message].map(_.enrichMessage)

FlowGraph { b =>
  k1 ~> merge
  k2 ~> merge
  k3 ~> merge
  merge ~> f ~> ioHandlerSink
}

这就是我将发布者连接到接收器的方式。但是这里我要解决的问题是缓慢的 IO。 IOHandler 演员处理消息的速度非常慢,所以我如何拥有多个 IOHandler 并且我应该能够分配任务。而且我还想保持背压,所以不要用火而忘记使用路由器。

我对 akka 流很陌生,所以请给我一个出路。

谢谢

【问题讨论】:

  • akka-streams 本身是新的,因此功能可能仍然缺失 :) 特别是 ticket #15959 建议的内容可能与您正在寻找的内容相似。它尚未实现,但如果您添加关于您希望它做什么的评论可能会有所帮助。

标签: akka akka-stream


【解决方案1】:

您可以使用FlowGraph 中的Balance 结点路由到多个IOHandler 接收器。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-08-25
    • 2012-01-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-06-09
    • 1970-01-01
    相关资源
    最近更新 更多