代码之家  ›  专栏  ›  技术社区  ›  Flame_Phoenix

如何使用akka流在图中限制请求?

  •  0
  • Flame_Phoenix  · 技术社区  · 7 年前

    出身背景

    在这个项目中,我有一个字符串流和一个对它们进行一些操作的图形。

    客观的

    在我的图表中,我想将该流广播给2个工人。其中一个将替换所有字符 'a' 具有 'A' 并在接收数据时实时发送数据。

    另一个将接收数据,每3个字符串,它将连接这3个字符串并将它们映射到数字。

    它将如下所示:

    akka-streams-buffer

    明显地 Sink 2 接收信息的速度不如 Sink 1 . 但这是预期的行为。这里有趣的部分是worker 2。

    问题

    做工人1很容易,但不难。这里的问题是做工人2。我知道akka有最多可以保存X条消息的缓冲区,但看起来我不得不选择一个现有的缓冲区 Overflow strategies 这通常会导致选择要删除的消息,或者选择是否要使流保持活动状态。

    但即使在阅读了 stream-rate akka的文档我找不到一种方法,至少使用Java。

    研究

    我还查了一个类似的问题, Selective request-throttling using akka-http stream 然而,已经过去一年多了,没有人回应。

    使用图形DSL,我将如何从以下位置创建路径:

    来源->bcast->工人2->水槽2

    ??

    1 回复  |  直到 7 年前
        1
  •  1
  •   0x26res    7 年前

    在你 bcast 应用 groupedWithin https://doc.akka.io/docs/akka/2.5/stream/operators/Source-or-Flow/groupedWithin.html

    List 并在每次到达3个元素时发出列表。

    import akka.stream.Attributes;
    import akka.stream.FlowShape;
    import akka.stream.Inlet;
    import akka.stream.Outlet;
    import akka.stream.stage.AbstractInHandler;
    import akka.stream.stage.GraphStage;
    import akka.stream.stage.GraphStageLogic;
    import com.google.common.collect.ImmutableList;
    import java.util.ArrayList;
    import java.util.List;
    
    public class RecordGrouper<T> extends GraphStage<FlowShape<T, List<T>>> {
    
      private final Inlet<T> inlet = Inlet.create("in");
      private final Outlet<List<T>> outlet = Outlet.create("out");
      private final FlowShape<T, List<T>> shape = new FlowShape<>(inlet, outlet);
    
      @Override
      public GraphStageLogic createLogic(Attributes inheritedAttributes) {
        return new GraphStageLogic(shape) {
          List<T> batch = new ArrayList<>(3);
    
          {
            setHandler(
                inlet,
                new AbstractInHandler() {
                  @Override
                  public void onPush() {
                    T record = grab(inlet);
                    batch.add(record);
                    if (batch.size() == 3) {
                      emit(outlet, ImmutableList.copyOf(batch));
                      batch.clear();
                    }
                    pull(inlet);
                  }
                });
          }
    
          @Override
          public void preStart() {
            pull(inlet);
          }
        };
      }
    
      @Override
      public FlowShape<T, List<T>> shape() {
        return shape;
      }
    }
    

    作为侧节点,我不认为 buffer https://doc.akka.io/docs/akka/2.5/stream/operators/Source-or-Flow/buffer.html