Spring MVC中Flux心跳异步异常的解决方法咨询
问题:Spring Web MVC中关闭客户端连接后心跳发送触发IllegalStateException异常
我使用的是Spring Web MVC而非Spring Webflux,现有如下代码实现:
控制器代码
@RestController public class MyController { private final Flux<T> flux1; private final Flux<T> flux2; @GetMapping(produces = APPLICATION_NDJSON_VALUE) public Flux<T> subscribe() { var mergedFlux = Flux.merge(this.flux1, this.flux2).doOnCancel(() -> { log.debug("canceled"); }).doOnError(e -> { log.error("error", e); }); return keepAlive(Duration.of(10, ChronoUnit.SECONDS), mergedFlux); } private Flux<T> keepAlive(Duration interval, Flux<T> flux) { Flux<HeartBeat> heartBeat = Flux.interval(interval) .map(t -> new HeartBeat()).doFinally(signalType -> log.debug("HeartBeat closed")); return Flux.merge(flux, heartBeat); } }
配置类代码
@Configuration public class Config implements WebMvcConfigurer { @Override public void configureAsyncSupport(AsyncSupportConfigurer configurer) { configurer.setTaskExecutor(poolTaskExecutor()); } @Bean(name = "AsyncExecutor") public AsyncTaskExecutor poolTaskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(1); executor.setMaxPoolSize(2); executor.setThreadNamePrefix("AsyncExecutor-"); executor.initialize(); return executor; } }
Sink配置代码
@Configuration public class SinkConfig { @Bean public Sinks.Many<T> flux1() { return Sinks.many().replay().latest(); } @Bean public Flux<T> flux1(Sinks.Many<T> sink) { return sink.asFlux(); } @Bean public Sinks.Many<T> flux2() { return Sinks.many().replay().latest(); } @Bean public Flux<T> flux2(Sinks.Many<T> sink) { return sink.asFlux(); } }
问题描述
每当我关闭客户端连接,下一次心跳即将发送时,会出现如下异常:
Exception in thread "AsyncExecutor-1" java.lang.IllegalStateException: A non-container (application) thread attempted to use the AsyncContext after an error had occurred and the call to AsyncListener.onError() had returned. This is not allowed to avoid race conditions
该异常仅在使用keepAlive方法发送心跳以维持连接避免超时的情况下出现,请问如何解决这个问题?
解决方案
这个问题的核心是:客户端连接关闭后,心跳Flux仍在持续发送数据,此时尝试往已失效的AsyncContext写数据,就会触发该异常。需要确保连接关闭后,心跳Flux能及时终止,不再尝试发送数据。
可以从以下几个方面调整:
1. 让心跳Flux感知连接状态,及时终止
Spring MVC处理异步响应时,客户端断开连接会触发请求的取消信号。修改keepAlive方法,将心跳Flux与原业务Flux的生命周期绑定,确保原Flux终止(包括取消、错误、完成)时,心跳Flux也立即停止:
private Flux<T> keepAlive(Duration interval, Flux<T> flux) { // 原flux终止时,心跳flux同步停止 Flux<HeartBeat> heartBeat = Flux.interval(interval) .map(t -> new HeartBeat()) .takeUntilOther(flux.then()) .doFinally(signalType -> log.debug("HeartBeat closed")); // 确保合并后的flux能正确传播取消信号 return Flux.merge(flux, heartBeat) .doOnCancel(() -> log.debug("Merged flux canceled")); }
2. 配置Spring MVC异步请求的超时与取消处理
在WebMvcConfigurer配置中,添加异步请求超时设置,同时注册监听器在连接关闭时主动清理异步上下文:
@Override public void configureAsyncSupport(AsyncSupportConfigurer configurer) { configurer.setTaskExecutor(poolTaskExecutor()); // 设置异步请求超时时间,与心跳间隔匹配或稍长 configurer.setDefaultTimeout(15000); // 注册异步监听器,处理连接关闭、超时等场景 configurer.registerAsyncInterceptor(new AsyncHandlerInterceptor() { @Override public boolean preHandle(HttpServletRequest request, HttpServletResponse response, Object handler) throws Exception { AsyncContext asyncContext = request.startAsync(); asyncContext.addListener(new AsyncListener() { @Override public void onComplete(AsyncEvent event) throws IOException {} @Override public void onTimeout(AsyncEvent event) throws IOException { event.getAsyncContext().complete(); } @Override public void onError(AsyncEvent event) throws IOException { event.getAsyncContext().complete(); } @Override public void onStartAsync(AsyncEvent event) throws IOException {} }); return true; } }); }
3. 调整Sink类型,减少无效订阅残留
当前使用的Sinks.many().replay().latest()会保留最新数据并重放给新订阅者,容易产生无效订阅残留。换成multicast类型的Sink,更适合实时推送场景:
@Bean public Sinks.Many<T> flux1() { return Sinks.many().multicast().onBackpressureBuffer(); } @Bean public Flux<T> flux1(Sinks.Many<T> sink) { return sink.asFlux().doOnCancel(() -> log.debug("flux1 subscription canceled")); } // flux2配置同理
4. 优化Flux合并逻辑,确保取消信号传播
在控制器方法中,使用mergeWith替代手动Flux.merge,它能更可靠地传播取消信号到所有子Flux:
@GetMapping(produces = APPLICATION_NDJSON_VALUE) public Flux<T> subscribe() { return Flux.merge(flux1, flux2) .doOnCancel(() -> log.debug("canceled")) .doOnError(e -> log.error("error", e)) .mergeWith(keepAliveFlux()) .doOnCancel(() -> log.debug("Total flux canceled")); } private Flux<T> keepAliveFlux() { return Flux.interval(Duration.of(10, ChronoUnit.SECONDS)) .map(t -> (T) new HeartBeat()) .doFinally(signalType -> log.debug("HeartBeat closed")); }
核心原理总结
客户端断开连接后,Spring MVC的异步上下文会进入错误/完成状态,此时再写入数据就会触发异常。通过让心跳Flux监听业务Flux的终止信号,在连接关闭时立即停止发送,同时配置异步监听器及时清理上下文,就能避免该异常。
内容的提问来源于stack exchange,提问作者Bioaim
相关产品推荐
相关产品推荐

