You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.16 12:27:47