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

使用Resilience4j @TimeLimiter无法终止CompletableFuture任务问题排查

问题:CompletableFuture超时后无法终止客户端任务执行

我遇到了无法取消/中断CompletableFuture任务的问题:控制器已触发超时,但客户端任务仍会继续执行直至完成。想知道缺少哪些配置或代码,才能实现超时后终止客户端执行。

以下是我的相关代码与配置:

服务层代码

@TimeLimiter(name = "ly-service-timelimiter", fallbackMethod = "fallFn")
@Bulkhead(name = "ly-service-bulkhead", fallbackMethod = "fallFn", type = Bulkhead.Type.THREADPOOL)
@Override
public CompletableFuture<Void> myMethod(Request request) throws Exception {
    try {
        log.info("MyMethod Service: {}", request);
        return client.myMethod(request);
    } catch (RuntimeException e) {
        log.info("Exception", request);
        throw new RuntimeException(e);
    }
}

客户端代码

public CompletableFuture<Void> myMethod(Request request) {
    CompletableFuture<Void> future = new CompletableFuture<>();

    CompletableFuture.runAsync(() -> {
        if (future.isCancelled()) {
            log.info("MyMethod was cancelled before execution.");
            return;
        }

        try {
            log.info("Processing request", request);
            ThreadUtil.fakeRandomSleep(10000); // Simulating work

            if (future.isCancelled()) {
                log.info("Processing was cancelled during execution.");
            } else {
                log.info("Completed routing with TimeOut");
                future.complete(null);
            }
        } catch (Exception e) {
            log.info("completeExceptionally......");
            future.completeExceptionally(e);
        }
    });

    return future;
}

控制器代码

@GetMapping("/runTimeOut")
public @ResponseBody String executeSample() throws ExecutionException, InterruptedException {

    log.info("Execute TimeOut EndPoint");
    Request request = Request.builder().build(); //Any class
    try {
        myService.myMethod(request).get();
    }catch (ExecutionException  ex){
        if(ex.getCause() instanceof java.util.concurrent.TimeoutException){
            log.info("TimeoutException occurred");
        }
        return "Failed";
    }

    return "OK";
}

配置文件

resilience4j:
  bulkhead:
    configs:
      default:
        max-concurrent-calls: 2
        max-wait-duration: 0ms
    instances:
      ly-service-bulkhead:
        base-config: default
  timelimiter:
    instances:
      ly-service-timelimiter:
        timeoutDuration: 900ms
        cancel-running-future: true

依赖配置

<parent>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-parent</artifactId>
    <version>3.1.2</version>
    <relativePath/> 
</parent>
...

    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-aop</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-actuator</artifactId>
    </dependency>
    <dependency>
        <groupId>io.github.resilience4j</groupId>
        <artifactId>resilience4j-spring-boot2</artifactId>
        <version>2.1.0</version>
    </dependency>

解决方案

核心问题分析

  1. CompletableFuture的cancel(true)仅标记任务为取消状态,不会主动中断运行中的线程,必须手动在任务中处理中断信号。
  2. 客户端代码中的ThreadUtil.fakeRandomSleep未处理中断,导致睡眠期间无法响应取消请求。
  3. Resilience4j的cancel-running-future: true仅调用future的cancel方法,无法直接中断线程,需要配合线程中断逻辑。

一、修改客户端代码,支持线程中断

替换不可中断的睡眠方法,添加线程中断监听与处理逻辑:

public CompletableFuture<Void> myMethod(Request request) {
    CompletableFuture<Void> future = new CompletableFuture<>();
    Thread workerThread = null;

    CompletableFuture.runAsync(() -> {
        workerThread = Thread.currentThread();
        try {
            // 执行前检查中断状态
            if (Thread.currentThread().isInterrupted()) {
                log.info("MyMethod was cancelled before execution.");
                future.cancel(true);
                return;
            }

            log.info("Processing request: {}", request);
            // 使用支持中断的Thread.sleep,替换原有的不可中断睡眠
            Thread.sleep(10000); 

            // 执行后检查中断状态
            if (Thread.currentThread().isInterrupted()) {
                log.info("Processing was cancelled during execution.");
                future.cancel(true);
            } else {
                log.info("Completed routing without TimeOut");
                future.complete(null);
            }
        } catch (InterruptedException e) {
            // 捕获中断信号,终止任务
            log.info("Processing was interrupted due to timeout.");
            future.cancel(true);
            // 重置中断状态,让上层逻辑感知
            Thread.currentThread().interrupt();
        } catch (Exception e) {
            log.info("completeExceptionally......");
            future.completeExceptionally(e);
        }
    });

    // 当future被取消时,主动中断工作线程
    future.whenComplete((v, t) -> {
        if (future.isCancelled() && workerThread != null) {
            workerThread.interrupt();
        }
    });

    return future;
}

二、确认Resilience4j配置有效性

当前配置中的cancel-running-future: true是正确的,需确保:

  • timeoutDuration: 900ms小于客户端任务的预估执行时间(当前为10s),确保能触发超时逻辑。
  • THREADPOOL类型的Bulkhead配置无冲突,当前的max-concurrent-calls:2与max-wait-duration:0ms是合理的。

三、控制器代码优化(可选)

补充中断异常处理,避免线程中断状态丢失:

@GetMapping("/runTimeOut")
public @ResponseBody String executeSample() {
    log.info("Execute TimeOut EndPoint");
    Request request = Request.builder().build();
    try {
        myService.myMethod(request).get();
    } catch (ExecutionException ex) {
        if (ex.getCause() instanceof TimeoutException) {
            log.info("TimeoutException occurred");
        }
        return "Failed";
    } catch (InterruptedException e) {
        log.info("Request was interrupted");
        Thread.currentThread().interrupt();
        return "Interrupted";
    }
    return "OK";
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 21:47:36