Phaser 可重用动态线程同步屏障

Phaser 可重用动态线程同步屏障:双栈编排机制、64 位状态字与多阶段调度全解析 🤔 一、道格·李为什么需要比 CyclicBarrier 更灵活的屏障 CountDownLatch 和 CyclicBarrier 分别覆盖了两种同步场景:前者是"一个线程等 N 个线程完成",后者是"N 个线程彼此等到齐后一起走"。但道格·李在后续实践中发现了一个覆盖盲区:多阶段计算中,每阶段的参与者数量可能并不相同。 举个例子:分片计算一个大型数据集,第一阶段 8 个线程并行处理各自的分片;第二阶段某些分片的数据已经为空,对应的线程应该退出,剩下 5 个线程继续;第三阶段可能又加入 2 个新线程处理汇总结果。 CountDownLatch 做不到——它是一次性的,三个阶段需要三个实例。 CyclicBarrier 也做不到——它的 parties 数量在构造时固定,运行期间不能增删参与者。如果有线程中途退出,CyclicBarrier 会永远等不到第 N 个线程而永久阻塞(或者触发 BrokenBarrierException)。 道格·李因此在 Java 7 引入了 Phaser:一个支持动态参与者数量 + 多阶段循环使用的同步屏障。线程可以在运行时通过 register() 加入、通过 arriveAndDeregister() 退出,Phaser 自动调整每轮的等待计数。内部用一个 64 位的 state 字段打包了阶段号、已到达计数、未到达计数等所有状态信息,通过 CAS 无锁操作更新——这是 JUC 中最复杂的一个状态字设计。 🎚️ 二、Phaser 核心概念(术语定义) 在深入源码前,先明确几个关键术语: 术语 定义 类比理解 Phase(阶段号) 从 0 开始递增的整数,每轮同步完成后 +1 表示"第几轮同步" Party(参与者) 注册到 Phaser 中的一个线程/任务 需要等待的对象 Unarrived(未到达数) 当前阶段尚未调用 arrive() 的参与者数量 每到达一个就减 1 Arrive(到达) 线程调用 arrive() 表示完成当前阶段工作 通知 Phaser"我到了" Advance(推进) 当 unarrived 归零时,phase 自增,进入下一轮 所有人都到了,开始下一阶段 Register(注册) 增加一个参与者(parties + 1, unarrived + 1) 动态加入 Deregister(注销) 减少一个参与者(parties - 1, unarrived - 1) 动态退出 Termination(终止) Phaser 进入终止态,所有操作立即返回负数 强制结束,不再同步 🏗️ 三、数据结构展开 🔢 3.1 64 位状态字:所有信息的原子载体 Phaser 没有使用 AQS,而是直接将全部状态压缩在一个 AtomicLong(字段名 state)中。这是理解 Phaser 的根基。 ...

九月 7, 2022 · 11 分钟 · 2232 字 · yaomingye

CyclicBarrier 可循环屏障

CyclicBarrier 可循环屏障:源码解析、代际机制与 CountDownLatch 对比全解析 🤔 一、道格·李为什么需要一个可循环的屏障 CountDownLatch 解决了一个问题:一个线程等待多个线程完成操作。但道格·李在设计 JSR 166 时意识到,还有一种更复杂的同步场景没有覆盖:多个线程彼此等待——所有线程都到达同一个"集合点"后,再一起继续往下走。这在分片并行计算中非常常见:N 个线程各算各的,算完之后需要"对表"(交叉校验、汇总),然后继续算下一阶段。 CountDownLatch 做不了这件事——它是一次性的,计数器归零后无法重置。而且它的语义是"一个线程等 N 个线程",不是"N 个线程彼此等"。 Thread.join() 也做不了——join() 等的是线程终止,不是线程到达某个执行点。如果线程需要继续执行(而不是终止),join() 完全不对路。 道格·李因此设计了 CyclicBarrier:一组线程各自执行到某个"屏障点"后调用 await(),先到的线程阻塞等待,直到最后一个线程也到达屏障,所有线程同时被唤醒,继续往下执行。屏障打开后自动重置,可以用于下一个阶段——这就是 Cyclic(可循环)的含义。 与 CountDownLatch 的核心设计区别: CountDownLatch:外部协调者等待 N 个工人完成任务(一次性,一个等 N 个) CyclicBarrier:N 个工人彼此等到齐后一起行动(可循环,N 个彼此等) 🔄 二、数据结构展开:CyclicBarrier 的六大核心字段 📌 2.1 字段总览 // java.util.concurrent.CyclicBarrier public class CyclicBarrier { private final ReentrantLock lock = new ReentrantLock(); // ① 锁 private final Condition trip = lock.newCondition(); // ② 条件队列 private final int parties; // ③ 参与方总数 private final Runnable barrierCommand; // ④ 屏障动作 private Generation generation = new Generation(); // ⑤ 当前代际 private int count; // ⑥ 倒计数 // 内部类——代际 private static class Generation { Generation() {} // 默认 broken = false boolean broken; // 当前代是否被打破 } } 用一张结构图展示这些字段之间的关系: ...

九月 5, 2022 · 10 分钟 · 1922 字 · yaomingye

CompletableFuture 异步编排

CompletableFuture 异步编排:Completion 链表与多中间件应用全解析 🤔 一、道格·李为什么需要比 Future 更强大的异步工具 Java 5 引入了 Future 接口和 FutureTask 实现,解决了"异步执行、获取返回值"的基础需求。到 Java 7 时代,Future 的局限已经非常明显:它只是一个结果的容器,没有回调机制——你不能在结果就绪时自动触发下一步操作,只能调用 get() 阻塞等待。 这在简单的"提交任务→等待结果"场景中够用,但面对以下需求时完全无力: 链式编排:A 的结果作为 B 的输入,B 完成后触发 C。用 Future 只能嵌套 get(),代码缩进越来越深 多结果组合:等 A、B、C 三个结果全部就绪后做汇总。用 Future 只能逐个 get(),最慢的那个决定了总耗时 异常传播:Future.get() 把异常包装成 ExecutionException,调用方需要捕获后 getCause()——异常处理散落在各处 道格·李在设计 CompletableFuture(Java 8 引入)时参考了 JavaScript 的 Promise 模式和函数式编程中的 monad 概念。核心思路是:把异步计算的结果建模为一条流水线——每个阶段接受上一个阶段的输出,产生下一个阶段的输入,阶段之间通过回调串联。这样开发者不需要手动管理线程和等待,只需要"声明"各个步骤之间的关系: // 声明式:订单+用户拼好,再拼优惠券 orderFuture.thenCombine(userFuture, this::mergeOrderUser) .thenCombine(couponFuture, this::assembleFinal); CompletableFuture 的革新在于把异步编程从"命令式等结果"推到了"声明式编排"——关心的不是什么时候拿到结果,而是结果拿到之后要做什么。 🔮 二、数据结构:CompletableFuture 内部长什么样 ⚙️ 2.1 核心字段 从 JDK 源码中看 CompletableFuture<T> 的结构定义(Java 8,java.util.concurrent.CompletableFuture): ...

九月 3, 2022 · 9 分钟 · 1825 字 · yaomingye

FutureTask 源码深度解析

FutureTask 源码深度解析:从 Runnable 的局限到异步结果获取的完整实现 🤔 一、道格·李为什么需要一个"身兼两职"的任务对象 Java 的 Thread 构造函数接受 Runnable,但 Runnable.run() 返回值是 void——执行完就完了,拿不到结果。在 Java 1.0 ~ 1.4 时代,想在主线程拿到子线程的计算结果,只能靠共享变量(比如把结果写进一个 final int[] result = new int[1]),这种写法没有类型安全,也无法向调用方传递异常。 道格·李在 Java 5 的 JSR 166 中为这个问题设计了三个层次: 第一层:Callable<V>——任务接口。和 Runnable 功能等价,但 call() 有返回值且可抛异常。解决了"任务有结果"的问题。 第二层:Future<V>——结果句柄。提供 get()(阻塞获取结果)、cancel()(取消任务)、isDone()(判断完成)等方法。解决了"怎么拿到异步结果"的问题。但它只是一个接口,不知道任务在哪执行、怎么执行。 第三层:FutureTask——把两者粘在一起。它同时实现了 RunnableFuture<V> 接口(该接口同时继承 Runnable 和 Future),所以一个 FutureTask 对象既是可执行的任务(可以传给 Thread 或提交给 Executor),又是可查询的结果句柄(可以 get() 拿结果、cancel() 取消)。 道格·李这个设计的精巧之处在于:通过 FutureTask 这个"桥梁",ExecutorService.submit(Callable) 可以把任意 Callable 包装成 FutureTask,提交到线程池执行后立即返回 Future 句柄——调用方拿到了一个"未来的结果承诺",可以继续干别的事,需要结果时再 get()。 🔮 二、类继承体系:RunnableFuture 的双重身份 🏗️ 2.1 继承结构图 flowchart LR %% 半暗底色 + 高亮描边:完美适配博客深色/浅色双主题 %% classDef reject fill:#450a0a,stroke:#dc2626,stroke-width:2px,color:#fecaca,font-weight:bold; classDef process fill:#1e1e24,stroke:#6b7280,stroke-width:2px,color:#e5e7eb; R[Runnable\nvoid run] F[Future\nget/cancel/isDone] RF[RunnableFuture\n同时继承Runnable+Future] FT[FutureTask] C[Callable\nV call throws Exception] R --> RF F --> RF RF --> FT C -->|组合-而非继承| FT class F,FT,R,RF process; class C reject; 接口/类 角色 核心方法 Runnable 可执行任务 void run() Callable<V> 有结果的任务 V call() throws Exception Future<V> 结果句柄 get(), cancel(), isDone() RunnableFuture<V> 二者的桥接接口 继承 Runnable + Future,无新增方法 FutureTask<V> 具体实现 组合 Callable,实现所有逻辑 📌 2.2 RunnableFuture 接口 public interface RunnableFuture<V> extends Runnable, Future<V> { void run(); } 这个接口只有 3 行,没有新增任何方法,只是将 Runnable 和 Future 合并。它的价值在于类型层面的统一——一个 RunnableFuture 实例可以同时作为任务提交给线程池和作为 Future 供调用方查询结果。 ...

八月 31, 2022 · 11 分钟 · 2292 字 · yaomingye

Semaphore 源码解析

Semaphore 源码解析:AQS 共享模式、许可传播机制与公平策略实现 🚀 道格·李为什么需要一个信号量 信号量(Semaphore)是操作系统教科书里最古老的并发原语之一,由 Edsger Dijkstra 在 1960 年代提出。但在 Java 1.0 ~ 1.4 时代,JDK 里并没有信号量——开发者只能用 synchronized 加一个计数器模拟,代码又长又容易出错。 道格·李在设计 JSR 166 时,需要将信号量引入 Java,原因很简单:synchronized 是互斥的(同一时刻只能一个线程进入),而很多并发控制场景需要的是"限制并发数量"而不是"限制到只有 1 个"。比如数据库连接池最多 10 个并发连接、文件读取最多 3 个线程同时打开、API 限流每秒 100 个请求——这些场景用 synchronized 无法表达。 Semaphore 的思路是许可计数:构造时定义 N 个"许可证",线程调用 acquire() 拿走一个许可(不够就阻塞),用完调用 release() 归还。许可与线程没有绑定关系——线程 A 获取的许可可以由线程 B 释放。这个设计很关键:它让 Semaphore 不仅可以用作"限流器",还可以用作"对象池管理器"或"同步器"。 在 AQS 框架中,Semaphore 使用的是共享模式(Shared Mode)——多个线程可以同时获取许可,不像 ReentrantLock 的独占模式那样一次只唤醒一个线程。 🏗️ 核心数据结构 🏗️ 整体类层次结构 Semaphore 和 ReentrantLock 的内部结构高度相似——都是通过内部类 Sync 间接继承 AQS,并通过两个子类实现公平/非公平策略。但有一个关键区别:Semaphore 使用的是 AQS 的 共享模式(Shared Mode),而非独占模式。 ...

八月 26, 2022 · 10 分钟 · 2051 字 · yaomingye

CountDownLatch 设计思想解析:一次性的共享门闩为什么这样设计

CountDownLatch:一次性的共享门闩 🤔 道格·李为什么需要一个倒计时门闩 多线程编程里有一个反复出现的问题:主线程需要等待若干个子任务全部完成,然后汇总结果继续执行。在 JUC 出现之前,Java 只有两种方式应对: Thread.join()——但 join() 等的是线程终止,不是任务完成。如果线程来自线程池(被复用,不会终止),join() 完全不适用。它是"等人死了"而不是"等人把活干完"。 自己写自旋等待——用一个 volatile 计数器,主线程循环检查。但自旋空转浪费 CPU,加 sleep 又会引入延迟,而且 count++ 本身不是原子操作。 道格·李在设计 JSR 166 时看到了这个空缺:需要一个轻量级的同步辅助工具,让一个(或多个)线程能够等待其他线程完成一组操作,不依赖线程终止,不浪费 CPU,而且足够简单。 这就是 CountDownLatch 的诞生背景。它用 AQS 的共享模式实现了一个一次性的倒计时器:计数器从 N 开始,每个子任务完成时减 1(countDown()),主线程在 await() 上阻塞直到计数器归零。 道格·李给它设计了几个刻意的约束: 一次性——计数器归零后无法重置。这个约束简化了实现(不需要考虑"重置时正在等待的线程怎么办"),也迫使使用者为可重复场景选用 CyclicBarrier 只减不增——countDown() 只能减少计数,无法增加。这也简化了状态机——计数器状态只有 N→0 这一个方向 基于 AQS 共享模式——让多个等待线程可以同时被唤醒(当计数器归零时),而不是排他锁那样只唤醒一个 设计理念 1:一次性的约束 CountDownLatch 最核心的设计约束是 一次性 ——一旦计数器从 N 减到 0,门闩永久打开,无法再关闭。 这个约束绝非"能力不足",而是 刻意的设计取舍 : 如果支持重置 代价 需要处理"已有线程在 await 上等待"和"重置后的新等待者"两种状态 状态机复杂度翻倍 countDown 和 reset 并发时的语义难以定义 需要额外的同步机制 每个等待者都要知道自己是"老批次"还是"新批次" 需要代际(generation)标记 一次性约束消除了一个巨大的设计空间:时间维度上的状态管理。 CountDownLatch 只有两个有意义的状态: state > 0 (门闩关闭)和 state == 0 (门闩打开)。从关闭到打开只需要 state 单向递减,不需要考虑"打开了又关上"的复杂路径。 ...

八月 26, 2022 · 5 分钟 · 986 字 · yaomingye
Cat Radio