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