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

如何基于自定义条件实现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异步非阻塞特性,采用递归异步重试+有序队列的组合方式:

  1. 用阻塞队列维护待发送消息,保证顺序性;
  2. 用异步递归调用实现持续重试,替代同步while循环避免阻塞事件循环;
  3. 通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 22:25:30