如何基于自定义条件实现Vertx消息重试的弹性机制?
Vertx主备实例消息重试与有序队列实现方案
需求背景
我需要在Vertx应用中实现弹性机制:场景为主、备两个独立应用实例,主实例需向备实例发送消息,当指定布尔条件未满足时持续重试(比如主备网络断连时反复发送同一条消息)。单条消息的核心逻辑伪代码如下:
while(condition not met) { boolean messageSendingSuccess = tryToSendMessage(secondaryAddress); if(messageSendingSuccess) return; }
实际业务需要构建有序消息队列,要求队列必须等第一条消息发送成功后,再处理下一条新消息。已知CircuitBreaker无法满足需求(仅支持指定次数重试),我不需要按次数或时间停止重试,而是直到自定义条件满足才停止。
曾看到某博客的EventBus重试方案,但它仅适用于应用内通信,无法指定跨网络的目标地址,相关代码参考:
public class RetryWithHandlerVerticle extends AbstractVerticle { @Override public void start() throws Exception { vertx.eventBus() .send("hello.handler.failure.retry", "Hello", new Handler<AsyncResult<Message<String>>>() { private int count = 1; @Override public void handle(final AsyncResult<Message<String>> aResult) { if (aResult.succeeded()) { System.out.printf("received: \"%s\"\n", aResult.result().body()); } else if (count < 3) { System.out.printf("retry count %d, received error \"%s\"\n", count, aResult.cause().getMessage()); vertx.eventBus().send("hello.handler.failure.retry", "Hello", this); count = count + 1; } else { aResult.cause().printStackTrace(); } } }); } }
实现方案
核心思路
基于Vertx异步非阻塞特性,采用递归异步重试+有序队列的组合方式:
- 用阻塞队列维护待发送消息,保证顺序性;
- 用异步递归调用实现持续重试,替代同步
while循环避免阻塞事件循环; - 通过
NetClient/NetServer实现跨实例网络通信,替代仅支持应用内的EventBus。
具体代码实现
1. 主实例:消息队列+持续重试逻辑
import io.vertx.core.AbstractVerticle; import io.vertx.core.AsyncResult; import io.vertx.core.Future; import io.vertx.core.Handler; import io.vertx.core.Vertx; import io.vertx.core.buffer.Buffer; import io.vertx.core.net.NetClient; import io.vertx.core.net.NetClientOptions; import java.util.concurrent.LinkedBlockingQueue; public class PrimaryVerticle extends AbstractVerticle { // 备实例网络地址配置 private static final String SECONDARY_HOST = "备实例IP"; private static final int SECONDARY_PORT = 8080; // 有序消息队列 private final LinkedBlockingQueue<String> messageQueue = new LinkedBlockingQueue<>(); // 跨实例通信客户端 private NetClient netClient; // 防止并发处理的标记 private volatile boolean isProcessing = false; @Override public void start() { // 初始化NetClient,配置连接超时 NetClientOptions clientOpts = new NetClientOptions().setConnectTimeout(5000); netClient = vertx.createNetClient(clientOpts); // 启动队列消费循环 processNextMessage(); } // 对外提供添加消息到队列的方法 public void addMessage(String message) { try { messageQueue.put(message); if (!isProcessing) { processNextMessage(); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } // 处理队列中的下一条消息 private void processNextMessage() { if (isProcessing) return; String message = messageQueue.poll(); if (message == null) { isProcessing = false; return; } isProcessing = true; // 开始持续重试发送当前消息 retrySendMessage(message); } // 持续重试发送消息,直到条件满足 private void retrySendMessage(String message) { sendToSecondary(message, ar -> { if (ar.succeeded()) { // 发送成功,继续处理下一条 isProcessing = false; processNextMessage(); } else { // 检查是否需要停止重试 if (!shouldStopRetry()) { // 延迟1秒后重试(可根据业务调整间隔) vertx.setTimer(1000, timerId -> retrySendMessage(message)); } else { // 满足停止条件,处理失败逻辑 isProcessing = false; System.err.println("停止重试,消息发送失败: " + message); // 可选:将消息重新放回队列尾部 messageQueue.offer(message); } } }); } // 向备实例发送消息的具体实现 private void sendToSecondary(String message, Handler<AsyncResult<Void>> resultHandler) { netClient.connect(SECONDARY_PORT, SECONDARY_HOST, connectAr -> { if (connectAr.succeeded()) { connectAr.result().write(Buffer.buffer(message), writeAr -> { if (writeAr.succeeded()) { resultHandler.handle(Future.succeededFuture()); } else { resultHandler.handle(Future.failedFuture(writeAr.cause())); } connectAr.result().close(); }); } else { resultHandler.handle(Future.failedFuture(connectAr.cause())); } }); } // 自定义停止重试的布尔条件,根据业务需求实现 private boolean shouldStopRetry() { // 示例:主实例集群状态失效时停止重试 return vertx.isClustered() && !vertx.getClusterManager().isActive(); // 可替换为自定义逻辑,比如"备实例已下线"的判断 } }
2. 备实例:消息接收逻辑
import io.vertx.core.AbstractVerticle; import io.vertx.core.net.NetServer; public class SecondaryVerticle extends AbstractVerticle { private static final int LISTEN_PORT = 8080; @Override public void start() { vertx.createNetServer() .connectHandler(socket -> { socket.handler(buffer -> { String message = buffer.toString(); System.out.println("收到主实例消息: " + message); // 业务处理逻辑... }); }) .listen(LISTEN_PORT, listenAr -> { if (listenAr.succeeded()) { System.out.println("备实例监听端口启动成功: " + LISTEN_PORT); } else { System.err.println("备实例监听端口启动失败: " + listenAr.cause().getMessage()); } }); } }
关键说明
- 跨实例通信:用
NetClient/NetServer实现跨网络的主备通信,完全支持指定目标地址和端口; - 持续重试:通过递归调用+定时器实现异步重试,不会阻塞Vertx事件循环,重试间隔可自定义;
- 队列顺序保证:用
LinkedBlockingQueue维护消息顺序,isProcessing标记确保同一时间仅处理一条消息,严格遵循"先成功再处理下一条"的要求; - 停止条件灵活:
shouldStopRetry()方法完全由业务逻辑控制停止时机,不受重试次数或时间限制。
内容的提问来源于stack exchange,提问作者Zygmuntix
相关产品推荐
相关产品推荐

