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

Spring Cloud Gateway v2过滤器Mono完成与MQ重试实现咨询

你的现有方案存在明确的设计缺陷,核心问题有两个:

  • ScheduledFuture.get()是阻塞调用,即便包装在Mono.fromRunnable里,默认调度下会直接占用Netty事件循环线程,高并发场景下会触发线程饥饿,导致网关吞吐量骤降甚至假死;即使手动将阻塞逻辑调度到弹性线程池,也存在不必要的线程上下文切换开销,不符合响应式编程的非阻塞原则。
  • 手动维护全局映射表、调度超时任务、管理Future状态的逻辑冗余度极高,很容易因异常分支漏写清理逻辑导致内存泄漏,并发场景下的状态一致性也需要额外处理。

这类MQ请求-响应+超时+重试的逻辑,完全可以基于Reactor原生操作符实现全链路非阻塞,不需要自己管理线程池和任务状态,具体实现方案如下:

核心设计思路

  • 基于唯一请求关联ID(correlationId)做请求和响应的匹配,MQ生产端发消息时携带该ID,消费端回传响应时原样带回该ID。
  • 用Mono.create()桥接MQ的异步响应逻辑,替代自定义的Runnable、Future和全局映射逻辑,通过sink的回调完成响应通知。
  • 超时、重试逻辑直接使用Reactor内置的timeout、retryWhen操作符实现,无需手动调度定时任务。
  • 所有资源清理逻辑绑定到流的终止事件,确保请求正常完成、异常、超时、客户端断连、重试等所有场景下,监听器都能被自动移除,避免内存泄漏。

具体代码实现

1. MQ响应监听器

用并发安全的Map存储请求ID和对应的响应回调,MQ客户端收到响应消息后直接触发对应回调即可,不需要额外的状态管理:

@Component
public class MqResponseListener {
    // 存储请求ID和对应的响应回调
    private static final ConcurrentHashMap<String, Consumer<MqResponse>> CALLBACK_MAP = new ConcurrentHashMap<>();

    /**
     * 注册请求响应回调
     */
    public static void addListener(String correlationId, Consumer<MqResponse> callback) {
        CALLBACK_MAP.put(correlationId, callback);
    }

    /**
     * 移除请求响应回调
     */
    public static void removeListener(String correlationId) {
        CALLBACK_MAP.remove(correlationId);
    }

    /**
     * MQ消费端监听响应队列,这里以RabbitMQ为例,其他MQ替换对应监听器注解即可
     */
    @RabbitListener(queues = "${mq.config.response-queue}")
    public void handleMqResponse(MqResponse response) {
        // 匹配到对应请求的回调后直接触发,同时移除回调避免重复执行
        Consumer<MqResponse> callback = CALLBACK_MAP.remove(response.getCorrelationId());
        if (callback != null) {
            callback.accept(response);
        }
    }
}

2. Gateway过滤器核心逻辑

全程无阻塞调用,所有异步逻辑通过Reactor操作符编排:

@Override
public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) {
    ServerHttpResponse response = exchange.getResponse();
    ServerHttpRequest request = exchange.getRequest();

    return Mono.<MqResponse>create(sink -> {
        // 每轮请求生成唯一关联ID
        String correlationId = UUID.randomUUID().toString();
        // 注册当前请求的响应回调
        MqResponseListener.addListener(correlationId, sink::success);
        // 构造并发送MQ消息,发送失败直接向下游传递异常
        mqProducer.send(buildRequestMessage(correlationId, request))
                .doOnError(sink::error)
                .subscribe();
        // 流无论因为什么原因终止,都自动清理当前请求的回调
        sink.onDispose(() -> MqResponseListener.removeListener(correlationId));
    })
    // 单轮请求超时配置:5秒未收到响应直接抛出超时异常
    .timeout(Duration.ofSeconds(5))
    // 重试配置:最多重试2次,每次间隔1秒,仅对超时、MQ发送失败类可重试异常生效
    .retryWhen(Retry.fixedDelay(2, Duration.ofSeconds(1))
            .filter(ex -> ex instanceof TimeoutException || ex instanceof MqSendException)
    )
    // 收到正常响应后,将结果写入HTTP返回给客户端
    .flatMap(mqResponse -> {
        response.setStatusCode(HttpStatus.OK);
        response.getHeaders().setContentType(MediaType.APPLICATION_JSON);
        DataBuffer respBuffer = response.bufferFactory().wrap(mqResponse.getRespBodyBytes());
        return response.writeWith(Mono.just(respBuffer));
    })
    // 重试耗尽或遇到不可重试异常时,统一返回500状态码
    .onErrorResume(ex -> {
        response.setStatusCode(HttpStatus.INTERNAL_SERVER_ERROR);
        return response.setComplete();
    });
}

方案优势

  • 全链路无阻塞,完全适配Spring Cloud Gateway的响应式运行模型,不会占用Netty事件循环线程,性能表现远高于阻塞式实现。
  • 不需要手动管理线程池、定时任务、Future状态,超时、重试、资源清理全部由Reactor原生能力实现,代码量比原有方案减少60%以上,异常分支覆盖更全。
  • 重试触发时会自动执行完整的发消息、注册监听器流程,上一轮请求的监听器会被自动清理,不需要手动编写状态重置逻辑。
  • 不存在内存泄漏风险:所有回调的移除逻辑绑定到流的onDispose事件,覆盖请求正常返回、超时、客户端主动断连、重试、异常等所有终止场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 12:57:14