如何避免Reactor中调用block()时出现IllegalStateException
问题描述
我们有一个同步流程,需要调用两个REST接口,且第二个接口依赖第一个的返回结果。使用Spring WebClient时,调用.block()会抛出如下异常:
java.lang.IllegalStateException: block()/blockFirst()/blockLast() are blocking, which is not supported in thread parallel-2
该如何避免这个问题?
简化代码片段
var job = webClient.createJob().block(); if (job == null || StringUtils.isBlank(job.getId())) { throw new Exception("WebClient did not return with a job id"); } batchRecords(job.getId(), records);// 该方法内部也会调用WebClient
这段代码在单元测试中正常运行,但通过@RestController调用时会抛出上述异常。
补充说明
batchRecords方法中同样存在阻塞式Mono调用,还包含延迟逻辑:
public void batchRecords(final String jobId, final List<InventoryRecord> records) { var recordCount = 0; var inventoryPositions = new ArrayList<InventoryPosition>(); var recordIterator = records.iterator(); while (recordIterator != null && recordIterator.hasNext()) { var inventoryRecord = recordIterator.next(); inventoryPositions.add(mapInventoryPosition(inventoryRecord)); recordCount++; if (inventoryPositions.size() == batchSize) { var response = createBatch(jobId, inventoryPositions); Thread.sleep(sleepTime); response.block(); inventoryPositions = new ArrayList<>(); } } }
解决方案
出现这个异常的核心原因是:你在**Reactor的非阻塞线程(parallel-2属于Reactor并行调度器线程)**上调用了阻塞方法block(),而Reactor的设计原则是禁止在其非阻塞线程上执行阻塞操作,否则会破坏响应式流的非阻塞特性,甚至导致线程池耗尽。
正确做法:全程使用响应式编程,避免阻塞调用
需要把整个流程改成纯响应式的,用Mono/Flux的链式调用替代所有block()和Thread.sleep():
- 改造主流程:将
block()替换为响应式链式调用,控制器方法返回Mono<Void>
// 控制器方法返回Mono<Void>,而非void public Mono<Void> processRecords(List<InventoryRecord> records) { return webClient.createJob() .flatMap(job -> { if (job == null || StringUtils.isBlank(job.getId())) { return Mono.error(new Exception("WebClient did not return with a job id")); } // 调用改造后的batchRecords方法,返回Mono<Void> return batchRecords(job.getId(), records); }); }
- 改造batchRecords方法:用
Mono.delay()替代Thread.sleep(),用flatMap替代block(),全程保持响应式
public Mono<Void> batchRecords(final String jobId, final List<InventoryRecord> records) { // 将列表转为Flux,逐个处理并按批次分组 return Flux.fromIterable(records) .map(this::mapInventoryPosition) .buffer(batchSize) .flatMap(batch -> // 先调用createBatch,再延迟指定时间,最后忽略响应返回Mono<Void> createBatch(jobId, batch) .delayElement(Duration.ofMillis(sleepTime)) .then() ) .then(); // 所有批次处理完成后返回Mono<Void> }
为什么单元测试正常?
单元测试中,WebClient的调用通常是同步模拟的,线程环境是普通的JUnit测试线程,不属于Reactor的非阻塞调度器线程,所以block()不会触发异常。但在@RestController中,请求处理线程是Reactor的非阻塞线程,此时调用block()就会触发禁止阻塞的检查。
额外注意点
- 控制器方法必须返回
Mono或Flux类型,Spring Web会自动处理响应式流的订阅和执行 - 避免在响应式流程中使用任何阻塞操作(
block()、Thread.sleep()、同步IO等),所有操作优先转为响应式API - 如果必须调用阻塞代码(比如第三方同步API),可以使用
Mono.fromCallable()配合subscribeOn(Schedulers.boundedElastic())将阻塞操作隔离到专用线程池,但这是最后的妥协方案
内容的提问来源于stack exchange,提问作者Thomas
相关产品推荐
相关产品推荐

