使用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>
解决方案
核心问题分析
- CompletableFuture的
cancel(true)仅标记任务为取消状态,不会主动中断运行中的线程,必须手动在任务中处理中断信号。 - 客户端代码中的
ThreadUtil.fakeRandomSleep未处理中断,导致睡眠期间无法响应取消请求。 - 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
相关产品推荐
相关产品推荐

