Java ExecutorService异常时重提交Callable任务问题排查与优化
问题场景
需要实现批量接口请求能力:向指定接口端点发起n(n≥0)次请求,针对500错误等服务端故障导致的失败请求执行重试(暂未实现重试间隔控制),重试终止条件为拿到请求成功响应,或是遇到非服务端故障导致的不可重试错误。
基于Java 11提供的Executors并发工具实现该功能时,实际运行效果不符合预期:无法持续重提交失败任务直到拿到全部响应,执行10次请求时最多只能拿到3~5个有效响应,无法定位问题原因。
初始实现代码
核心实现方法如下:
private void DoMyTasks(List<MyRequest> requests) { final ExecutorService executorService = Executors.newFixedThreadPool(10); final ExecutorCompletionService<MyReqResDto> completionService = new ExecutorCompletionService<>(executorService); for (final MyRequest MyRequest : requests) { completionService.submit(new MyCallableRequest(webClient, MyRequest)); } List<MyReqResDto> responses = new ArrayList<>(); for (int i = 0; i < requests.size(); ++i) { try { final Future<MyReqResDto> future = completionService.take(); if (future.get().getEx() != null) { completionService.submit(new MyCallableRequest(webClient, future.get().getMyRequest())); } responses.add(future.get()); } catch (ExecutionException | InterruptedException e) { log.warn("Error")); } catch (Exception exception) { log.error("Other error"); } finally { executorService.shutdown(); try { if (!executorService.awaitTermination(10, TimeUnit.MINUTES)) { executorService.shutdownNow(); } } catch (InterruptedException e) { executorService.shutdownNow(); } } } responses.size(); }
原有重试逻辑为判断返回结果携带异常时,重新提交对应请求任务:
if (future.get().getEx() != null) { completionService.submit(new MyCallableRequest(webClient, future.get().getMyRequest())); }
该逻辑运行后始终无法拿到全部请求的响应。
自定义Callable任务实现
自定义请求任务类代码如下:
public class MyCallableRequest implements Callable<MyReqResDto> { private final WebClient webClient; private final MyRequest myRequest; public MyCallableRequest(WebClient webClient, MyRequest myRequest) { this.webClient = webClient; this.myRequest = myRequest; } @Override public MyReqResDto call() throws Exception { try { if (new Random().nextInt(10) % 2 == 0) { throw new TestException(); } if (new Random().nextInt(10) % 7 == 0) { throw new RuntimeException(); } WebClient.UriSpec<WebClient.RequestBodySpec> uriSpec = webClient.post(); WebClient.RequestBodySpec bodySpec = uriSpec.uri( s -> s.path("/myEndpoint").build()); MyRequestDto myMyRequestDto = new MyRequestDto(); WebClient.RequestHeadersSpec<?> headersSpec = bodySpec.body(Mono.just(myMyRequestDto), MyRequestDto.class); ResponseDto responseDto = headersSpec.exchangeToMono(s -> { if (s.statusCode().equals(HttpStatus.OK)) { return s.bodyToMono(ResponseDto.class); } else if (s.statusCode().is1xxInformational()) { return s.createException().flatMap(Mono::error); } else if (s.statusCode().is3xxRedirection()) { return s.createException().flatMap(Mono::error); } else if (s.statusCode().is4xxClientError()) { return s.createException().flatMap(Mono::error); } else if (s.statusCode().is5xxServerError()) { return s.createException().flatMap(Mono::error); } else { return s.createException().flatMap(Mono::error); } //return null; }).block(); return new MyReqResDto(myRequestRequest, responseDto, null); } catch (Exception exception) { return new MyReqResDto(myRequest, null, exception); } } }
临时调整方案与疑问
参考社区建议将固定次数for循环改为while循环后,功能可正常运行,能够持续重试直到拿到全部无错响应,但担心该实现属于临时拼凑的脆弱方案,需要明确这种Executor使用方式存在哪些需要注意的线程相关问题。调整后代码如下:
while (true) { Future < MyReqResDto > future = null; try { future = completionService.take(); if (future.get().getEx() != null /*and check exception if possible to handle, if not break from a loop*/) { completionService.submit(new MyCallableRequest(webClient, future.get().getRequestCT()); } else { responseDtos.add(future.get()); } } catch (ExecutionException | InterruptedException e) { log.warn("Error while downloading", e.getCause()); // test if I can recover from these exceptions if no break; } } if (responseDtos.size() == requests.size()) { executorService.shutdown(); try { if (!executorService.awaitTermination(10, TimeUnit.MINUTES)) { executorService.shutdownNow(); } } catch (InterruptedException e) { executorService.shutdownNow(); } break; }
问题根因与实现注意事项
初始代码失效的核心原因
- 线程池关闭逻辑位置错误:初始代码将
executorService.shutdown()及后续等待终止的逻辑放在了for循环的finally块中,意味着第一次处理完单个任务结果后,就会触发线程池关闭。线程池进入shutdown状态后不再接收新提交的任务,后续重试提交的请求会直接被拒绝,队列中未执行的任务也会在等待终止超时后被强制中断,自然无法拿到全部响应。这也是10次请求最多只能返回3~5个结果的核心原因。 - 循环次数与实际待处理任务数不匹配:初始for循环的执行次数固定等于初始请求数,但重试操作会产生额外的待处理任务,固定次数的循环只会取前N个任务结果就退出,剩余重试任务的结果不会被处理,即使线程池正常运行也会漏结果。
- 无效结果被计入成功响应:初始代码无论任务返回成功结果还是带异常的失败结果,都会直接加入响应列表,最终返回的结果集包含无效的失败数据,不符合业务预期。
调整后while方案的潜在线程风险
- 无限循环风险:当前
while(true)没有全局兜底终止逻辑,如果所有请求持续返回可重试的5xx错误,代码会无限提交重试任务,不仅永远无法执行结束,还可能产生流量洪峰打垮服务端。必须补充两层兜底:单请求设置最大重试次数,全局任务设置最长执行超时时间。 - 异常分类逻辑缺失:当前仅判断结果携带异常就重试,没有区分异常类型。按照业务需求,4xx客户端错误、参数非法、序列化失败等不可重试的异常,应该直接标记对应请求失败,终止该请求的重试流程,不能无脑重试所有异常场景。
- 线程池资源泄漏风险:如果循环执行过程中抛出未捕获的异常,会直接跳出循环,此时不会触发线程池关闭逻辑,线程池会一直驻留内存占用系统资源。必须把线程池关闭逻辑放在整个任务处理流程最外层的
finally块中,无论执行成功、失败还是抛出异常,都保证线程池能被正确回收。 - 终止条件逻辑缺陷:当前仅靠
responseDtos.size() == requests.size()判断是否结束循环,如果出现不可重试的失败请求,成功响应数永远达不到初始请求总数,会直接进入死循环。需要单独维护两个计数器:已成功请求数、已判定为不可重试的失败请求数,两者相加等于初始请求总数时必须终止循环。 - 中断状态处理不规范:
completionService.take()是阻塞方法,响应线程中断时会抛出InterruptedException,当前代码捕获中断后直接break,既没有重置线程的中断标记,也没有做对应收尾处理,会导致上层逻辑无法正确感知线程中断状态。 - 响应式组件使用不当:WebClient是响应式HTTP客户端,调用链中随意使用
.block()阻塞当前线程的做法并不合理,尤其运行在线程池环境时,很容易因为线程被长时间占满导致吞吐量骤降,甚至触发死锁。如果使用WebClient,优先使用框架自带的重试算子(如retryWhen)实现重试逻辑,无需手动维护线程池提交任务,可大幅降低并发问题出现概率。
内容的提问来源于stack exchange,提问作者treo
相关产品推荐
相关产品推荐

