如何让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的广播机制,结合分布式锁实现消费组逻辑:
- 同一服务的所有副本使用相同的消费组标识
- 消费者收到事件后,先尝试获取以「事件唯一标识+消费组名」为key的分布式锁
- 成功获取锁的实例处理事件,其他实例直接忽略
- 处理完成后释放锁
代码示例
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
相关产品推荐
相关产品推荐

