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

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并发模型的设计初衷。

额外注意事项

  1. 线程安全的状态更新:Request的status字段是多线程访问的,必须保证线程安全,建议用AtomicReference<RequestStatus>替代普通字段,或者对状态更新操作加同步锁。
  2. 超时处理:必须为异步请求设置超时,避免Saga无限等待无响应的请求,可使用CompletableFuture.orTimeout()方法实现:
    CompletableFuture.supplyAsync(() -> executeRemoteCall(request), customPool)
        .orTimeout(5, TimeUnit.SECONDS) // 设置5秒超时
        .exceptionally(ex -> {
            request.setStatus(RequestStatus.FAILED);
            return request;
        });
    
  3. Spring环境适配:在Spring组件中,可直接使用TaskExecutor或@Async注解配合自定义线程池,无需手动创建ThreadPoolExecutor,更符合Spring生态的最佳实践。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 17:54:56