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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.08 13:52:38