电商聚合查询应使用completablefuture实现异步并行调用商品、库存、价格等服务,指定自定义线程池,通过allof组合结果、exceptionally处理异常、thencompose链式编排、ortimeout设置超时与降级。

在电商系统中,聚合查询通常需要同时调用多个服务(如商品、库存、价格、优惠券、用户画像等),再合并结果返回给前端。用 CompletableFuture 可以高效实现异步并行、错误隔离和灵活编排,避免传统同步串行导致的响应延迟。
并行拉取多个服务数据
每个子查询封装为独立的 CompletableFuture,用 supplyAsync 提交到线程池执行,互不阻塞:
CompletableFuture<product> productFuture = CompletableFuture.supplyAsync(
() -> productService.findById(productId), executor);
CompletableFuture<stock> stockFuture = CompletableFuture.supplyAsync(
() -> stockService.getStock(productId), executor);
CompletableFuture<price> priceFuture = CompletableFuture.supplyAsync(
() -> priceService.getCurrentPrice(productId), executor);
</price></stock></product>
注意:务必指定自定义线程池(如 ForkJoinPool.commonPool() 不适合 IO 密集型任务),避免耗尽 Tomcat 线程。
组合结果并处理异常
用 allOf 等待全部完成,再用 join() 获取结果;对单个失败不中断整体流程,可设默认值或记录告警:
CompletableFuture<void> allFutures = CompletableFuture.allOf(
productFuture, stockFuture, priceFuture);
AggregatedResult result = allFutures.thenApply(v -> {
Product p = productFuture.join(); // join 不抛异常,失败时抛 CompletionException
Stock s = stockFuture.exceptionally(t -> {
log.warn("库存查询失败", t);
return new Stock(0, "UNKNOWN");
}).join();
Price pr = priceFuture.join();
return new AggregatedResult(p, s, pr);
}).join();
</void>
- 不要用
get(),它会抛受检异常且可能阻塞;join()更简洁,配合exceptionally或handle处理失败 -
allOf返回CompletableFuture<void></void>,需手动提取各 future 结果 - 若某个服务超时,可在
supplyAsync内部加TimeoutException包装,或用orTimeout(JDK 9+)
按依赖关系编排(如先查商品再查关联优惠)
用 thenCompose 实现异步链式调用,前序结果作为后序入参:
CompletableFuture<aggregatedresult> fullResult =
CompletableFuture.supplyAsync(() -> productService.findById(productId))
.thenCompose(product -> {
if (product == null) {
return CompletableFuture.completedFuture(
new AggregatedResult(null, null, null));
}
// 根据商品类目查对应优惠
return CompletableFuture.supplyAsync(
() -> couponService.getApplicableCoupon(product.getCategoryId()))
.thenApply(coupon -> new AggregatedResult(product, coupon));
});
</aggregatedresult>
-
thenApply用于同步转换;thenCompose用于返回新的CompletableFuture,避免嵌套 - 链路中任意环节失败,整个 future 会失败,可用
exceptionally统一兜底
超时控制与降级策略
电商场景对响应时间敏感,必须设置合理超时,并提供轻量降级结果:
CompletableFuture<product> productWithTimeout =
CompletableFuture.supplyAsync(() -> productService.findById(productId), executor)
.orTimeout(800, TimeUnit.MILLISECONDS)
.exceptionally(t -> {
if (t instanceof TimeoutException) {
log.warn("商品服务超时,启用缓存降级");
return productCache.get(productId); // 本地缓存或空对象
}
log.error("商品查询异常", t);
return null;
});
</product>
-
orTimeout(JDK 9+)会主动完成 future 并抛TimeoutException;低版本可用completeOnTimeout+orTimeout模拟 - 降级逻辑应尽量轻量,避免再触发远程调用形成雪崩
- 建议配合熔断器(如 Resilience4j)做失败率统计与自动熔断
Java免费学习笔记:立即使用
解锁 Java 大师之旅:从入门到精通的终极指南











