使用Mono.fromCallable()仍遇Reactor线程不支持阻塞错误的原因及解决
你遇到的核心问题其实不是findTenantStorageFactory里的实现,而是调用它的地方在Reactor的IO线程中执行了block()——这是Reactor明确禁止的操作,因为Netty的IO线程(也就是错误里的reactor-http-nio-4)是整个应用的性能核心,阻塞它会直接拖垮系统的并发能力。
为什么你的代码会报错?
看一下blobStorageClientBuilder的get()方法:
private val blobStorageClientBuilder: AzureBlobClientBuilder get() = findTenantStorageFactory(tenantId).block()!! .fold({ throw it }, { it.blobStorageClient })!!
这里直接调用了block(),而这个get()方法是在IO线程中被触发的(比如处理HTTP请求的时候)。哪怕你在findTenantStorageFactory里用了boundedElastic(),也改变不了调用block()的线程是IO线程这个事实——Reactor会严格检查这一点,直接抛出异常。
另外,findTenantStorageFactory里还有一个问题:你把tenantId这个Mono直接block()了,这完全违背了响应式编程的原则,Mono应该用链式操作来处理,而不是强制阻塞获取值。
正确的解决方案
我们需要把整个调用链改成全响应式,彻底移除所有不必要的block()调用:
1. 重构findTenantStorageFactory方法
不要阻塞tenantId,而是用flatMap来处理它,同时把阻塞的lookupTenantBlobStorageFactory放到boundedElastic线程池执行:
fun findTenantStorageFactory(tenantId: Mono<TenantId>): Mono<Either<BaseAzureBlobStorageException, MultitenantAzureBlobStorageFactory>> { return tenantId.flatMap { id -> // 把阻塞的lookup操作包装到Mono.fromCallable,并指定boundedElastic调度器 Mono.fromCallable { lookupTenantBlobStorageFactory(id, factories) }.subscribeOn(Schedulers.boundedElastic()) } }
这样我们就完全避免了在方法内部阻塞tenantId,而是用响应式的方式处理它。
2. 修改blobStorageClientBuilder的获取方式
不要再用val加get()的方式(因为这会强制同步返回值),而是把它改成返回Mono<AzureBlobClientBuilder>的方法:
private fun getBlobStorageClientBuilder(): Mono<AzureBlobClientBuilder> { return findTenantStorageFactory(tenantId) .map { either -> either.fold( { exception -> throw exception }, // 异常直接抛出,会被Reactor处理 { factory -> factory.blobStorageClient } ) } }
3. 在业务代码中响应式地使用它
在需要用到AzureBlobClientBuilder的地方,继续用响应式操作符链式调用,比如在WebFlux的Controller里:
@GetMapping("/some-endpoint") fun getSomeData(): Mono<ResponseEntity<Data>> { return getBlobStorageClientBuilder() .flatMap { builder -> // 这里用builder执行你的Blob操作,同样要保持响应式 Mono.fromCallable { builder.buildClient().getBlobContainerClient("container").getBlobClient("blob").downloadContent() }.subscribeOn(Schedulers.boundedElastic()) } .map { content -> ResponseEntity.ok(content) } .onErrorResume { e -> ResponseEntity.internalServerError().body(e.message) } }
关键注意事项
- 永远不要在IO线程中调用阻塞方法:Reactor的IO线程(
reactor-http-nio-*)是用来处理非阻塞IO的,任何阻塞操作都必须放到boundedElastic或者自定义的阻塞线程池中。 - 保持调用链全响应式:一旦进入Reactor的响应式世界,就尽量用
flatMap、map、filter等操作符来串联逻辑,不要中途用block()打断链式调用。 - 谨慎处理异常:响应式中的异常可以用
onErrorResume、onErrorReturn等操作符来处理,避免直接抛出未捕获的异常。
内容的提问来源于stack exchange,提问作者A Bit of Help

