【问题标题】:How to attach multiple actors as sources to an Akka stream?如何将多个演员作为源附加到 Akka 流?
【发布时间】:2015-05-06 13:09:44
【问题描述】:

我正在尝试构建和运行一个 akka 流(在 Java DSL 中),其中 2 个 Actor 作为源,然后是一个合并结点,然后是 1 个接收器:

    Source<Integer, ActorRef> src1 = Source.actorRef(100, OverflowStrategy.backpressure());
    Source<Integer, ActorRef> src2 = Source.actorRef(100, OverflowStrategy.backpressure());
    Sink<Integer, BoxedUnit> sink = Flow.of(Integer.class).to(Sink.foreach(System.out::println));

    RunnableFlow<BoxedUnit> closed = FlowGraph.factory().closed(sink, (b, out) -> {
        UniformFanInShape<Integer, Integer> merge = b.graph(Merge.<Integer>create(2));
        b.from(src1).via(merge).to(out);
        b.from(src2).to(merge);
    });

    closed.run(mat);

我的问题是如何获取对源演员的 ActorRef 引用以便向他们发送消息?如果有 1 个演员,我不会使用图形生成器,然后 .run() 或 runWith() 方法将返回 ActorRef 对象。但是如果有很多源演员怎么办?有没有可能实现这样的流程?

【问题讨论】:

  • 您需要将需要访问物化值的元素传递给closed,然后提供一个组合物化值的函数。像这样的东西:closed(src1, src2, (actorRef1, actorRef2) -&gt; SomethingContainingBothActorRefs, (b, s1, s2) -&gt; ...)

标签: akka akka-stream


【解决方案1】:

回答我自己的问题以防有人需要。

根据 jrudolph 的建议,我能够使用这样的演员(在实际代码中,我做了比 2 个 ActorRef 列表更好的事情):

    Source<Integer, ActorRef> src1 = Source.actorRef(100, OverflowStrategy.fail());
    Source<Integer, ActorRef> src2 = Source.actorRef(100, OverflowStrategy.fail());
    Sink<Integer, BoxedUnit> sink = Flow.of(Integer.class).to(Sink.foreach(System.out::println));

    RunnableFlow<List<ActorRef>> closed = FlowGraph.factory().closed(src1, src2, (a1, a2) -> Arrays.asList(a1, a2), (b, s1, s2) -> {
        UniformFanInShape<Integer, Integer> merge = b.graph(Merge.<Integer>create(2));
        b.from(s1).via(merge).to(sink);
        b.from(s2).to(merge);
    });

    List<ActorRef> stream = closed.run(mat);
    ActorRef a1 = stream.get(0);
    ActorRef a2 = stream.get(1);

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-12-15
    • 1970-01-01
    • 2015-10-01
    • 1970-01-01
    • 2014-11-15
    • 2011-11-17
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多