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)); } }
问题分析与修复方案
核心问题
- 线程池资源耗尽:
thenAcceptAsync默认使用ForkJoinPool公共线程池,回调中的休眠会占用线程池线程。第一次休眠后启动新流程,后续线程池资源被占满时,新的回调无法执行,导致休眠逻辑失效。 - 列表非线程安全:
ArrayList不支持并发修改,多线程下add/remove操作会引发异常或数据不一致。 - 休眠方式不可靠:
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
相关产品推荐
相关产品推荐

