从零开始,用启发式教学带你彻底搞懂"非阻塞延迟"这个利器。
每个 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] 程序结束观察与思考:
- 任务在 2 秒后才执行——延迟生效
- 主线程没有等待——
runAsync是异步的 - 执行线程是
ForkJoinPool.commonPool-worker-1——不是主线程
3.2 理解执行线程
问题:任务在哪个线程上执行?
delayedExecutor内部有两步:
- 计时:由 JDK 内部的 Delayer 守护线程负责
- 执行:计时到期后,任务被丢到
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→ 延迟清理 |