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

CompletableFuture异步逻辑异常:sleep仅生效一次,循环执行失效

问题描述

使用CompletableFuture的supplyAsync、thenApplyAsync、thenAcceptAsync实现异步支付逻辑,期望完成支付后休眠10分钟,再无限循环执行整个流程,但目前休眠仅生效一次,后续不再触发。

设计逻辑
  • 获取PaymentWorker实例,调用其call方法完成支付操作
  • 支付完成后将该worker从列表移除,避免重复处理
  • 休眠60秒后创建新的PaymentWorker,重复执行上述流程
  • 要求所有流程异步执行,且每个流程使用独立线程
相关代码

PaymentWorker类

public class PaymentWorker implements Callable<AutoPayResult> {
    private boolean isActive;
    private UUID workerId;

    @Override
    public AutoPayResultV2 call() {
        // 支付逻辑实现
    }

    public void stopWorker() {
        isActive = false;
    }

    // 补充必要的getter方法
    public UUID getWorkerId() {
        return workerId;
    }

    public boolean isActive() {
        return isActive;
    }
}

AutoPaymentServiceV2Impl类

@Service
@Slf4j
@RequiredArgsConstructor
public class AutoPaymentServiceV2Impl implements AutoPaymentServiceV2 {
    private final List<PaymentWorker> workers = new ArrayList<>();
    private volatile boolean isActive = false;

    private void runProcess(PaymentWorker z) {
        log.info("Starting worker with id: {}", z.getWorkerId());

        CompletableFuture
                .supplyAsync(z::call)
                .thenApplyAsync(s -> {
                    stopWorkerById(z);
                    return s;
                })
                .thenAcceptAsync(d -> {
                    if (isActive) {
                        var pw = new PaymentWorker(true, UUID.randomUUID());
                        workers.add(pw);
                        sleepAfterFinishWorker(pw.getBranch());
                        runProcess(pw);
                    } else {
                        log.warn("Auto payment V2 is not active and will be stopped");
                    }
                });
    }

    private void stopWorkerById(PaymentWorker paymentWorker) {
        log.info("Stop worker with id: {}", paymentWorker.getWorkerId());

        workers.stream()
                .filter(z -> z.getWorkerId().equals(paymentWorker.getWorkerId()))
                .findFirst()
                .ifPresent(w -> {
                    log.info("Worker found for removing");
                    if (w.isActive())
                        w.stopWorker();
                    workers.remove(w);
                });
    }

    private void sleepAfterFinishWorker(String branch) {
        var seconds = 60;
        log.info("Sleep worker after finishing for: {} seconds", seconds);
        LockSupport.parkNanos(TimeUnit.SECONDS.toNanos(seconds));
    }
}
问题分析与修复方案

核心问题

  1. 线程池资源耗尽:thenAcceptAsync默认使用ForkJoinPool公共线程池,回调中的休眠会占用线程池线程。第一次休眠后启动新流程,后续线程池资源被占满时,新的回调无法执行,导致休眠逻辑失效。
  2. 列表非线程安全:ArrayList不支持并发修改,多线程下add/remove操作会引发异常或数据不一致。
  3. 休眠方式不可靠:LockSupport.parkNanos可能被意外唤醒,导致休眠提前结束,且无法直观处理中断。

修复方案

1. 改用线程安全的列表

将ArrayList替换为CopyOnWriteArrayList,避免并发修改问题:

private final List<PaymentWorker> workers = new CopyOnWriteArrayList<>();

2. 分离休眠与流程启动逻辑

不在CompletableFuture的回调线程中执行休眠,而是使用CompletableFuture.delayedExecutor实现延迟启动新流程:

private void runProcess(PaymentWorker z) {
    log.info("Starting worker with id: {}", z.getWorkerId());

    CompletableFuture
            .supplyAsync(z::call)
            .thenApplyAsync(s -> {
                stopWorkerById(z);
                return s;
            })
            .thenAcceptAsync(d -> {
                if (isActive) {
                    var pw = new PaymentWorker(true, UUID.randomUUID());
                    workers.add(pw);
                    // 延迟60秒后启动新流程
                    CompletableFuture.runAsync(() -> runProcess(pw), 
                        CompletableFuture.delayedExecutor(60, TimeUnit.SECONDS));
                } else {
                    log.warn("Auto payment V2 is not active and will be stopped");
                }
            });
}

3. 替换休眠实现(可选)

如果需要保留显式休眠逻辑,改用Thread.sleep并处理中断:

private void sleepAfterFinishWorker(String branch) {
    var seconds = 60;
    log.info("Sleep worker after finishing for: {} seconds", seconds);
    try {
        Thread.sleep(TimeUnit.SECONDS.toMillis(seconds));
    } catch (InterruptedException e) {
        log.warn("Sleep interrupted", e);
        Thread.currentThread().interrupt();
    }
}

4. 确保线程独立性(可选)

如果需要完全独立的线程,可以自定义线程池,避免使用公共ForkJoinPool:

private final ExecutorService customExecutor = Executors.newCachedThreadPool();

private void runProcess(PaymentWorker z) {
    log.info("Starting worker with id: {}", z.getWorkerId());

    CompletableFuture
            .supplyAsync(z::call, customExecutor)
            .thenApplyAsync(s -> {
                stopWorkerById(z);
                return s;
            }, customExecutor)
            .thenAcceptAsync(d -> {
                if (isActive) {
                    var pw = new PaymentWorker(true, UUID.randomUUID());
                    workers.add(pw);
                    CompletableFuture.runAsync(() -> {
                        sleepAfterFinishWorker(pw.getBranch());
                        runProcess(pw);
                    }, customExecutor);
                } else {
                    log.warn("Auto payment V2 is not active and will be stopped");
                }
            }, customExecutor);
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 13:44:57