Java实现并行方法调用且异常不终止全部处理的方案咨询
问题
现有一段同步循环执行的代码,逻辑是遍历ID列表逐个调用disablePackXYZ方法,异常时仅记录日志不中断循环:
private void disableXYZ(Long rId, List<Long> disableIds, String requestedBy) { for (Long disableId : disableIds) { try { disablePackXYZ(UnSubscribeRequest.unsubscriptionRequest() .requestedBy(requestedBy) .cancellationReason("system") .id(disableId) .build()); } catch (Exception e) { log.error("Failed to disable pack. id: {}, rId: {}. Error: {}", disableId, rId, e.getMessage()); } } }
现在需要改成并行调用disablePackXYZ,要求不管单个或多个调用抛异常,都不能终止其他并行任务,必须等所有任务执行完,异常只做捕获记录,不能打断其他线程或后续任务。
之前试过用Java并行流实现,但发现只要有一个线程抛异常,forEach就会立刻把异常抛给调用方,导致没法等所有线程执行完:
final CompletableFuture<ParseException> thrownException = new CompletableFuture<>(); Stream.of(columns).parallel().forEach(column -> { try { result[column.index] = parseColumn(valueCache[column.index], column.type); } catch (ParseException e) { thrownException.complete(e); }});
可行方案
方案一:用CompletableFuture.allOf(推荐)
借助CompletableFuture把每个任务异步提交,然后用allOf等待所有任务完成,每个任务内部单独捕获异常打日志,保证所有任务都能跑完。
private void disableXYZ(Long rId, List<Long> disableIds, String requestedBy) { // 把每个ID对应的任务转成异步Future List<CompletableFuture<Void>> futures = disableIds.stream() .map(disableId -> CompletableFuture.runAsync(() -> { try { disablePackXYZ(UnSubscribeRequest.unsubscriptionRequest() .requestedBy(requestedBy) .cancellationReason("system") .id(disableId) .build()); } catch (Exception e) { log.error("Failed to disable pack. id: {}, rId: {}. Error: {}", disableId, rId, e.getMessage()); } })) .collect(Collectors.toList()); // 等待所有异步任务执行完毕,不管成功失败 CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join(); }
- 好处:Java 8+原生API,不用手动管线程池(默认用公共ForkJoinPool,也可以传自定义线程池),代码简洁还能保证所有任务执行完。
- 提示:如果业务对线程池有特殊要求,比如需要隔离线程,就在
runAsync里传自己定义的Executor,别占用公共线程池影响其他业务。
方案二:用ExecutorService手动管理线程
自己创建线程池提交所有任务,然后调用awaitTermination等所有任务跑完,每个任务内部自己捕获异常。
private void disableXYZ(Long rId, List<Long> disableIds, String requestedBy) { // 线程池大小可以根据业务调整,比如取ID数量和10的最小值,避免创建太多线程 ExecutorService executor = Executors.newFixedThreadPool(Math.min(disableIds.size(), 10)); try { for (Long disableId : disableIds) { executor.submit(() -> { try { disablePackXYZ(UnSubscribeRequest.unsubscriptionRequest() .requestedBy(requestedBy) .cancellationReason("system") .id(disableId) .build()); } catch (Exception e) { log.error("Failed to disable pack. id: {}, rId: {}. Error: {}", disableId, rId, e.getMessage()); } }); } executor.shutdown(); // 等待所有任务完成,超时时间可以根据业务设置,比如1小时 executor.awaitTermination(1, TimeUnit.HOURS); } catch (InterruptedException e) { log.error("Task execution was interrupted", e); Thread.currentThread().interrupt(); } }
- 好处:线程池参数完全可控,能根据业务场景调整核心线程数、队列大小这些配置。
- 注意:必须调用
shutdown()和awaitTermination(),不然线程池不会关闭,而且没法确保等所有任务执行完。
方案三:改进并行流实现
如果一定要用并行流,核心是绝对不能让异常逃出lambda表达式,而且用collect替代forEach当终端操作——因为forEach遇到未捕获异常会提前终止,而collect会等所有元素处理完。
private void disableXYZ(Long rId, List<Long> disableIds, String requestedBy) { disableIds.parallelStream() .map(disableId -> { try { disablePackXYZ(UnSubscribeRequest.unsubscriptionRequest() .requestedBy(requestedBy) .cancellationReason("system") .id(disableId) .build()); return null; // 没异常就返回占位符 } catch (Exception e) { log.error("Failed to disable pack. id: {}, rId: {}. Error: {}", disableId, rId, e.getMessage()); return e; // 捕获异常返回,后续如果要统一处理可以用这个(可选) } }) .collect(Collectors.toList()); // 这个终端操作会确保所有元素都处理完 }
- 注意:并行流用的是公共ForkJoinPool,如果是CPU密集型任务,可能会影响其他用并行流的业务,所以核心业务里还是更推荐用自定义线程池的异步方案。
内容的提问来源于stack exchange,提问作者SaurabhVS
相关产品推荐
相关产品推荐

