如何在非阻塞场景下记录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
相关产品推荐
相关产品推荐

