Java 11中基于CompletableFuture实现Saga扇出模式的并发优化问题
问题描述
我正尝试在Java 11 Spring组件中建模扇出Saga模式,目标是利用Java并发类跟踪列表中大量请求的聚合结果状态,以便在全部成功、任一失败或无响应时执行对应操作。
我的实体类定义如下,Request对象的状态会根据远程响应更新:
class SagaObject { private long sagaId; private List<Request> requests; } class Request { private long requestId; private RequestStatus status = RequestStatus.PENDING; } enum RequestStatus { PENDING,SUCCESS,FAILED; }
查阅示例后,常见做法是为每个Request创建CompletableFuture,再通过allOf()分组:
// 伪代码 // 为每个请求创建CompletableFuture List<CompletableFuture<Request>> completableFutures = sagaObject.requests.stream().collect(Collectors.toList()); // 用于跟踪所有单个请求的CompletableFuture CompletableFuture<Void> allFutures = CompletableFuture .allOf(completableFutures.toArray(new CompletableFuture[sagaObject.requests.size()])); allFutures.get(); // 根据allOf()的结果执行操作 if (allFutures.isDone()) { // 所有CompletableFuture完成后,仍需遍历列表检查结果 }
但我的场景中请求数量可能很大(约10万条),我担心为每个请求创建单独的CompletableFuture会带来过大的开销,这种顾虑是否合理?
我考虑让SagaObject实现Callable/Supplier接口,如下所示:
// 伪代码 class SagaObject implements Callable<RequestStatus> { private long sagaId; private List<Request> requests; public RequestStatus call() throws Exception { // 遍历请求列表,若仍有PENDING则返回异常,否则返回SUCCESS/FAILED return ..; } }
随后通过ExecutorService调用:
// 伪代码 ExecutorService executorService = Executors.newCachedThreadPool(); try { Future<RequestStatus> sagaFuture = executorService.submit(sagaObject); if (sagaFuture.isDone()) { sagaFuture.get(); } } finally { executorService.shutdown(); }
我不确定哪种方案更符合JVM并发模型,希望得到相关建议。
解决方案与建议
关于CompletableFuture的开销顾虑
你的顾虑有一定合理性,但核心风险并非CompletableFuture实例本身:
- 单个CompletableFuture实例的内存开销极小(约几十字节),10万个实例的总内存占用仅几MB,不会对JVM造成显著负担。
- 真正的风险在于线程资源耗尽:如果未限制线程池大小,直接为10万请求分配并发线程,每个线程默认1MB栈内存的话,总内存需求会达到100GB,直接压垮JVM。只要用自定义线程池控制并发数,就能规避这个问题。
两种方案的适用性对比
1. CompletableFuture方案(推荐)
这是Java 8+官方推荐的异步编程模型,完全契合JVM并发设计理念:
- 优势:
- 支持异步回调,单个请求完成时可立即更新状态,无需等待全部请求结束,适合需要实时跟踪部分请求状态的场景。
- 内置
allOf()、anyOf()方法,可直接实现“全部成功触发后续逻辑”“任一失败立即启动补偿”的Saga核心逻辑,代码简洁高效。 - 结合自定义线程池可精准控制并发度,避免资源浪费。
- 优化要点:
必须使用自定义线程池执行异步任务,示例代码如下:// 自定义线程池,根据系统配置调整参数 ThreadPoolExecutor customPool = new ThreadPoolExecutor( Runtime.getRuntime().availableProcessors() * 2, // 核心线程数 Runtime.getRuntime().availableProcessors() * 4, // 最大线程数 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>(10000), // 任务队列 Executors.defaultThreadFactory(), new ThreadPoolExecutor.CallerRunsPolicy() // 拒绝策略 ); // 为每个请求创建带线程池的CompletableFuture List<CompletableFuture<Request>> futures = sagaObject.getRequests().stream() .map(request -> CompletableFuture.supplyAsync(() -> { // 执行远程请求并更新状态 request.setStatus(executeRemoteCall(request) ? RequestStatus.SUCCESS : RequestStatus.FAILED); return request; }, customPool)) .collect(Collectors.toList()); // 监听全部完成事件 CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])) .thenRun(() -> { // 遍历检查所有请求状态,执行后续逻辑 boolean allSuccess = futures.stream() .map(CompletableFuture::join) .allMatch(req -> RequestStatus.SUCCESS.equals(req.getStatus())); if (allSuccess) { // 全部成功逻辑 } else { // 存在失败请求,执行补偿逻辑 } }); // 监听任一失败事件,提前触发补偿 CompletableFuture.anyOf(futures.stream() .map(future -> future.exceptionally(ex -> { // 处理异常,标记请求为失败 request.setStatus(RequestStatus.FAILED); return request; })) .toArray(CompletableFuture[]::new)) .thenAccept(req -> { // 任一请求失败,立即启动Saga补偿 compensateSaga(sagaObject); });
2. SagaObject实现Callable方案
这种方案存在明显缺陷,完全不适合10万请求的场景:
- 本质是单线程轮询所有请求状态,没有利用并发优势,轮询逻辑会成为性能瓶颈,无法实时响应请求的完成/失败事件。
- 如果
call()方法中采用阻塞等待,会占用线程资源却不做有效工作;如果是非阻塞轮询,会导致CPU空转,浪费系统资源。 - 仅适合请求数量极少、无需并发执行的边缘场景,完全不符合JVM并发模型的设计初衷。
额外注意事项
- 线程安全的状态更新:Request的
status字段是多线程访问的,必须保证线程安全,建议用AtomicReference<RequestStatus>替代普通字段,或者对状态更新操作加同步锁。 - 超时处理:必须为异步请求设置超时,避免Saga无限等待无响应的请求,可使用
CompletableFuture.orTimeout()方法实现:CompletableFuture.supplyAsync(() -> executeRemoteCall(request), customPool) .orTimeout(5, TimeUnit.SECONDS) // 设置5秒超时 .exceptionally(ex -> { request.setStatus(RequestStatus.FAILED); return request; }); - Spring环境适配:在Spring组件中,可直接使用
TaskExecutor或@Async注解配合自定义线程池,无需手动创建ThreadPoolExecutor,更符合Spring生态的最佳实践。
内容的提问来源于stack exchange,提问作者emeraldjava
相关产品推荐
相关产品推荐

