我会考虑使用akka-cluster 和akka-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 发送消息,而是使用MyEventReceiver 和deviceRegion Actors,您使用分片系统跟踪特定 EventActor 生活在集群中的哪个节点在。如果之前没有创建过,它将创建一个,或者如果有,则将消息路由给它。每个 EventActor 必须有一个唯一的 id:从 message 中提取(所以 UUID 非常适合,但它可以是其他一些 id,如 customerID 或 orderID ,或其他任何内容,只要它对于您要处理的 Actor 实例是唯一的)。
(我省略了 EventActor 代码,否则它是一个非常普通的 Actor,取决于你用它做什么,“魔法”在上面的代码中)。
分片系统根据您选择的算法自动知道创建EventActor并将其分配给一个分片(在这种特殊情况下,它基于唯一ID的hashCode,这就是全部我用过)。此外,对于任何给定的唯一 ID,您保证只有一个 Actor。消息被透明地路由到正确的节点和分片,无论它在哪里;从任何节点和分片发送它。
Akka 站点和文档中有更多信息和示例代码。
这是确保同一个 Entity/Actor 始终处理为其指定的消息的一种非常好的方法。集群和分片会自动负责正确分配 Actor,以及故障转移等(如果 Actor 有一堆与之关联的严格状态,则必须添加 akka-persistence 以获得钝化、再水化和故障转移(必须恢复))。