如何解决Spring WebFlux异步接口超时异常并实现定时批量POST请求
问题
需求:编写异步API,接收请求后返回202 Accepted,随后在halt-time秒内每秒发送n个POST请求。
现有实现代码:
@PostMapping(value = "/stress") @ResponseStatus(HttpStatus.ACCEPTED) private Mono<?> stressPostLog(@RequestBody StressTest request) throws InterruptedException { return Mono.fromCallable(() -> { sendBatchRequest(request); return null; }).subscribeOn(Schedulers.boundedElastic()); } private void sendBatchRequest(StressTest requestParams) throws InterruptedException { successfulRequest.set(0); long haltTimeMillis = requestParams.getHalt_time() * 1000; log.info("Halt Time Millis = " + haltTimeMillis ); Integer count = 0; long initTime = System.currentTimeMillis(); long currTime; do { Thread.sleep(1000); //! Sleep time is subjected to system for(int i = 0; i < requestParams.getRequest_per_second(); i++) { AuditSchema request = AuditSchema.generateRandomValue(requestParams.getRequest_size()); Mono<?> response = postAndIncrement(request); response.publishOn(Schedulers.boundedElastic()).subscribe(); count++; } currTime = System.currentTimeMillis(); log.info("Time duration = " + ( currTime - initTime)); }while((currTime - initTime) < haltTimeMillis); log.info("Request Send = {} and Response Succeed = {}", count, successfulRequest); } private Mono<?> postAndIncrement(AuditSchema request) { return this.client.post() .uri("v1/log") .body(BodyInserters.fromValue(request)) .accept(MediaType.APPLICATION_JSON) .exchangeToMono(clientResponse -> { log.info("Request suceeded {}", clientResponse.statusCode()); successfulRequest.getAndIncrement(); return Mono.empty(); }); }
出现的问题:请求未异步执行,且5秒后触发超时错误:
java.lang.IllegalStateException: Timeout on blocking read for 5000000000 NANOSECONDS
请问如何使用Spring WebFlux正确实现该功能?
解决方案
问题根源
- 控制器方法使用
private修饰,Spring无法正确代理处理请求,导致异步逻辑失效。 sendBatchRequest中的Thread.sleep是阻塞操作,会占用线程池资源,拖慢响应返回,触发超时。- 手动调用
subscribe()缺乏统一调度和错误处理,不符合WebFlux响应式编程规范。
正确实现代码
import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.core.scheduler.Schedulers; import java.time.Duration; @PostMapping(value = "/stress") @ResponseStatus(HttpStatus.ACCEPTED) public Mono<Void> stressPostLog(@RequestBody StressTest request) { // 立即返回202响应,后台异步执行压测逻辑 Mono.fromRunnable(() -> startStressTest(request)) .subscribeOn(Schedulers.boundedElastic()) .subscribe(); return Mono.empty(); } private void startStressTest(StressTest requestParams) { successfulRequest.set(0); int totalSeconds = requestParams.getHalt_time(); int requestsPerSecond = requestParams.getRequest_per_second(); // 非阻塞定时任务,每秒执行一次,持续totalSeconds秒 Flux.interval(Duration.ofSeconds(1)) .take(totalSeconds) .flatMap(second -> // 每秒生成指定数量的请求,批量异步发送 Flux.range(0, requestsPerSecond) .flatMap(i -> { AuditSchema auditRequest = AuditSchema.generateRandomValue(requestParams.getRequest_size()); return postAndIncrement(auditRequest); }) ) .doOnComplete(() -> log.info("Total Requests Sent = {}, Successful Responses = {}", totalSeconds * requestsPerSecond, successfulRequest.get()) ) .subscribeOn(Schedulers.boundedElastic()) .subscribe( null, error -> log.error("Stress test encountered error", error) ); } private Mono<Void> postAndIncrement(AuditSchema request) { return this.client.post() .uri("v1/log") .body(BodyInserters.fromValue(request)) .accept(MediaType.APPLICATION_JSON) .exchangeToMono(clientResponse -> { if (clientResponse.statusCode().is2xxSuccessful()) { successfulRequest.getAndIncrement(); } log.info("Request status: {}", clientResponse.statusCode()); return Mono.empty(); }) .onErrorResume(error -> { log.error("Request failed", error); return Mono.empty(); }); }
关键改进点
- 控制器方法改为
public,确保Spring能正确处理请求。 - 用
Mono.fromRunnable包装后台任务,立即返回Mono.empty()保证202响应快速返回。 - 替换
Thread.sleep为Flux.interval,实现非阻塞定时调度,避免线程资源浪费。 - 用
flatMap批量处理请求,统一管理异步任务的调度和错误处理,避免手动subscribe()的弊端。 - 增加全局错误处理逻辑,防止单个请求异常导致整个压测任务中断。
- 保留原子变量
successfulRequest确保计数的线程安全性。
内容的提问来源于stack exchange,提问作者denisssj
相关产品推荐
相关产品推荐

