Akka Classic与Typed混合集群是否支持分布式Pub-Sub?求示例代码
Akka Classic与Typed混合集群的分布式Pub-Sub互通方案
核心结论
Akka 2.6的共存模式支持Classic与Typed节点间的分布式Pub-Sub消息互通。两者的Pub-Sub实现底层依赖相同的集群分布式数据结构,只需通过API适配+统一配置,即可实现跨类型节点的消息发布与订阅。
实现方案与示例代码
以下基于Java语言提供双向互通的代码示例,核心是保证两端使用相同的主题名称,并确保消息序列化兼容。
通用配置前提
所有集群节点需启用共存模式(Akka 2.6+默认开启),且集群配置(种子节点、角色、序列化规则)保持一致。
场景1:Classic节点发布消息,Typed节点订阅
Typed端订阅实现
import akka.actor.typed.ActorSystem; import akka.actor.typed.Behavior; import akka.actor.typed.javadsl.AbstractBehavior; import akka.actor.typed.javadsl.ActorContext; import akka.actor.typed.javadsl.Behaviors; import akka.cluster.typed.Subscribe; import akka.cluster.typed.SubscribeOptions; import akka.cluster.typed.Topic; // 定义跨节点共享的消息类型(确保序列化兼容) public interface PubSubMsg { record ClassicMsg(String content) implements PubSubMsg {} } public class TypedSubscriber extends AbstractBehavior<PubSubMsg> { public static Behavior<PubSubMsg> create() { return Behaviors.setup(TypedSubscriber::new); } private TypedSubscriber(ActorContext<PubSubMsg> context) { super(context); ActorSystem<?> system = context.getSystem(); // 创建与Classic端同名的主题 Topic<PubSubMsg> sharedTopic = Topic.create(system, "cluster-shared-topic", PubSubMsg.class); // 订阅主题 context.getSelf().tell(Subscribe.create(sharedTopic.ref(), SubscribeOptions.withAllHistory())); } @Override public Receive<PubSubMsg> createReceive() { return newReceiveBuilder() .onMessage(PubSubMsg.ClassicMsg.class, this::handleClassicMessage) .build(); } private Behavior<PubSubMsg> handleClassicMessage(PubSubMsg.ClassicMsg msg) { getContext().getLog().info("Typed节点收到Classic消息: {}", msg.content()); return this; } } // 启动Typed节点 public class TypedNodeBootstrap { public static void main(String[] args) { ActorSystem.create(TypedSubscriber.create(), "akka-cluster"); } }
Classic端发布实现
import akka.actor.ActorRef; import akka.actor.ActorSystem; import akka.cluster.pubsub.DistributedPubSub; import akka.cluster.pubsub.DistributedPubSubMediator; import akka.japi.pf.ReceiveBuilder; public class ClassicPublisher extends akka.actor.AbstractActor { private final ActorRef mediator; public ClassicPublisher() { mediator = DistributedPubSub.get(getContext().system()).mediator(); } public static akka.actor.Props props() { return akka.actor.Props.create(ClassicPublisher.class); } @Override public Receive createReceive() { return ReceiveBuilder.create() .match(String.class, msg -> { // 向共享主题发布消息 mediator.tell(new DistributedPubSubMediator.Publish("cluster-shared-topic", msg), getSelf()); }) .build(); } // 启动Classic节点并发布测试消息 public static void main(String[] args) { ActorSystem system = ActorSystem.create("akka-cluster"); ActorRef publisher = system.actorOf(ClassicPublisher.props(), "classic-publisher"); publisher.tell("来自Classic节点的问候", ActorRef.noSender()); } }
场景2:Typed节点发布消息,Classic节点订阅
Classic端订阅实现
import akka.actor.ActorRef; import akka.actor.ActorSystem; import akka.cluster.pubsub.DistributedPubSub; import akka.cluster.pubsub.DistributedPubSubMediator; import akka.japi.pf.ReceiveBuilder; public class ClassicSubscriber extends akka.actor.AbstractActor { public ClassicSubscriber() { ActorRef mediator = DistributedPubSub.get(getContext().system()).mediator(); // 订阅共享主题 mediator.tell(new DistributedPubSubMediator.Subscribe("cluster-shared-topic", getSelf()), getSelf()); } public static akka.actor.Props props() { return akka.actor.Props.create(ClassicSubscriber.class); } @Override public Receive createReceive() { return ReceiveBuilder.create() .match(PubSubMsg.TypedMsg.class, msg -> { getContext().getLog().info("Classic节点收到Typed消息: {}", msg.content()); }) .build(); } // 启动Classic节点 public static void main(String[] args) { ActorSystem system = ActorSystem.create("akka-cluster"); system.actorOf(ClassicSubscriber.props(), "classic-subscriber"); } }
Typed端发布实现
import akka.actor.typed.ActorSystem; import akka.actor.typed.Behavior; import akka.actor.typed.javadsl.AbstractBehavior; import akka.actor.typed.javadsl.ActorContext; import akka.actor.typed.javadsl.Behaviors; import akka.cluster.typed.Topic; public interface PubSubMsg { record TypedMsg(String content) implements PubSubMsg {} } public class TypedPublisher extends AbstractBehavior<String> { private final Topic<PubSubMsg> sharedTopic; public static Behavior<String> create() { return Behaviors.setup(TypedPublisher::new); } private TypedPublisher(ActorContext<String> context) { super(context); ActorSystem<?> system = context.getSystem(); // 创建与Classic端同名的主题 sharedTopic = Topic.create(system, "cluster-shared-topic", PubSubMsg.class); } @Override public Receive<String> createReceive() { return newReceiveBuilder() .onMessage(String.class, this::publishMessage) .build(); } private Behavior<String> publishMessage(String msg) { // 发布消息到共享主题 sharedTopic.ref().tell(new PubSubMsg.TypedMsg(msg)); getContext().getLog().info("Typed节点发送消息: {}", msg); return this; } // 启动Typed节点并发布测试消息 public static void main(String[] args) { ActorSystem<String> system = ActorSystem.create(TypedPublisher.create(), "akka-cluster"); system.tell("来自Typed节点的问候"); } }
关键注意事项
- 主题名称一致性:Classic与Typed端必须使用完全相同的主题名称,否则无法互通。
- 消息序列化兼容:自定义消息类型需在两端配置相同的序列化规则(如Jackson、Protobuf),避免反序列化失败。
- 集群配置统一:所有节点的种子节点、角色、共存模式配置需保持一致,确保集群拓扑正常。
内容的提问来源于stack exchange,提问作者skinnydan72
相关产品推荐
相关产品推荐

