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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 23:27:09