Akka分布式Pub/Sub背压

bit*_*tan 5 java publish-subscribe akka akka-cluster akka-stream

我正在使用Akka Distributed 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();
    }
}
Run Code Online (Sandbox Code Playgroud)

订户代码:

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();
    }
}
Run Code Online (Sandbox Code Playgroud)

Ram*_*gil 1

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

切换到询问

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

切换到流

“背压”是 akka 流的一个主要特征。因此,通过切换到流实现,您将能够实现您想要的目标。

如果可以Source从原始数据创建流,那么您可以使用Flow.throttleSink.actorRef来创建流Sinkmediator使用Flow.throttle来控制流向中介器的速率。