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

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均未关联有效问题")));
}

关键逻辑说明

  1. concatMap替代flatMap:concatMap保证按ID的顺序串行处理,避免并行调用带来的不必要资源消耗,符合"遍历直到找到有效结果"的需求。
  2. filter过滤空结果:自动跳过返回空列表的ID,无需手动抛错。
  3. next()终止流:拿到第一个有效结果后立即停止处理后续ID,提升效率。
  4. switchIfEmpty处理无有效结果场景:当所有ID都返回空时,可自定义错误或返回空Flux,适配不同业务需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 22:26:01