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

如何确保ExecutorService提交的Future存入Map后再触发/callback接口?

问题场景与解决方案

问题描述

我有两个接口/issue和/callback,调用/issue后会立即触发/callback。在/issue接口中,我向ExecutorService提交任务,并将返回的Future以订单ID为键存入ConcurrentHashMap;/callback接口被调用时,会通过订单ID从Map中获取Future并调用get()等待任务完成后再继续后续处理。

当前存在异常场景:调用/issue提交任务后,/callback在Future存入Map前就被调用,导致Map.get(orderId).get()抛出空指针异常。需要确保Future成功存入Map后,/callback才执行后续逻辑。

可行解决方案

方案1:使用CountDownLatch实现同步

核心思路是在/issue方法中先初始化计数为1的CountDownLatch,将其与Future的包装类提前存入Map;/callback获取包装类后先调用await()等待,直到/issue完成Future的存入并触发countDown(),再继续后续操作。

修改后的代码示例:

public class CallbackServiceImpl implements CallbackService {

    private final OrderService orderService;
    private final OrderRepository orderRepository;
    private final ExecutorService executorService = Executors.newCachedThreadPool();
    // 存储Future和CountDownLatch的包装类
    private final Map<String, FutureLatchWrapper> issueCertificateMap = new ConcurrentHashMap<>();
    private final Map<String, FutureLatchWrapper> renewCertificateMap = new ConcurrentHashMap<>();

    // 自定义包装类
    private static class FutureLatchWrapper {
        private Future<?> future;
        private final CountDownLatch latch = new CountDownLatch(1);

        public void setFuture(Future<?> future) {
            this.future = future;
            latch.countDown(); // 设置Future后触发计数减1
        }

        public Future<?> getFuture() throws InterruptedException {
            latch.await(); // 等待Future被设置完成
            return future;
        }
    }

    @Override
    public IssueCertificateAsyncResponseDto issueCertificateExecutorService(ConfirmSSLOrderRequestDto confirmSSLOrderRequestDto){
        String orderId = confirmSSLOrderRequestDto.getOrderid();
        // 先存入空的包装类占位
        FutureLatchWrapper wrapper = new FutureLatchWrapper();
        issueCertificateMap.put(orderId, wrapper);

        Runnable issueCertificateRunnable = () -> {
          try {
              orderService.issueCertificate(confirmSSLOrderRequestDto);
              log.info("ISSUE ORDER COMPLETED for order {}", confirmSSLOrderRequestDto.getOrderid());
          } catch (Exception e){
              log.error(e.getMessage());
              throw new AppException(e.getMessage());
          }
        };

        Future<?> issueCertificateFuture = executorService.submit(issueCertificateRunnable);
        wrapper.setFuture(issueCertificateFuture); // 设置Future并触发latch
        return IssueCertificateAsyncResponseDto.builder()
                .caOrderId(confirmSSLOrderRequestDto.getOrderid())
                .build();
    }

    private void handleIssueOrderCallbackRequest(SSLIssueOrderCallBackDto sslIssueOrderCallBackDto){
        String orderId = sslIssueOrderCallBackDto.getCaOrderId();
        Runnable issueCertificateCallBackRunnable = () -> {
            FutureLatchWrapper wrapper = null;
            try{
                wrapper = issueCertificateMap.get(orderId);
                if (wrapper == null) {
                    throw new AppException("订单不存在或已处理");
                }
                wrapper.getFuture().get(); // 等待Future被存入后再获取
                // 执行回调任务
            } catch (InterruptedException | ExecutionException e) {
                Thread.currentThread().interrupt();
                throw new AppException(e.getMessage());
            } finally {
                if (wrapper != null) {
                    issueCertificateMap.remove(orderId);
                }
            }
        };
        executorService.submit(issueCertificateCallBackRunnable);
    }
}

方案2:使用CompletableFuture替代Future(推荐)

CompletableFuture原生支持异步状态同步与异常处理,无需额外同步工具。我们可以先将未完成的CompletableFuture存入Map,/issue的任务完成后标记Future状态,/callback直接等待Future完成即可。

修改后的代码示例:

public class CallbackServiceImpl implements CallbackService {

    private final OrderService orderService;
    private final OrderRepository orderRepository;
    private final ExecutorService executorService = Executors.newCachedThreadPool();
    private final Map<String, CompletableFuture<Void>> issueCertificateFutureMap = new ConcurrentHashMap<>();
    private final Map<String, CompletableFuture<Void>> renewCertificateFutureMap = new ConcurrentHashMap<>();

    @Override
    public IssueCertificateAsyncResponseDto issueCertificateExecutorService(ConfirmSSLOrderRequestDto confirmSSLOrderRequestDto){
        String orderId = confirmSSLOrderRequestDto.getOrderid();
        // 先创建未完成的CompletableFuture并存入Map
        CompletableFuture<Void> future = new CompletableFuture<>();
        issueCertificateFutureMap.put(orderId, future);

        executorService.submit(() -> {
          try {
              orderService.issueCertificate(confirmSSLOrderRequestDto);
              log.info("ISSUE ORDER COMPLETED for order {}", orderId);
              future.complete(null); // 任务完成后标记Future完成
          } catch (Exception e){
              log.error(e.getMessage());
              future.completeExceptionally(e); // 异常时标记失败
              throw new AppException(e.getMessage());
          }
        });

        return IssueCertificateAsyncResponseDto.builder()
                .caOrderId(orderId)
                .build();
    }

    private void handleIssueOrderCallbackRequest(SSLIssueOrderCallBackDto sslIssueOrderCallBackDto){
        String orderId = sslIssueOrderCallBackDto.getCaOrderId();
        executorService.submit(() -> {
            CompletableFuture<Void> future = issueCertificateFutureMap.get(orderId);
            if (future == null) {
                throw new AppException("订单不存在或已处理");
            }
            try{
                future.join(); // 等待任务完成,join()会抛出未检查异常
                // 执行回调任务
            } catch (CompletionException e) {
                Thread.currentThread().interrupt();
                throw new AppException(e.getCause().getMessage());
            } finally {
                issueCertificateFutureMap.remove(orderId);
            }
        });
    }
}

方案3:使用AtomicReference配合忙等同步

利用ConcurrentHashMap的原子性computeIfAbsent方法,先存入AtomicReference占位,/issue完成任务提交后设置Future;/callback获取AtomicReference后循环等待,直到Future被填充完成。

修改后的代码示例:

public class CallbackServiceImpl implements CallbackService {

    private final OrderService orderService;
    private final OrderRepository orderRepository;
    private final ExecutorService executorService = Executors.newCachedThreadPool();
    private final Map<String, AtomicReference<Future<?>>> issueCertificateFutureMap = new ConcurrentHashMap<>();
    private final Map<String, AtomicReference<Future<?>>> renewCertificateFutureMap = new ConcurrentHashMap<>();

    @Override
    public IssueCertificateAsyncResponseDto issueCertificateExecutorService(ConfirmSSLOrderRequestDto confirmSSLOrderRequestDto){
        String orderId = confirmSSLOrderRequestDto.getOrderid();
        // 原子性创建AtomicReference并存入Map
        AtomicReference<Future<?>> futureRef = issueCertificateFutureMap.computeIfAbsent(orderId, k -> new AtomicReference<>());

        Runnable issueCertificateRunnable = () -> {
          try {
              orderService.issueCertificate(confirmSSLOrderRequestDto);
              log.info("ISSUE ORDER COMPLETED for order {}", orderId);
          } catch (Exception e){
              log.error(e.getMessage());
              throw new AppException(e.getMessage());
          }
        };

        Future<?> issueCertificateFuture = executorService.submit(issueCertificateRunnable);
        futureRef.set(issueCertificateFuture); // 设置Future

        return IssueCertificateAsyncResponseDto.builder()
                .caOrderId(orderId)
                .build();
    }

    private void handleIssueOrderCallbackRequest(SSLIssueOrderCallBackDto sslIssueOrderCallBackDto){
        String orderId = sslIssueOrderCallBackDto.getCaOrderId();
        executorService.submit(() -> {
            AtomicReference<Future<?>> futureRef = issueCertificateFutureMap.get(orderId);
            if (futureRef == null) {
                throw new AppException("订单不存在或已处理");
            }
            try{
                Future<?> future;
                // 循环等待Future被设置
                while ((future = futureRef.get()) == null) {
                    Thread.yield(); // 让出CPU,避免忙等消耗资源
                }
                future.get(); // 等待任务完成
                // 执行回调任务
            } catch (InterruptedException | ExecutionException e) {
                Thread.currentThread().interrupt();
                throw new AppException(e.getMessage());
            } finally {
                issueCertificateFutureMap.remove(orderId);
            }
        });
    }
}

方案对比

  • CountDownLatch方案:逻辑清晰,适合需要明确同步点的场景,但需要额外定义包装类。
  • CompletableFuture方案:最推荐,Java 8+原生支持,API简洁,自带异步状态管理与异常处理,无需额外同步组件。
  • AtomicReference+忙等方案:实现简单,但忙等可能占用CPU资源,适合低并发场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 11:07:23