Java + Quarkus 动态并行调用多REST API并合并结果方案问询
解决方案:Quarkus 动态并行调用 REST API
核心思路
先建立provider 名称到 API 调用逻辑的映射表,运行时根据传入的 providers 列表动态匹配对应的调用逻辑,再借助 Mutiny 的并行组合 API 实现并发执行,同时单独处理每个调用的失败场景,确保单个失败不影响整体流程。
具体实现步骤
1. 定义 Provider 与 API 调用的映射
创建线程安全的映射表,将每个 provider 名称关联到返回 Uni<ApiResult> 的调用方法(ApiResult 为自定义的统一结果类型):
import io.smallrye.mutiny.Uni; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.function.Supplier; // 自定义统一返回结果类型 class ApiResult { private String providerName; private Object data; // 构造器、getter/setter 省略 public ApiResult(String providerName, Object data) { this.providerName = providerName; this.data = data; } } @ApplicationScoped public class DynamicApiCaller { // 存储 provider 到对应 API 调用的映射 private final Map<String, Supplier<Uni<ApiResult>>> apiMapping = new ConcurrentHashMap<>(); // 初始化映射,可根据需求扩展更多 provider public DynamicApiCaller() { apiMapping.put("A", this::callApiA); apiMapping.put("B", this::callApiB); apiMapping.put("C", this::callApiC); } // 示例 API 调用方法,实际替换为你的 REST 客户端调用 private Uni<ApiResult> callApiA() { return RestClientA.fetchData() .map(data -> new ApiResult("A", data)); } private Uni<ApiResult> callApiB() { return RestClientB.getResource() .map(data -> new ApiResult("B", data)); } private Uni<ApiResult> callApiC() { return RestClientC.queryInfo() .map(data -> new ApiResult("C", data)); } // 对外暴露映射表,供后续过滤使用 public Map<String, Supplier<Uni<ApiResult>>> getApiMapping() { return apiMapping; } }
2. 动态构建并行调用并处理结果
在资源类中接收 providers 列表,动态生成 API 调用流,并行执行后过滤成功结果返回:
import io.smallrye.mutiny.Multi; import io.smallrye.mutiny.Uni; import jakarta.inject.Inject; import jakarta.ws.rs.POST; import jakarta.ws.rs.Path; import java.util.List; import java.util.logging.Logger; @Path("/aggregate") public class AggregateResource { private static final Logger LOG = Logger.getLogger(AggregateResource.class.getName()); @Inject DynamicApiCaller apiCaller; @POST public Uni<List<ApiResult>> aggregate(List<String> providers) { // 1. 过滤有效 provider,生成对应的 API 调用 Uni 集合 List<Uni<ApiResult>> apiCalls = providers.stream() .filter(apiCaller.getApiMapping()::containsKey) .map(provider -> apiCaller.getApiMapping().get(provider).get() // 单个调用失败时记录日志,恢复为 null 以便后续过滤 .onFailure().invoke(err -> LOG.severe("Provider " + provider + " 调用失败: " + err.getMessage())) .onFailure().recoverWithNull()) .toList(); // 2. 并行执行所有调用,收集非空的成功结果 return Multi.createFrom().iterable(apiCalls) .merge() // 按调用完成顺序收集结果,用 concatenate() 可保持原顺序 .filter(result -> result != null) .collect().asList(); } }
3. 关键细节说明
- 动态扩展性:映射表支持运行时修改,可通过配置加载、事件监听等方式动态添加/更新 provider 对应的调用逻辑。
- 失败隔离:每个 API 调用单独处理失败,仅记录日志并返回 null,不会中断其他调用的执行。
- 并行控制:Mutiny 的
merge()/concatenate()会自动并行执行所有Uni,无需手动管理线程池,Quarkus 会根据配置优化资源占用。 - 无效 Provider 过滤:提前过滤映射表中不存在的 provider,避免无效调用和空指针问题。
替代方案:等待所有调用完成后合并结果
如果需要等待所有调用(包括失败恢复后的)完成再返回结果,可使用 Uni.combine():
return Uni.combine().all().unis(apiCalls) .combinedWith(results -> results.stream() .filter(result -> result != null) .toList());
内容的提问来源于stack exchange,提问作者Akhil Prajapati
相关产品推荐
相关产品推荐

