Spring WebFlux:遍历Flux直至获取非空列表的实现问题
问题:WebFlux下遍历ID获取非空结果的非阻塞实现
现有代码与返回示例
获取ID列表的方法
public Flux<IdEntrevista> getIdEntrevista(String perfil){ return this.webClient.baseUrl("http://localhost:8081").build() .get() .uri("/api/orquestador/v1/entrevistador/public/entrevista_muestra_id?perfil="+perfil) .retrieve() .bodyToFlux(IdEntrevista.class); }
返回示例:
[ { "id": "6639711af44f9905dbdcd889" }, { "id": "663971fdd44ffdd5dbdcd88b" }, ...(29 elements more) ]
获取问题的仓库方法调用逻辑
public Flux<SoloPreguntaImp> getPreguntas(String perfil, int limit) { // 待实现逻辑:从ID列表中获取有效ID,调用仓库方法 return interviewTestDao.getPreguntas("6639711af44f9905dbdcd889",3); }
仓库方法返回示例(有效结果):
[ { "pregunta": "¿Cuál es tu experiencia trabajando con tecnologías de Java y React en proyectos de desarrollo de software?" }, { "pregunta": "¿Has liderado el diseño e implementación de aplicaciones de software complejas utilizando JavaScript, ReactJS y Java?" }, { "pregunta": "¿Cómo garantizas la calidad del código, la mantenibilidad y escalabilidad de las aplicaciones de software en las que trabajas?" } ]
需求
遍历getIdEntrevista返回的Flux中的每个ID,依次传入仓库方法interviewTestDao.getPreguntas,直到获取到非空的问题列表,全程遵循非阻塞原则。
错误尝试代码
public Flux<SoloPreguntaImp> getPreguntas(String perfil, int limit) { return getIdEntrevista(perfil) .flatMap(idEntrevista -> interviewTestDao.getPreguntas(idEntrevista.getId(), limit) .collectList() .flatMapMany(list -> { if (list.isEmpty()) { return Mono.error(new RuntimeException("Lista vacía")); } else { return Flux.fromIterable(list); } })) .repeatWhenEmpty(repeat -> repeat.delayElements(Duration.ofSeconds(5))) }
解决方案
public Flux<SoloPreguntaImp> getPreguntas(String perfil, int limit) { return getIdEntrevista(perfil) // 按顺序逐个处理ID,保证遍历逻辑的顺序性 .concatMap(idEntrevista -> interviewTestDao.getPreguntas(idEntrevista.getId(), limit) .collectList() // 过滤掉空结果,只保留有效列表 .filter(list -> !list.isEmpty()) // 将有效列表转为Flux输出 .flatMapMany(Flux::fromIterable) ) // 获取第一个有效结果流后立即终止,无需处理后续ID .next() // 将单个结果转为Flux(如果需要保留多个结果则去掉next()) .flatMapMany(Flux::just) // 所有ID均无有效结果时的处理,可根据业务调整为返回空Flux或错误 .switchIfEmpty(Mono.error(new RuntimeException("所有ID均未关联有效问题"))); }
关键逻辑说明
- concatMap替代flatMap:
concatMap保证按ID的顺序串行处理,避免并行调用带来的不必要资源消耗,符合"遍历直到找到有效结果"的需求。 - filter过滤空结果:自动跳过返回空列表的ID,无需手动抛错。
- next()终止流:拿到第一个有效结果后立即停止处理后续ID,提升效率。
- switchIfEmpty处理无有效结果场景:当所有ID都返回空时,可自定义错误或返回空Flux,适配不同业务需求。
内容的提问来源于stack exchange,提问作者DESTROYER GG
相关产品推荐
相关产品推荐

