☰
从零开始,带你彻底搞懂“非阻塞延迟“这个利器
2026/10/8 6:53:49 网站建设 项目流程

从零开始,用启发式教学带你彻底搞懂"非阻塞延迟"这个利器。
每个 API 先讲清功能,再上代码,注释逐行拆解。


第一章:我们遇到了什么问题?

1.1 一个真实场景

假设你在写一个订单服务。调用第三方支付接口时,偶尔会失败。你的第一反应是什么?

"失败了就重试啊。"

好,那重试之间要不要等一等?毕竟对方服务可能正在抖动,立刻重试大概率还是失败。

"等个 1 秒再试。"

很合理。那你会怎么写?

public String callPayment() { for (int i = 0; i < 3; i++) { try { return doCallPayment(); // 调用支付接口 } catch (Exception e) { Thread.sleep(1000); // 等 1 秒再重试 } } throw new RuntimeException("支付调用失败"); }

这段代码有什么问题?

1.2 问题出在哪?

Thread.sleep(1000)做了什么?它让当前线程冻结 1 秒。

思考:如果这段代码运行在线程池里——比如一个 Web 服务器处理请求的线程池——会发生什么?

线程池(共 4 个线程): T1: [处理请求A] → 支付失败 → sleep(1s) → 重试 → 成功 T2: [处理请求B] → 支付失败 → sleep(1s) → 重试 → 成功 T3: [处理请求C] → 支付失败 → sleep(1s) → 重试 → 成功 T4: [处理请求D] → 支付失败 → sleep(1s) → 重试 → 成功 问题:4 个线程全在 sleep!新的请求 E 来了,没人处理!

核心矛盾:线程在sleep期间什么都没干,但占着位置不放。就像 4 个厨师都在等锅热,但灶台被他们占着,别的菜做不了。

1.3 我们想要什么?

我们想要的其实是:"等 1 秒,然后继续做"——但等待期间,线程应该被释放回去干别的活。

理想情况: T1: [处理请求A] → 支付失败 → 释放线程!→ (1秒后回来) → 重试 → 成功 → 释放 T2: [处理请求B] → 支付失败 → 释放线程!→ (1秒后回来) → 重试 → 成功 → 释放 ↑ 等待期间 T1、T2 可以去处理请求 E、F、G...

怎么实现"释放线程,1 秒后自动回来"?

这就是delayedExecutor要解决的问题。


第二章:认识 delayedExecutor

2.1 它是什么?

CompletableFuture.delayedExecutor是 JDK 9 引入的一个静态方法。

public static Executor delayedExecutor(long delay, TimeUnit unit)

一句话:它返回一个特殊的Executor——你往这个 Executor 提交的任务,不会立即执行,而是等指定时间后才执行。

2.2 先搞懂 Executor 是什么

在 Java 中,Executor是一个接口,只有一个方法:

public interface Executor { void execute(Runnable command); // 执行一个任务 }

常见的 Executor 实现:

实现类行为
ThreadPoolExecutor提交后立即交给线程池执行
ForkJoinPool.commonPool()提交后立即交给公共池执行
delayedExecutor(1, SECONDS)提交后等 1 秒再执行

关键理解:delayedExecutor返回的也是一个Executor,只是它的execute方法多了一个"延迟"行为。

2.3 它内部怎么工作?

你的代码 JDK 内部 ───────── ────────── ┌──────────────────────────────┐ │ Delayer(守护线程,只有 1 个) │ │ 职责:专门负责计时 │ └──────────────────────────────┘ │ delayedExecutor(2, SECONDS) │ .execute(myTask) ──────────▶ 告诉 Delayer:"2 秒后执行 myTask" │ 你的线程池继续干别的事 Delayer 开始计时... 2 秒到了 → 把 myTask 丢到线程池执行

要点:

  • 等待期间,你的线程池完全不受影响——线程该干嘛干嘛
  • 计时由 JDK 内部唯一的 1 个守护线程负责——它不是你的线程池里的线程
  • "零额外线程占用"指的是不占你的用户线程池

2.4 它解决什么问题?不适合解决什么问题?

适合的场景:

场景说明
非阻塞重试退避失败后等 N 秒重试,等待期间不占线程
延迟任务调度"5 秒后发一条通知"、"10 秒后检查状态"
限流/冷却两次操作之间强制间隔
超时补偿超时后延迟执行清理逻辑

不适合的场景:

场景为什么不适合更好的选择
精确定时任务(每天 8 点执行)delayedExecutor只做相对延迟,不做绝对时间ScheduledExecutorService.scheduleAtFixedRate
周期性重复任务它只延迟一次,不自动重复ScheduledExecutorService
需要取消/修改延迟一旦提交无法取消ScheduledFuture

第三章:从最简单的例子开始

3.1 Hello delayedExecutor

先写一个最简单的程序,感受"延迟执行"的效果。

import java.util.concurrent.*; /** * 最基础的 delayedExecutor 使用演示。 * * 目标:提交一个任务,让它 2 秒后才执行。 */ public class HelloDelayedExecutor { public static void main(String[] args) throws Exception { // 记录当前时间,用于计算耗时 long start = System.currentTimeMillis(); System.out.println("[" + elapsed(start) + "ms] 准备提交延迟任务"); // ─── 核心代码:创建延迟执行器 ─── // delayedExecutor(2, SECONDS) 返回一个 Executor // 这个 Executor 的特点是:提交的任务会延迟 2 秒才真正执行 Executor delayed = CompletableFuture.delayedExecutor(2, TimeUnit.SECONDS); // ─── 提交任务 ─── // runAsync 把一个 Runnable 任务提交到指定的 Executor 上执行 // 因为我们的 Executor 是 delayedExecutor,所以任务会延迟 2 秒 CompletableFuture.runAsync(() -> { // 这个 lambda 就是我们要执行的任务 // 它会在 2 秒后才真正运行 System.out.println("[" + elapsed(start) + "ms] 任务执行了!" + " 线程名:" + Thread.currentThread().getName()); }, delayed); // ↑ 注意第二个参数:指定在哪个 Executor 上执行 // 这里传入的是 delayedExecutor,所以会延迟 2 秒 // ─── 验证:主线程没有被阻塞 ─── // runAsync 是异步的——它立即返回,不会等待任务完成 // 所以这行代码会立刻执行,不会等 2 秒 System.out.println("[" + elapsed(start) + "ms] 主线程继续执行,没有被阻塞"); // 等待程序退出(否则主线程结束,JVM 退出,延迟任务来不及执行) Thread.sleep(4000); System.out.println("[" + elapsed(start) + "ms] 程序结束"); } /** * 计算从 start 到现在经过了多少毫秒 */ static long elapsed(long start) { return System.currentTimeMillis() - start; } }

运行结果:

[0ms] 准备提交延迟任务 [0ms] 主线程继续执行,没有被阻塞 ← 立即执行,没有等 2 秒 [2001ms] 任务执行了! 线程名:ForkJoinPool.commonPool-worker-1 [4001ms] 程序结束

观察与思考:

  1. 任务在 2 秒后才执行——延迟生效
  2. 主线程没有等待——runAsync是异步的
  3. 执行线程是ForkJoinPool.commonPool-worker-1——不是主线程

3.2 理解执行线程

问题:任务在哪个线程上执行?

delayedExecutor内部有两步:

  1. 计时:由 JDK 内部的 Delayer 守护线程负责
  2. 执行:计时到期后,任务被丢到ForkJoinPool.commonPool()执行

所以任务实际运行在commonPool的线程上,而不是 Delayer 线程上。

如果你想让任务在你自己的线程池上执行呢?

import java.util.concurrent.*; /** * 演示:让延迟任务在自定义线程池上执行。 */ public class DelayedWithCustomPool { public static void main(String[] args) throws Exception { // 创建我们自己的线程池(方便观察线程名) ExecutorService myPool = Executors.newFixedThreadPool(2, r -> { Thread t = new Thread(r); t.setName("my-pool-" + t.getId()); // 自定义线程名 return t; }); long start = System.currentTimeMillis(); // 方案:先用 delayedExecutor 延迟,再用 supplyAsync 指定线程池 // // 执行流程: // 1. delayedExecutor 等 1 秒 // 2. 1 秒后,把任务丢到 myPool 执行 CompletableFuture.supplyAsync(() -> { System.out.println("[" + elapsed(start) + "ms] 任务执行!线程:" + Thread.currentThread().getName()); return "结果"; }, CompletableFuture.delayedExecutor(1, TimeUnit.SECONDS, myPool)); // ↑ 第三个参数:指定执行线程池 // 延迟到期后,任务在这个线程池上执行 Thread.sleep(3000); myPool.shutdown(); } static long elapsed(long start) { return System.currentTimeMillis() - start; } }

运行结果:

[1001ms] 任务执行!线程:my-pool-14

关键发现:delayedExecutor有一个三参数重载版本:

// 两参数版:延迟后在 commonPool 执行 delayedExecutor(1, TimeUnit.SECONDS) // 三参数版:延迟后在指定 Executor 执行 delayedExecutor(1, TimeUnit.SECONDS, myPool)

第四章:核心场景——非阻塞重试

4.1 先用阻塞方式实现重试(对比基线)

import java.util.concurrent.*; import java.util.concurrent.atomic.*; /** * 阻塞式重试:Thread.sleep 方式。 * * 问题:sleep 期间线程被冻结,无法服务其他任务。 */ public class BlockingRetry { // 模拟计数器:前 2 次调用失败,第 3 次成功 static final AtomicInteger counter = new AtomicInteger(0); public static void main(String[] args) throws Exception { // 只有 1 个线程的线程池——问题更明显 ExecutorService pool = Executors.newFixedThreadPool(1); long start = System.currentTimeMillis(); // 提交 2 个任务 pool.submit(() -> retryTask("任务A", start)); pool.submit(() -> retryTask("任务B", start)); // ↑ 任务B 必须等任务A 完全结束(包括 sleep)才能开始 pool.shutdown(); pool.awaitTermination(30, TimeUnit.SECONDS); } static void retryTask(String name, long start) { for (int attempt = 1; attempt <= 3; attempt++) { counter.incrementAndGet(); System.out.println("[" + elapsed(start) + "ms] " + name + " 第" + attempt + "次尝试(线程:" + Thread.currentThread().getName() + ")"); // 模拟:前 4 次调用失败,第 5 次成功 if (counter.get() < 5) { System.out.println(" → 失败,等待 1 秒后重试..."); try { // ★ 问题所在:线程在 sleep 期间完全冻结 ★ Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); return; } } else { System.out.println(" → 成功!"); return; } } } static long elapsed(long start) { return System.currentTimeMillis() - start; } }

运行结果:

[0ms] 任务A 第1次尝试(线程:pool-1-thread-1) → 失败,等待 1 秒后重试... ↑ 线程在 sleep,任务B 进不来! [1001ms] 任务A 第2次尝试(线程:pool-1-thread-1) → 失败,等待 1 秒后重试... [2001ms] 任务A 第3次尝试(线程:pool-1-thread-1) → 成功! ↑ 任务A 终于结束,任务B 可以开始了 [3001ms] 任务B 第1次尝试(线程:pool-1-thread-1) → 失败,等待 1 秒后重试... [4001ms] 任务B 第2次尝试(线程:pool-1-thread-1) → 成功!

问题:任务 B 等了 3 秒才开始!因为线程池只有 1 个线程,被任务 A 的sleep占满了。

4.2 用 delayedExecutor 实现非阻塞重试

现在用delayedExecutor来解决这个问题。核心思路:

失败 → 不 sleep → 线程立即释放 → 用 delayedExecutor 设一个"1秒后的闹钟" → 闹钟响了 → 再提交一次任务到线程池
import java.util.concurrent.*; import java.util.concurrent.atomic.*; /** * 非阻塞重试:delayedExecutor 方式。 * * 对比上一节:同样的场景(1 线程池 + 2 个任务 + 重试), * 观察任务 B 是否还需要等待。 */ public class NonBlockingRetry { static final AtomicInteger counter = new AtomicInteger(0); public static void main(String[] args) throws Exception { ExecutorService pool = Executors.newFixedThreadPool(1); long start = System.currentTimeMillis(); // 提交 2 个任务——使用非阻塞重试 CompletableFuture.runAsync(() -> nonBlockingRetry("任务A", pool, start), pool); CompletableFuture.runAsync(() -> nonBlockingRetry("任务B", pool, start), pool); Thread.sleep(10000); pool.shutdown(); } /** * 非阻塞重试方法。 * * @param name 任务名 * @param pool 线程池(重试时要把任务提交回同一个池) * @param start 起始时间(用于打印耗时) */ static void nonBlockingRetry(String name, ExecutorService pool, long start) { int count = counter.incrementAndGet(); System.out.println("[" + elapsed(start) + "ms] " + name + " 尝试(线程:" + Thread.currentThread().getName() + ")"); if (count < 5) { // 失败!但不 sleep——线程立即释放! System.out.println(" → 失败,1 秒后重试(线程已释放)..."); // ★ 核心代码 ★ // // runAsync(() -> {}, ...) 中的 () -> {} 是空操作。 // 我们不需要它"做事"——我们只需要 delayedExecutor 的"延迟"效果。 // // 整个表达式的含义: // "等 1 秒后,执行 whenComplete 里的代码" // // 等 1 秒期间: // - 当前线程已经释放,可以去处理其他任务 // - JDK 内部的 Delayer 守护线程在计时 CompletableFuture.runAsync(() -> {}, CompletableFuture.delayedExecutor(1, TimeUnit.SECONDS)) .whenComplete((v, ex) -> { // 1 秒后到达这里——重新尝试 // 把任务提交回线程池 pool.submit(() -> nonBlockingRetry(name, pool, start)); }); } else { System.out.println(" → 成功!"); } } static long elapsed(long start) { return System.currentTimeMillis() - start; } }

运行结果:

[0ms] 任务A 尝试(线程:pool-1-thread-1) → 失败,1 秒后重试(线程已释放)... ↑ 线程释放了!任务B 可以进来了! [0ms] 任务B 尝试(线程:pool-1-thread-1) → 失败,1 秒后重试(线程已释放)... ↑ 两个任务的"闹钟"同时在计时 [1001ms] 任务A 尝试(线程:pool-1-thread-1) → 失败,1 秒后重试(线程已释放)... [2001ms] 任务A 尝试(线程:pool-1-thread-1) → 成功! [2001ms] 任务B 尝试(线程:pool-1-thread-1) → 成功!

对比:

阻塞版: 任务A 开始:0ms 任务B 开始:3001ms ← 等了 3 秒! 全部完成:5001ms 非阻塞版: 任务A 开始:0ms 任务B 开始:0ms ← 立即开始! 全部完成:3001ms ← 快了将近 2 秒!

4.3 更优雅的实现:Promise 递归模式

上面的实现有个问题:每次重试都重新调用nonBlockingRetry方法,逻辑分散。更好的方式是用Promise 模式——一个方法返回 Future,内部递归链接重试。

先理解 Promise 模式:

import java.util.concurrent.*; /** * Promise 模式入门。 * * 核心思想:创建一个 CompletableFuture("承诺"), * 稍后在某个时机完成它(complete 或 completeExceptionally)。 * 调用方拿到这个 Future,用 whenComplete/get 等待结果。 */ public class PromiseDemo { public static void main(String[] args) throws Exception { long start = System.currentTimeMillis(); // 创建一个"承诺"——稍后会被完成 CompletableFuture<String> promise = new CompletableFuture<>(); // 2 秒后完成这个承诺 CompletableFuture.runAsync(() -> {}, CompletableFuture.delayedExecutor(2, TimeUnit.SECONDS)) .whenComplete((v, ex) -> { // 2 秒后到达这里 System.out.println("[" + elapsed(start) + "ms] 延迟到期,完成 promise"); promise.complete("Hello from delayedExecutor!"); // ↑ 完成 promise,调用方拿到结果 }); // 调用方:等待 promise 完成 // get() 会阻塞当前线程,直到 promise 被 complete String result = promise.get(); System.out.println("[" + elapsed(start) + "ms] 拿到结果:" + result); } static long elapsed(long start) { return System.currentTimeMillis() - start; } }

运行结果:

[2001ms] 延迟到期,完成 promise [2001ms] 拿到结果:Hello from delayedExecutor!

现在用 Promise 模式重写非阻塞重试:

import java.util.concurrent.*; import java.util.concurrent.atomic.*; import java.util.function.*; /** * 用 Promise + delayedExecutor 实现优雅的非阻塞重试。 * * 这个 retryAsync 方法可以直接作为工具类使用。 */ public class RetryWithPromise { static final AtomicInteger counter = new AtomicInteger(0); public static void main(String[] args) throws Exception { ExecutorService pool = Executors.newFixedThreadPool(2); long start = System.currentTimeMillis(); // 使用 retryAsync 工具方法 CompletableFuture<String> result = retryAsync( () -> { // 要重试的操作 int c = counter.incrementAndGet(); System.out.println("[" + elapsed(start) + "ms] 尝试 #" + c + "(线程:" + Thread.currentThread().getName() + ")"); if (c < 3) { throw new RuntimeException("失败"); // 模拟失败 } return "第 " + c + " 次成功!"; }, 3, // 最大重试次数 500, // 退避延迟(毫秒) pool // 线程池 ); // 等待最终结果 System.out.println("最终结果:" + result.get()); pool.shutdown(); } /** * 非阻塞重试工具方法。 * * @param action 要执行的操作(可能抛异常) * @param maxRetries 最大重试次数 * @param backoffMs 退避延迟(毫秒) * @param executor 线程池 * @return 最终结果的 Future */ static <T> CompletableFuture<T> retryAsync( Callable<T> action, int maxRetries, long backoffMs, ExecutorService executor) { // 创建"承诺"——调用方持有这个 Future 等待最终结果 CompletableFuture<T> promise = new CompletableFuture<>(); // 启动第一次尝试,然后递归链接重试 attemptAndRetry(action, maxRetries, backoffMs, executor, promise, 0); return promise; } /** * 执行一次尝试,失败时用 delayedExecutor 延迟后递归重试。 * * @param action 要执行的操作 * @param maxRetries 最大重试次数 * @param backoffMs 退避延迟 * @param executor 线程池 * @param promise 最终承诺 * @param attempt 当前是第几次尝试(从 0 开始) */ static <T> void attemptAndRetry( Callable<T> action, int maxRetries, long backoffMs, ExecutorService executor, CompletableFuture<T> promise, int attempt) { // supplyAsync:把任务提交到线程池异步执行 // 返回一个 Future,任务完成时 Future 自动完成,异常时异常完成 CompletableFuture.supplyAsync(() -> { try { return action.call(); // 执行用户操作 } catch (Exception e) { // 把受检异常包装为 CompletionException,让 supplyAsync 的 Future 异常完成 throw new CompletionException(e); } }, executor) // whenComplete:无论成功还是异常,回调都执行 // 参数:result = 成功时的值(异常时为 null),ex = 异常(成功时为 null) .whenComplete((result, ex) -> { if (ex == null) { // ─── 成功路径 ─── // 任务成功了,完成 promise,调用方拿到结果 promise.complete(result); return; } // ─── 异常路径 ─── if (attempt >= maxRetries) { // 重试次数耗尽,传播异常 promise.completeExceptionally(ex); return; } // 还有重试机会——用 delayedExecutor 延迟后重试 System.out.println(" 第 " + (attempt + 1) + " 次重试,延迟 " + backoffMs + "ms"); // ★ 核心:delayedExecutor 延迟 → whenComplete 触发 → 递归重试 ★ CompletableFuture.runAsync(() -> {}, CompletableFuture.delayedExecutor(backoffMs, TimeUnit.MILLISECONDS)) .whenComplete((v, delayEx) -> { // 延迟到期!递归调用自身,进行下一次尝试 attemptAndRetry(action, maxRetries, backoffMs, executor, promise, attempt + 1); }); }); } static long elapsed(long start) { return System.currentTimeMillis() - start; } }

运行结果:

[0ms] 尝试 #1(线程:pool-1-thread-1) 第 1 次重试,延迟 500ms [501ms] 尝试 #2(线程:pool-1-thread-1) 第 2 次重试,延迟 500ms [1001ms] 尝试 #3(线程:pool-1-thread-1) 最终结果:第 3 次成功!

执行流程图:

supplyAsync(尝试#1) → 失败 ↓ whenComplete: 失败,还有重试机会 ↓ runAsync(() -> {}, delayedExecutor(500ms)) ← 设置 500ms 闹钟 ↓ (线程释放!) ... 500ms 后 ... ↓ Delayer: 叮! whenComplete: 延迟到期 ↓ supplyAsync(尝试#2) → 失败 ↓ whenComplete: 失败,还有重试机会 ↓ runAsync(() -> {}, delayedExecutor(500ms)) ← 再设一个闹钟 ↓ (线程释放!) ... 500ms 后 ... ↓ Delayer: 叮! supplyAsync(尝试#3) → 成功! ↓ promise.complete("第 3 次成功!")

第五章:进阶用法

5.1 指数退避

实际生产中,重试延迟通常逐次增加(指数退避),避免对故障服务造成压力。

import java.util.concurrent.*; /** * 指数退避重试:每次重试延迟翻倍。 * * 第 1 次重试:100ms * 第 2 次重试:200ms * 第 3 次重试:400ms * 第 4 次重试:800ms * ... */ public class ExponentialBackoff { public static void main(String[] args) throws Exception { ExecutorService pool = Executors.newFixedThreadPool(2); long start = System.currentTimeMillis(); retryWithExponentialBackoff( () -> { int c = Counter.next(); System.out.println("[" + elapsed(start) + "ms] 尝试 #" + c); if (c < 4) throw new RuntimeException("失败"); return "成功"; }, 5, // 最大重试 5 次 100, // 初始延迟 100ms pool ).get(); pool.shutdown(); } static <T> CompletableFuture<T> retryWithExponentialBackoff( Callable<T> action, int maxRetries, long initialDelayMs, ExecutorService executor) { return doRetry(action, maxRetries, initialDelayMs, executor, new CompletableFuture<>(), 0); } static <T> CompletableFuture<T> doRetry( Callable<T> action, int maxRetries, long currentDelayMs, ExecutorService executor, CompletableFuture<T> promise, int attempt) { CompletableFuture.supplyAsync(() -> { try { return action.call(); } catch (Exception e) { throw new CompletionException(e); } }, executor).whenComplete((result, ex) -> { if (ex == null) { promise.complete(result); return; } if (attempt >= maxRetries) { promise.completeExceptionally(ex); return; } // ★ 指数退避:当前延迟 = 初始延迟 × 2^attempt ★ // attempt=0: 100 × 2^0 = 100ms // attempt=1: 100 × 2^1 = 200ms // attempt=2: 100 × 2^2 = 400ms long delay = currentDelayMs; System.out.println(" 重试 #" + (attempt + 1) + ",延迟 " + delay + "ms"); CompletableFuture.runAsync(() -> {}, CompletableFuture.delayedExecutor(delay, TimeUnit.MILLISECONDS)) .whenComplete((v, e) -> // ★ 下次延迟翻倍 ★ doRetry(action, maxRetries, currentDelayMs * 2, executor, promise, attempt + 1)); }); return promise; } // 简单的全局计数器 static class Counter { static int n = 0; static int next() { return ++n; } } static long elapsed(long start) { return System.currentTimeMillis() - start; } }

运行结果:

[0ms] 尝试 #1 重试 #1,延迟 100ms [101ms] 尝试 #2 重试 #2,延迟 200ms [302ms] 尝试 #3 重试 #3,延迟 400ms [703ms] 尝试 #4 最终结果:成功

5.2 与其他 CompletableFuture 操作组合

delayedExecutor可以和其他 CF 操作自由组合:

import java.util.concurrent.*; /** * delayedExecutor 与其他 CF 操作的组合用法。 */ public class ComposingDelayed { public static void main(String[] args) throws Exception { long start = System.currentTimeMillis(); // ─── 场景 1:延迟后转换结果 ─── // "等 1 秒,然后计算 42 的 2 倍" CompletableFuture.runAsync(() -> {}, CompletableFuture.delayedExecutor(1, TimeUnit.SECONDS)) .thenApply(v -> 42 * 2) // 延迟后转换 .whenComplete((result, ex) -> System.out.println("[" + elapsed(start) + "ms] 场景1结果:" + result)); // 输出:[1001ms] 场景1结果:84 // ─── 场景 2:延迟后执行异步操作 ─── // "等 1 秒,然后发起一个异步查询" CompletableFuture.runAsync(() -> {}, CompletableFuture.delayedExecutor(1, TimeUnit.SECONDS)) .thenCompose(v -> CompletableFuture.supplyAsync(() -> "查询结果")) .whenComplete((result, ex) -> System.out.println("[" + elapsed(start) + "ms] 场景2结果:" + result)); // 输出:[1001ms] 场景2结果:查询结果 // ─── 场景 3:多个延迟任务并行 ─── // "1 秒后做 A,2 秒后做 B,等两个都完成" CompletableFuture<String> taskA = CompletableFuture.runAsync(() -> {}, CompletableFuture.delayedExecutor(1, TimeUnit.SECONDS)) .thenApply(v -> "A完成"); CompletableFuture<String> taskB = CompletableFuture.runAsync(() -> {}, CompletableFuture.delayedExecutor(2, TimeUnit.SECONDS)) .thenApply(v -> "B完成"); CompletableFuture.allOf(taskA, taskB) .thenApply(v -> taskA.join() + " + " + taskB.join()) .whenComplete((result, ex) -> System.out.println("[" + elapsed(start) + "ms] 场景3结果:" + result)); // 输出:[2001ms] 场景3结果:A完成 + B完成 Thread.sleep(4000); } static long elapsed(long start) { return System.currentTimeMillis() - start; } }

第六章:常见陷阱

陷阱 1:忘记把任务提交到 delayedExecutor

// ❌ 错误:只是创建了 delayedExecutor,什么都没做 CompletableFuture.delayedExecutor(1, TimeUnit.SECONDS); // 这行代码没有任何效果! // ✅ 正确:必须配合 runAsync / supplyAsync 使用 CompletableFuture.runAsync(() -> doSomething(), CompletableFuture.delayedExecutor(1, TimeUnit.SECONDS));

陷阱 2:用 supplyAsync 但 lambda 没有返回值

// ❌ 编译错误:supplyAsync 要求 Supplier 有返回值 CompletableFuture.supplyAsync(() -> { System.out.println("执行"); // 没有 return!编译报错 }, CompletableFuture.delayedExecutor(1, TimeUnit.SECONDS)); // ✅ 正确:无返回值用 runAsync CompletableFuture.runAsync(() -> { System.out.println("执行"); }, CompletableFuture.delayedExecutor(1, TimeUnit.SECONDS));

陷阱 3:以为 delayedExecutor 会中断已有任务

// ❌ 误解:以为 1 秒后会中断正在运行的 supplyAsync 任务 CompletableFuture.supplyAsync(() -> { Thread.sleep(10000); // 任务要跑 10 秒 return "done"; }, pool); // delayedExecutor 跟上面这个任务没有任何关系! // 它只负责"延迟后执行新提交的任务" // ✅ 正确理解:delayedExecutor 只是一个定时器 // 它不取消、不中断、不影响任何已经在执行的任务

陷阱 4:在延迟回调里阻塞等待

// ❌ 错误:延迟后又去 get() 阻塞——白用了非阻塞延迟 CompletableFuture.runAsync(() -> {}, CompletableFuture.delayedExecutor(1, TimeUnit.SECONDS)) .thenApply(v -> someOtherFuture.get()); // 阻塞! // ✅ 正确:继续用异步操作 CompletableFuture.runAsync(() -> {}, CompletableFuture.delayedExecutor(1, TimeUnit.SECONDS)) .thenCompose(v -> someOtherFuture); // 非阻塞

陷阱 5:忽略返回的 Future

// ❌ 问题:提交了延迟任务,但没有持有 Future 引用 // 如果主线程退出,延迟任务可能来不及执行 CompletableFuture.runAsync(() -> {}, CompletableFuture.delayedExecutor(5, TimeUnit.SECONDS)); // 主线程立即结束 → JVM 退出 → 延迟任务没机会执行 // ✅ 正确:持有 Future 引用,确保程序不会提前退出 CompletableFuture<?> delayed = CompletableFuture.runAsync(() -> {}, CompletableFuture.delayedExecutor(5, TimeUnit.SECONDS)); delayed.get(); // 等待延迟任务完成 // 或者:delayed.join();

总结

一句话记忆

delayedExecutor(delay, unit)= 一个"定时闹钟"

告诉它"X 秒后做某事",然后你就可以走开干别的了。
等待期间,你的线程池零占用。

使用模板

// 基本用法:延迟后执行操作 CompletableFuture.runAsync(() -> { // 这里是要做的事 }, CompletableFuture.delayedExecutor(延迟时间, 时间单位)); // 指定线程池:延迟后在自定义线程池上执行 CompletableFuture.supplyAsync(() -> { return 结果; }, CompletableFuture.delayedExecutor(延迟时间, 时间单位, 自定义线程池)); // 链式操作:延迟后继续做其他事 CompletableFuture.runAsync(() -> {}, CompletableFuture.delayedExecutor(1, TimeUnit.SECONDS)) .thenCompose(v -> 下一个异步操作) .whenComplete((result, ex) -> 处理结果);

适用场景速查

场景用法
非阻塞重试退避失败 →delayedExecutor→ 重新提交任务
延迟通知"5 秒后发消息" →runAsync(() -> send(), delayedExecutor(5, SECONDS))
延迟健康检查"启动 30 秒后检查" →supplyAsync(() -> check(), delayedExecutor(30, SECONDS))
操作冷却间隔两次调用间强制间隔 → 链式delayedExecutor
超时后清理orTimeout异常 →whenComplete→delayedExecutor→ 延迟清理

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询