【问题标题】:Retrieve an Akka actor or create it if it does not exist检索一个 Akka actor,如果它不存在则创建它
【发布时间】:2019-08-09 07:04:52
【问题描述】:

我正在开发一个应用程序,它创建一些 Akka Actor 来管理和处理来自 Kafka 主题的消息。具有相同密钥的消息由相同的参与者处理。我也使用消息键来命名相应的演员。

当从主题中读取一条新消息时,我不知道id等于消息键的actor是否已经由actor系统创建。因此,我尝试使用它的名字来解析这个actor,如果它还不存在,我就创建它。我需要管理有关参与者解析的并发性。所以有可能不止一个客户端询问actor系统是否存在actor。

我现在使用的代码如下:

private CompletableFuture<ActorRef> getActor(String uuid) {
    return system.actorSelection(String.format("/user/%s", uuid))
                 .resolveOne(Duration.ofMillis(1000))
                 .toCompletableFuture()
                 .exceptionally(ex -> 
                     system.actorOf(Props.create(MyActor.class, uuid), uuid))
                 .exceptionally(ex -> {
                     try {
                         return system.actorSelection(String.format("/user/%s",uuid)).resolveOne(Duration.ofMillis(1000)).toCompletableFuture().get();
                     } catch (InterruptedException | ExecutionException e) {
                         throw new RuntimeException(e);
                     }
                 });
}

以上代码没有优化,异常处理可以做得更好。

但是,在 Akka 中是否有一种更惯用的方式来解析一个 actor,或者如果它不存在则创建它?我错过了什么吗?

【问题讨论】:

    标签: java scala akka actor


    【解决方案1】:

    考虑创建一个actor,将消息ID 映射到ActorRefs 作为其状态。这个“接待员”参与者将处理所有请求以获取消息处理参与者。当接待员收到一个演员的请求(该请求将包含消息 ID)时,它会尝试在其映射中查找关联的演员:如果找到了这样的演员,它会将ActorRef 返回给发送者;否则,它会创建一个新的处理参与者,将该参与者添加到其映射中,并将该参与者引用返回给发送者。

    【讨论】:

    • 感谢您的回复。这可能是个好主意。但是,我说的是数以百万计的演员 :( 我知道,它应该插入到问题中......
    • 在这种情况下,请查看 Akka Clustering。它有效地做同样的事情,但参与者通过您定义的某些策略跨多个分区进行分片。消息使用分区策略路由到特定分区,然后分区将它们定向到正确的参与者。如果演员不存在,它被创建。如果您的数百万演员需要,这也可以让您跨多个节点进行横向扩展。
    【解决方案2】:

    我会考虑使用akka-clusterakka-cluster-sharding。首先,这为您提供了吞吐量,以及可靠性。但是,它也会让系统管理“实体”actors 的创建。

    但你必须改变与这些演员交谈的方式。您创建一个处理所有消息的ShardRegion 演员:

    import akka.actor.AbstractActor;
    import akka.actor.ActorRef;
    import akka.actor.ActorSystem;
    import akka.actor.Props;
    import akka.cluster.sharding.ClusterSharding;
    import akka.cluster.sharding.ClusterShardingSettings;
    import akka.cluster.sharding.ShardRegion;
    import akka.event.Logging;
    import akka.event.LoggingAdapter;
    
    
    public class MyEventReceiver extends AbstractActor {
    
        private final ActorRef shardRegion;
    
        public static Props props() {
            return Props.create(MyEventReceiver.class, MyEventReceiver::new);
        }
    
        static ShardRegion.MessageExtractor messageExtractor 
          = new ShardRegion.HashCodeMessageExtractor(100) {
                // using the supplied hash code extractor to shard
                // the actors based on the hashcode of the entityid
    
            @Override
            public String entityId(Object message) {
                if (message instanceof EventInput) {
                    return ((EventInput) message).uuid().toString();
                }
                return null;
            }
    
            @Override
            public Object entityMessage(Object message) {
                if (message instanceof EventInput) {
                    return message;
                }
                return message; // I don't know why they do this it's in the sample
            }
        };
    
    
        public MyEventReceiver() {
            ActorSystem system = getContext().getSystem();
            ClusterShardingSettings settings =
               ClusterShardingSettings.create(system);
            // this is setup for the money shot
            shardRegion = ClusterSharding.get(system)
                    .start("EventShardingSytem",
                            Props.create(EventActor.class),
                            settings,
                            messageExtractor);
        }
    
        @Override
        public Receive createReceive() {
            return receiveBuilder().match(
                    EventInput.class,
                    e -> {
                        log.info("Got an event with UUID {} forwarding ... ",
                                e.uuid());
                        // the money shot
                        deviceRegion.tell(e, getSender());
                    }
            ).build();
        }
    }
    

    所以这个 Actor MyEventReceiver 在集群的所有节点上运行,并封装了 shardRegion Actor。您不再直接向EventActors 发送消息,而是使用MyEventReceiverdeviceRegion Actors,您使用分片系统跟踪特定 EventActor 生活在集群中的哪个节点在。如果之前没有创建过,它将创建一个,或者如果有,则将消息路由给它。每个 EventActor 必须有一个唯一的 id:从 message 中提取(所以 UUID 非常适合,但它可以是其他一些 id,如 customerID 或 orderID ,或其他任何内容,只要它对于您要处理的 Actor 实例是唯一的)。

    (我省略了 EventActor 代码,否则它是一个非常普通的 Actor,取决于你用它做什么,“魔法”在上面的代码中)。

    分片系统根据您选择的算法自动知道创建EventActor并将其分配给一个分片(在这种特殊情况下,它基于唯一ID的hashCode,这就是全部我用过)。此外,对于任何给定的唯一 ID,您保证只有一个 Actor。消息被透明地路由到正确的节点和分片,无论它在哪里;从任何节点和分片发送它。

    Akka 站点和文档中有更多信息和示例代码。

    这是确保同一个 Entity/Actor 始终处理为其指定的消息的一种非常好的方法。集群和分片会自动负责正确分配 Actor,以及故障转移等(如果 Actor 有一堆与之关联的严格状态,则必须添加 akka-persistence 以获得钝化、再水化和故障转移(必须恢复))。

    【讨论】:

      【解决方案3】:

      Jeffrey Chung 的回答确实是 Akka 方式。这种方法的缺点是性能低下。最高效的解决方案是使用 Java 的 ConcurrentHashMap.computeIfAbsent() 方法。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2016-10-09
        • 1970-01-01
        • 2018-07-17
        • 2014-07-20
        • 2017-10-15
        • 2022-01-11
        • 2021-02-28
        • 2016-09-18
        相关资源
        最近更新 更多