如何确保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

