【问题标题】:Akka broadcast: get the first reply and discard othersAkka广播:得到第一个回复并丢弃其他人
【发布时间】:2017-11-11 16:45:04
【问题描述】:

我有一个带有工作池的 CalculationSupervisor 演员。
每次我需要进行一些计算时,CalculationSupervisor 使用路由器向工作人员广播CalculationRequest

我需要得到最快计算的结果,而忽略其他结果。

CalculationSupervisor 如下所示:

public class CalculationSupervisor extends AbstractActor {

    private Router router = new Router(new RoundRobinRoutingLogic());

    public static Props props() {
        return Props.create(CalculationSupervisor.class, CalculationSupervisor::new);
    }

    @Override
    public Receive createReceive() {
        return receiveBuilder()
                .match(RegisterWorker.class, registration -> {
                    final String workerName = registration.name();
                    final ActorRef worker = 
                        context().actorOf(Worker.props(workerName), workerName);
                    router = router.addRoutee(worker);
                })
                .match(CalculationRequest.class, (request) -> {
                    router.route(new Broadcast(request), self());
                })
                .match(CalculationResult.class, (result) -> {
                    // process only the first (the fastest) result
                })
                .build();
    }
}

实现丢弃第一个(最快)结果后出现的消息的逻辑的最佳模式是什么?

【问题讨论】:

    标签: java akka broadcast


    【解决方案1】:

    如果您的主管正在处理多个请求,一个简单的方法是保留一个 Map,其中包含请求 ID 对和收到的请求响应数。主管检查此映射并仅在处理时该特定请求 ID 的回复数为零时才处理回复。此外,为了防止映射无限增长,如果结果数等于池大小,主管会从映射中删除条目:

    public class CalculationSupervisor extends AbstractActor {
        ...
        private int poolSize = 0;
        private Map<Long, Integer> numReplies = new HashMap<>();
    
        @Override
        public Receive createReceive() {
            return receiveBuilder()
                .match(RegisterWorker.class, registration -> {
                    ...
                    poolSize = poolSize + 1;
                })
                .match(CalculationRequest.class, request -> {
                    numReplies.putIfAbsent(request.getId(), 0);
                    router.route(new Broadcast(request), self());
                })
                .match(CalculationResult.class, result -> {
                    Long requestId = result.getRequestId();
                    if (numReplies.contains(requestId)) {
                        int num = numReplies.get(requestId);
                        if (num == 0) {
                            // process only the first (the fastest) result
                            ...
                            numReplies.put(requestId, 1);
                        } else {
                            if (num + 1 == poolSize)
                                numReplies.remove(requestId);
                            else
                                numReplies.put(requestId, num + 1);
                        }
                    }
                })
                .build();
        }
    }
    

    上述方法有两个假设:

    1. CalculationRequestCalculationResult 类中都提供了请求 ID(在本例中,ID 是 Long;使用任何合适的 ID)。
    2. 在发送请求之前,路由已向主管注册。

    一个更简单的解决方案是不使用路由器,在这种情况下CalculationSupervisor 不需要为同一个请求协调多个结果。由于您要丢弃每个请求的所有结果,但最早的结果除外,因此首先使用路由器是没有意义的。

    【讨论】:

    • 有很多请求进入 CalculationSupervisor,我想全部处理它们。在您的方法中,如果我得到任何传入 CalculationRequest 的 CalculationResult,我将停止处理其他请求的结果。可以使用一些地图存储计算请求标识符及其状态来扩展您的解决方案。这是我最初的想法,但我想找到更优雅的解决方案:)
    猜你喜欢
    • 1970-01-01
    • 2018-09-01
    • 1970-01-01
    • 2016-10-07
    • 1970-01-01
    • 2016-03-21
    • 2020-11-24
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多