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
相关产品推荐
相关产品推荐

