使用DataBufferUtils.write传至Azure Blob超250MB卡顿,Spring WebFlux代码是否正确?
问题分析与解决方案
核心问题
你的代码存在阻塞式IO与响应式代码混用的矛盾,这是大文件传输卡住的关键原因:
getBlobOutputStream()返回的是传统阻塞式输出流,所有写入操作都会阻塞当前线程DataBufferUtils.write()是异步响应式操作,直接将异步数据流写入阻塞流时,当数据量超过一定阈值(比如260MB),Reactor的背压机制无法正常协调异步生产与阻塞消费的速率,最终导致数据流停滞。
修复方案
方案1:使用Azure Blob Storage响应式SDK(推荐)
Azure Blob SDK提供了原生响应式API,直接支持Flux<DataBuffer>的写入,完全适配Reactor异步模型,无需手动处理背压:
@GetMapping(path = "/trigger-download-to-blob/{fileSizeInMb}") public Mono<Void> triggerDownloadToBlob(@PathVariable int fileSizeInMb) { log.info("triggerDownload"); Flux<DataBuffer> flux = this.webClient .get() .uri("/serve-file/" + fileSizeInMb) .accept(MediaType.APPLICATION_OCTET_STREAM) .exchangeToFlux(clientResponse -> clientResponse.body(BodyExtractors.toDataBuffers())); String destination = "TestDownloadToAzureBlobStorage" + System.currentTimeMillis() + ".pdf"; BlobClient blobClientTarget = this.containerClient.getBlobClient(destination); // 直接用响应式upload方法传入数据流,SDK自动处理分块与背压 return blobClientTarget.getBlockBlobClient() .uploadWithResponse(flux, null, null, null, null, null, null) .doOnSuccess(response -> log.info("!!!!!!!!!!!!!!!!!!!!!!!!!!!! end download of {}", destination)) .then(); }
方案2:修复原有代码的背压与线程模型
如果必须保留阻塞式OutputStream,需调整线程模型并确保背压处理正常:
@GetMapping(path = "/trigger-download-to-blob/{fileSizeInMb}") public void triggerDownloadToBlob(@PathVariable int fileSizeInMb) { log.info("triggerDownload"); Flux<DataBuffer> flux = this.webClient .get() .uri("/serve-file/" + fileSizeInMb) .accept(MediaType.APPLICATION_OCTET_STREAM) .exchangeToFlux(clientResponse -> clientResponse.body(BodyExtractors.toDataBuffers())); String destination = "TestDownloadToAzureBlobStorage" + System.currentTimeMillis() + ".pdf"; BlobClient blobClientTarget = this.containerClient.getBlobClient(destination); try (OutputStream outputStream = blobClientTarget.getBlockBlobClient().getBlobOutputStream(this.parallelTransferOptions, null, null, null, null)) { // 将阻塞IO操作放到专门的弹性线程池,避免阻塞Reactor NIO主线程 DataBufferUtils.write(flux.subscribeOn(Schedulers.boundedElastic()), outputStream) .map(DataBufferUtils::release) .blockLast(Duration.ofHours(22)); outputStream.flush(); } catch (IOException e) { throw new IllegalStateException(e); } log.info("!!!!!!!!!!!!!!!!!!!!!!!!!!!! end download of {}", destination); }
关键说明
- 方案1是最优解,完全遵循响应式编程范式,Azure SDK内部会自动处理分块上传、背压协调和资源释放,彻底避免阻塞与异步混用的风险
- 方案2中,
Schedulers.boundedElastic()用于隔离阻塞IO操作,防止Reactor核心线程被占用;默认的背压缓冲机制会确保数据流不会因阻塞消费而中断
内容的提问来源于stack exchange,提问作者user286974
相关产品推荐
相关产品推荐

