【问题标题】:Akka Distributed Pub/Sub back-pressureAkka 分布式 Pub/Sub 背压
【发布时间】:2018-05-24 12:33:33
【问题描述】:

我正在使用 Akka 分布式 Pub/Sub,并且只有一个发布者和一个订阅者。我的发布者比订阅者快得多。有没有办法在某个时间点后减慢发布者的速度?

发布者代码:

public class Publisher extends AbstractActor {
    private ActorRef mediator;

    static public Props props() {
        return Props.create(Publisher.class, () -> new Publisher());
    }

    public Publisher () {
        this.mediator = DistributedPubSub.get(getContext().system()).mediator();
        this.self().tell(0, ActorRef.noSender());
    }

    @Override
    public Receive createReceive() {
        return receiveBuilder()
            .match(Integer.class, msg -> {
                // Sending message to Subscriber
                mediator.tell(
                    new DistributedPubSubMediator.Send(
                        "/user/" + Subscriber.class.getName(),
                        msg.toString(),
                        false),
                    getSelf());

                getSelf().tell(++msg, ActorRef.noSender());
            })
            .build();
    }
}

订阅者代码:

public class Subscriber extends AbstractActor {
    static public Props props() {
        return Props.create(Subscriber.class, () -> new Subscriber());
    }

    public Subscriber () {
        ActorRef mediator = DistributedPubSub.get(getContext().system()).mediator();
        mediator.tell(new DistributedPubSubMediator.Put(getSelf()), getSelf());
    }

    @Override
    public Receive createReceive() {
        return receiveBuilder()
            .match(String.class, msg -> {
                System.out.println("Subscriber message received: " + msg);
                Thread.sleep(10000);
            })
            .build();
    }
}

【问题讨论】:

    标签: java akka publish-subscribe akka-stream akka-cluster


    【解决方案1】:

    不幸的是,按照目前的设计,我认为没有办法为原始发件人提供“背压”。由于您使用ActorRef.tell 将消息发送到mediator,因此无法获得下游接收器正在备份的信号。这是因为您正在使用的方法 tell 返回一个 void

    切换到提问

    如果您将tell 切换为ask,您可以设置一个适当的Timeout 值,该值至少可以让您知道您在特定时间内没有收到回复。

    切换到流

    "Back-pressure" is a primary feature of akka streams。因此,通过切换到流实现,您将能够实现您想要的目标。

    如果可以从原始数据创建流Source,那么您可以使用Sink.actorRefmediator 创建Sink,并使用Flow.throttle 控制流向中介的速率.

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2022-01-13
      • 2021-01-19
      • 2017-09-07
      • 2021-08-16
      • 1970-01-01
      • 2022-01-01
      • 2017-09-25
      相关资源
      最近更新 更多