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

如何让Vert.x Event Bus发布的消息仅被一个副本消费?

Vert.x集群模式下实现事件消费组机制(同一服务副本仅单实例消费,多服务类型可同时消费)

针对你的需求——同一服务的多个副本仅一个消费事件,同时允许不同类型服务共同消费该事件,以下是几种可行的实现方案:


方案1:利用Hazelcast Reliable Topic原生消费组特性

由于你已经集成了Hazelcast作为集群管理器,直接使用其原生的Reliable Topic消费组机制是最简洁可靠的方案。Hazelcast的Reliable Topic天然支持:

  • 同一消费组内,每个事件仅被一个实例处理
  • 不同消费组可独立消费同一事件,互不干扰

代码示例

package org.example;

import com.hazelcast.core.HazelcastInstance;
import com.hazelcast.core.ITopic;
import com.hazelcast.core.MessageListener;
import io.vertx.core.Vertx;
import io.vertx.core.VertxOptions;
import io.vertx.core.http.HttpServer;
import io.vertx.core.spi.cluster.ClusterManager;
import io.vertx.spi.cluster.hazelcast.HazelcastClusterManager;

public class SlaveNode {

    public static final String TOPIC_NAME = "my.topic";
    // 同一服务副本使用相同的消费组名称
    public static final String CONSUMER_GROUP = "service-slave-group"; 

    public static void main(String[] args) {
        ClusterManager mgr = new HazelcastClusterManager();
        VertxOptions options = new VertxOptions().setClusterManager(mgr);

        Vertx.clusteredVertx(options, res -> {
            if (res.succeeded()) {
                Vertx vertx = res.result();
                HazelcastInstance hazelcastInstance = ((HazelcastClusterManager) mgr).getHazelcastInstance();

                // 获取或创建Reliable Topic
                ITopic<String> topic = hazelcastInstance.getReliableTopic(TOPIC_NAME);

                // 订阅指定消费组:同一组内仅一个实例收到事件
                topic.addMessageListener(message -> {
                    // 切换到Vert.x上下文处理事件(保证线程模型一致性)
                    vertx.runOnContext(v -> {
                        System.out.println("receive event: " + message.getMessageObject());
                    });
                }, CONSUMER_GROUP);

                HttpServer server = vertx.createHttpServer();
                server.requestHandler(event -> {
                    String eventBody = "publish: " + event.absoluteURI();
                    // 发布事件到Topic
                    topic.publish(eventBody);
                    event.response().end("success");
                });
                server.listen(18081);

            } else {
                res.cause().printStackTrace();
            }
        });
    }
}

扩展说明

如果有其他类型的服务需要消费该事件,只需为其配置不同的消费组名称即可,Hazelcast会自动将事件推送给所有消费组的实例,且每个消费组内仅一个实例处理事件。


方案2:Vert.x Event Bus + 分布式锁实现自定义消费组

若不想直接操作Hazelcast API,可基于Vert.x Event Bus的广播机制,结合分布式锁实现消费组逻辑:

  1. 同一服务的所有副本使用相同的消费组标识
  2. 消费者收到事件后,先尝试获取以「事件唯一标识+消费组名」为key的分布式锁
  3. 成功获取锁的实例处理事件,其他实例直接忽略
  4. 处理完成后释放锁

代码示例

package org.example;

import io.vertx.core.Vertx;
import io.vertx.core.VertxOptions;
import io.vertx.core.eventbus.EventBus;
import io.vertx.core.http.HttpServer;
import io.vertx.core.shareddata.Lock;
import io.vertx.core.spi.cluster.ClusterManager;
import io.vertx.spi.cluster.hazelcast.HazelcastClusterManager;

import java.util.UUID;

public class SlaveNode {

    public static final String MY_ADDRESS = "my.address";
    // 同一服务副本使用相同的消费组名称
    public static final String CONSUMER_GROUP = "service-slave-group"; 

    public static void main(String[] args) {
        ClusterManager mgr = new HazelcastClusterManager();
        VertxOptions options = new VertxOptions().setClusterManager(mgr);

        Vertx.clusteredVertx(options, res -> {
            if (res.succeeded()) {
                Vertx vertx = res.result();
                EventBus eventBus = vertx.eventBus();

                // 订阅Event Bus事件
                eventBus.consumer(MY_ADDRESS, event -> {
                    String eventBody = (String) event.body();
                    // 生成事件唯一标识(可根据事件内容哈希或UUID生成)
                    String eventId = UUID.randomUUID().toString();
                    String lockKey = String.format("%s-%s-%s", MY_ADDRESS, CONSUMER_GROUP, eventId);

                    // 获取分布式锁
                    vertx.sharedData().getLockWithTimeout(lockKey, 5000, lockRes -> {
                        if (lockRes.succeeded()) {
                            Lock lock = lockRes.result();
                            // 仅拿到锁的实例处理事件
                            System.out.println("receive event: " + eventBody);
                            // 处理完成后释放锁
                            lock.release();
                        }
                    });
                });

                HttpServer server = vertx.createHttpServer();
                server.requestHandler(event -> {
                    vertx.eventBus().publish(MY_ADDRESS, "publish: " + event.absoluteURI());
                    event.response().end("success");
                });
                server.listen(18081);

            } else {
                res.cause().printStackTrace();
            }
        });
    }
}

注意事项

  • 事件唯一标识需保证每个事件的唯一性,避免不同事件抢占同一锁
  • 使用getLockWithTimeout设置锁超时时间,防止实例崩溃导致锁无法释放

方案3:多地址路由+Event Bus send

若服务类型固定,可给每种服务分配独立的Event Bus地址,发布事件时对每个地址调用send方法:

  • 服务SlaveNode的地址为my.address.slave,用send发送(仅一个副本消费)
  • 服务ServiceB的地址为my.address.serviceb,用send发送(仅一个副本消费)

这种方案无需额外依赖,但扩展性较差,适合服务类型固定的场景。


方案推荐

优先选择方案1,Hazelcast的Reliable Topic原生支持消费组逻辑,无需手动实现锁机制,性能和可靠性更优,完全匹配你的需求。

内容的提问来源于stack exchange,提问作者张Del

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 23:03:10