【发布时间】:2018-12-17 01:22:06
【问题描述】:
想了解如何使用 Akka.net Streams 从 Source.Queue 并行处理项目,并在 Actor 中完成处理。
我已经能够使用 Sink.ForEachParallel 调用一个函数,并且它按预期工作。
是否可以与 Sink.ActorRefWithAck 并行处理项目(因为我希望它利用背压)?
【问题讨论】:
标签: akka-stream akka.net
想了解如何使用 Akka.net Streams 从 Source.Queue 并行处理项目,并在 Actor 中完成处理。
我已经能够使用 Sink.ForEachParallel 调用一个函数,并且它按预期工作。
是否可以与 Sink.ActorRefWithAck 并行处理项目(因为我希望它利用背压)?
【问题讨论】:
标签: akka-stream akka.net
当试图将之前的尝试和中提琴结合起来时,即将按下 Post!
当我尝试在其中创建参与者时,之前使用 ForEachParallel 的尝试失败,但无法在异步函数中执行此操作。如果我使用之前声明的单个演员,那么 Tell 可以工作,但我无法获得我想要的并行度。
我让它与具有循环配置的路由器一起使用。
var props = new RoundRobinPool(5).Props(Props.Create<MyActor>());
var actor = Context.ActorOf(props);
flow = Source.Queue<Element>(2000,OverflowStrategy.Backpressure)
.Select(x => {
return new Wrapper() { Element = x, Request = ++cnt };
})
.To(Sink.ForEachParallel<Wrapper>(5, (s) => { actor.Tell(s); }))
.Run(materializer);
Request ++cnt 用于控制台输出,以验证请求是否按要求处理。
MyActor 每隔 10 个请求就会有很长的延迟,以验证背压是否正常工作。
【讨论】:
actor.Tell 发送消息,它不会产生背压 - 消息只会开始堆积在演员的邮箱中。如果你想要真正的背压,你可以用.SelectAsync(5, s => actor.Ask<Ack>(s, timeout)).To(Sink.Ignore<Ack>()) 改变 ForEachParallel 并让 Actor 将 Ack 消息发送回背压。