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

如何在非阻塞场景下记录bodyToFlux返回的Flux元素数量?

记录Flux返回元素数量的可行方案

你之前注释掉count().subscribe()的核心问题其实是冷流重复请求——原Flux是冷流,每次调用subscribe都会重新发起接口请求,而不是阻塞当前方法。下面两种方式可以解决这个问题,既记录数量又不影响原Flux的正常使用:

方案一:用share()共享流避免重复请求

public Flux<DdocumentBackofficeResponseRestDTO> queryPendingDocumentsToBeHashed() {
    Flux<DdocumentBackofficeResponseRestDTO> documents = webClient.get().uri("/pending-documents")
            .headers(h -> h.setBearerAuth(bearerToken)).retrieve()
            .onStatus(httpStatus -> !httpStatus.is2xxSuccessful(),
                    clientResponse -> handleErrorResponse(clientResponse))
            .bodyToFlux(DdocumentBackofficeResponseRestDTO.class)
            .share(); // 将冷流转为热流,多个订阅者共享同一份数据

    documents.count().subscribe(c -> logger.info("Will request a total of {} files", c));
    return documents;
}

share()会让原Flux变成热流:第一次订阅(这里的count().subscribe())会触发接口请求,后续调用方的订阅会复用已有的数据流,不会重复调用接口,完美解决重复请求的问题。

方案二:用计数器+生命周期操作符(无需额外订阅)

public Flux<DdocumentBackofficeResponseRestDTO> queryPendingDocumentsToBeHashed() {
    AtomicInteger counter = new AtomicInteger(0);
    return webClient.get().uri("/pending-documents")
            .headers(h -> h.setBearerAuth(bearerToken)).retrieve()
            .onStatus(httpStatus -> !httpStatus.is2xxSuccessful(),
                    clientResponse -> handleErrorResponse(clientResponse))
            .bodyToFlux(DdocumentBackofficeResponseRestDTO.class)
            .doOnNext(dto -> counter.incrementAndGet()) // 每处理一个元素就计数+1
            .doOnComplete(() -> logger.info("Will request a total of {} files", counter.get())); // 流完成时打印总数
}

这种方式不需要单独订阅count(),而是在流的处理生命周期中维护计数器:

  • doOnNext在每个元素通过时递增计数器
  • doOnComplete在所有元素处理完成后打印总数
    如果需要处理异常场景,可以额外添加doOnError来重置计数器或记录异常,避免计数不准确。

内容的提问来源于stack exchange,提问作者icordoba

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 04:52:39