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

如何解决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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 02:34:51