JUC 在中间件中的应用

JUC 在中间件中的应用:线程池与并发集合实战全景 问题切入:道格·李的组件在中间件里是如何落地的 道格·李设计的每一个 JUC 组件都有明确的定位:ThreadPoolExecutor 管理线程资源、ConcurrentHashMap 提供高并发下的安全容器、BlockingQueue 协调生产者与消费者。但这些组件本身只是"积木"——积木搭成什么,看用的人。 Tomcat、Netty、Dubbo、RocketMQ 这些中间件的作者,就是最高水平的积木搭手。他们在道格·李提供的基础上做了大量二次定制:继承 ThreadPoolExecutor 改写拒绝策略、用 ConcurrentHashMap 存储单例对象、用 BlockingQueue 实现异步日志缓冲。 翻开这些中间件的源码,你会发现:标准 JUC 组件很少被直接使用,几乎都被继承或组合包装。这不是因为标准组件不够好,而是因为每个中间件的场景都有自己的约束——Tomcat 的线程池需要在队列满时反过来创建线程(而不是拒绝),Netty 用 NioEventLoopGroup 把线程池拆成了事件循环。 本篇从源码层面逐一拆解道格·李的 JUC 积木如何在中间件中被定制、组合和落地。覆盖的中间件和对应的 JUC 组件如下: flowchart LR classDef root fill:#0f172a,stroke:#3b82f6,stroke-width:2px,color:#bfdbfe,font-weight:bold; classDef branch fill:#2d1a05,stroke:#f59e0b,stroke-width:2px,color:#fde68a,font-weight:bold; classDef leaf fill:#1e1e24,stroke:#6b7280,stroke-width:1.5px,color:#e5e7eb; classDef highlight fill:#450a0a,stroke:#dc2626,stroke-width:1.5px,color:#fecaca,font-weight:bold; ROOT[JUC 在中间件中的应用全景] ROOT --> THREAD["线程池 ThreadPoolExecutor"] THREAD --> T1["Tomcat: 请求处理线程池\n自定义 TaskQueue 配合拒绝策略"] THREAD --> T2["Netty: NioEventLoopGroup\nSingleThreadEventExecutor 模型"] THREAD --> T3["Dubbo: 多种线程池策略\nFixed/Cached/Limited/Eager"] THREAD --> T4["RocketMQ: Broker 线程池组\nSendMessage/PullMessage 等"] ROOT --> MAP["ConcurrentHashMap"] MAP --> M1["Spring IOC: singletonObjects\n所有单例 Bean 的存储容器"] MAP --> M2["Netty: DefaultChannelHandlerContext\nChannel 属性存储"] MAP --> M3["Tomcat: Servlet 映射表\nURL → Servlet 的路由缓存"] ROOT --> QUEUE["BlockingQueue"] QUEUE --> Q1["Logback: AsyncAppender\nArrayBlockingQueue 异步写日志"] QUEUE --> Q2["Tomcat: TaskQueue\n继承 LinkedBlockingQueue"] QUEUE --> Q3["Disruptor: RingBuffer\n虽非JUC但思想同源"] ROOT --> LIST["CopyOnWriteArrayList"] LIST --> L1["Tomcat: Session 监听器列表\n遍历时无需加锁"] LIST --> L2["Spring: ApplicationListener 集合\n事件多播时安全迭代"] class ROOT root; class THREAD,MAP,QUEUE,LIST branch; class T1,T2,T3,T4,M1,M2,M3,Q1,Q2,Q3,L1,L2 leaf; class T1,T2,M1,Q1,L1 highlight; 🏊 线程池在中间件中的应用 🏊 Tomcat:请求处理的线程池引擎 Tomcat 处理 HTTP 请求的核心是一个定制化的 ThreadPoolExecutor 。它没有直接用 JDK 的标准实现,而是继承了 ThreadPoolExecutor 并重写了其中的关键行为。 ...

十月 2, 2022 · 9 分钟 · 1834 字 · yaomingye

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

JUC 源码阅读路线图

JUC 源码阅读路线图:从 LockSupport 到 ForkJoinPool 的完整导读 ❓ 1️⃣ 一、道格·李的 JUC 有明确的分层设计——读源码必须按这个顺序 道格·李在设计 java.util.concurrent 包时,不是把二十几个类平铺在一个包里的。JUC 有严格的分层:底层原语(CAS、volatile、LockSupport)→ 核心框架(AQS)→ 具体实现(锁、同步器、集合、线程池)。 这个分层意味着:如果你一上来就读 ReentrantLock.lock() 的源码,三行之后就会遇到 tryAcquire(),再往下就是 CAS 修改 AQS state、LockSupport.park() 阻塞线程——全是底层 API。不知道 CAS 的语义,不知道 state 的 CLH 入队流程,不知道 park/unpark 的 permit 机制,每走一步都得暂停查资料,阅读体验极差。 反之,如果你按道格·李的设计顺序来读——先理解 CAS 和 LockSupport(地基),再啃透 AQS(骨架),然后逐一看 ReentrantLock、Semaphore、CountDownLatch 怎么在骨架上加肉(定制 tryAcquire / tryRelease)——整个 JUC 包的结构就一目了然了。 这篇博客提供的就是这份按设计分层排列的源码阅读路线图——告诉你每个组件在 JDK 中的位置、入口 API、核心函数调用链、推荐阅读顺序、以及读完这个组件你能学到什么设计思想。不会逐行分析源码(各组件详细分析见本系列其他文章)。 🏗️ 2️⃣ 二、总览:JUC 全景架构图 阅读之前,先建立全局坐标系。以下是所有 JUC 核心组件的逻辑关系图: flowchart TD classDef base fill:#450a0a,stroke:#dc2626,stroke-width:2px,color:#fecaca,font-weight:bold; classDef core fill:#1e1b4b,stroke:#4f46e5,stroke-width:2px,color:#e0e7ff,font-weight:bold; classDef lock fill:#1e1b4b,stroke:#4f46e5,stroke-width:2px,color:#e0e7ff,font-weight:bold; classDef sync fill:#2a1147,stroke:#a855f7,stroke-width:1.5px,color:#ede9fe,font-weight:bold; classDef coll fill:#052e16,stroke:#16a34a,stroke-width:1.5px,color:#bbf7d0,font-weight:bold; classDef pool fill:#0f172a,stroke:#3b82f6,stroke-width:1.5px,color:#bfdbfe,font-weight:bold; classDef tl fill:#1e1e24,stroke:#6b7280,stroke-width:1.5px,color:#e5e7eb; BASE[🔧 底层原语] BASE --> CAS[CAS\nUnsafe.compareAndSwapX] BASE --> VOL[volatile\n内存可见性] BASE --> PARK[LockSupport\npark/unpark] PARK --> AQS[AbstractQueuedSynchronizer\nAQS框架] CAS --> AQS VOL --> AQS AQS --> RL[ReentrantLock] AQS --> RW[ReentrantReadWriteLock] AQS --> SEM[Semaphore] AQS --> CDL[CountDownLatch] AQS --> FUT[FutureTask] AQS --> TPE[ThreadPoolExecutor] PARK --> CB[CyclicBarrier] CAS --> CHM[ConcurrentHashMap] CAS --> CLQ[ConcurrentLinkedQueue] VOL --> CHM RL --> COW[CopyOnWriteArrayList] AQS -.-> BQ[BlockingQueue] TPE --> STPE[ScheduledThreadPoolExecutor] TPE -.-> FJP[ForkJoinPool] CB --> TL[ThreadLocal] TL --> ITL[InheritableThreadLocal] ITL --> TTL_CLASS[TransmittableThreadLocal] class BASE base; class AQS core; class RL,RW lock; class SEM,CDL,CB,FUT sync; class CHM,CLQ,COW,BQ coll; class TPE,STPE,FJP pool; class TL,ITL,TTL_CLASS tl; 这张图揭示了 JUC 的设计分层: ...

九月 6, 2022 · 12 分钟 · 2397 字 · 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

ThreadLocal 线程池上下文传递

ThreadLocal 线程池上下文传递:从 InheritableThreadLocal 缺陷到 TransmittableThreadLocal 全解析 🤔 一、JDK 设计者为什么要给每个线程配一个"私房钱罐" 多线程编程中有一个经典矛盾:线程之间要共享一部分数据来协作,又要有各自私有的数据来隔离。共享数据靠锁来保护,私有数据呢?如果每建一个新线程都要手动传参数、写包装类,代码很快就变成意大利面条。 JDK 1.2 的设计者(Josh Bloch 等人)给出的方案是 ThreadLocal——每个线程维护一个私有的 ThreadLocalMap,key 是 ThreadLocal 实例,value 是你想隔离的数据。同一个 ThreadLocal 对象在不同线程中的值互不干扰。这个设计让"线程级上下文"(traceId、事务、用户 Session)变得自然:只要在主线程 set 一下,当前线程的任何方法都能 get 到,不需要在方法签名里一路传参。 但 JDK 设计者很快发现一个新问题:ThreadLocal 在线程间是完全隔离的——如果父线程 set 了值,新建子线程时,子线程拿不到。这就是为什么后来又有了 InheritableThreadLocal:它在 Thread 构造函数中触发 init(),将父线程 ThreadLocalMap 中标记为可继承的条目浅拷贝到子线程。 然而,ITL 的设计有一个致命缺陷——它只在 new Thread() 时触发传递。线程池复用已有线程,不再走 Thread 构造函数,ITL 的传递逻辑完全不执行。第一次提交任务时碰巧用的是刚创建的新线程(触发了一次传递),第二次复用同一个线程时,父线程的新值就传不过来了。 这就是阿里开源的 TransmittableThreadLocal 要解决的问题。它的核心思路是:不再依赖线程创建时的一次性传递,而是在每次提交任务时主动 capture 父线程的快照 → replay 到工作线程 → 任务完成后再 restore 还原。 阅读本篇文章的收获: InheritableThreadLocal 是如何在 Thread 构造函数中实现传递的?源码在哪一行触发? 为什么它在线程池中会失效?根源在 JDK 源码的哪一行? TransmittableThreadLocal 又是如何在源码层面解决这些缺陷的? capture() / replay() / restore() 三个方法各自做了什么? 🧵 二、ThreadLocal 基础回顾:数据到底存在哪里 在深入 ITL 和 TTL 之前,先快速回顾 ThreadLocal 的核心数据结构。如果你已经熟悉这部分,可以直接跳到第三章。 ...

九月 5, 2022 · 13 分钟 · 2586 字 · yaomingye

ScheduledThreadPoolExecutor 定时调度增强

ScheduledThreadPoolExecutor 定时调度增强:DelayedWorkQueue 二叉堆延时队列与 Spring 体系实战 🚀 道格·李为什么需要一个能定时的线程池 Java 1.3 引入的 java.util.Timer 是 JDK 最早提供的定时任务工具。但它有两个致命设计缺陷:① 单线程执行——一个任务执行时间过长,后面的所有任务都会延迟;② 异常吞没导致线程终止——任务抛了未捕获异常,Timer 线程静默死亡,剩余任务永远不会执行。第二个问题在生产环境尤其危险——线上定时取消超时订单的任务因为一个 NullPointerException 静默停止,几天后才被发现。 道格·李在设计 JSR 166 时,ScheduledThreadPoolExecutor 是 ThreadPoolExecutor 的直接扩展。它的设计策略是复用线程池的全部管理能力(线程生命周期、拒绝策略、钩子方法),只替换两个关键组件: 任务队列:用 DelayedWorkQueue(基于二叉堆的延时队列)替换 BlockingQueue,任务按触发时间排序,堆顶是最先到期的任务 任务类型:用 ScheduledFutureTask 替换普通的 FutureTask,增加了周期执行模式(固定频率 vs 固定延迟)和下次触发时间的计算逻辑 核心改进:线程池里有 N 个工作线程,一个任务异常不会影响其他线程和任务。Timer 的单线程弱点不再存在。 🏊 ScheduledThreadPoolExecutor 的整体架构 🔗 继承关系与组件概览 ScheduledThreadPoolExecutor 直接继承 ThreadPoolExecutor,在父类基础上替换了三个关键组件: flowchart LR %% ========================================== %% 样式定义 %% ========================================== classDef root fill:#0f172a,stroke:#3b82f6,stroke-width:2px,color:#bfdbfe,font-weight:bold; classDef branch fill:#2d1a05,stroke:#f59e0b,stroke-width:2px,color:#fde68a,font-weight:bold; classDef leaf fill:#1e1e24,stroke:#6b7280,stroke-width:1.5px,color:#e5e7eb; classDef highlight fill:#450a0a,stroke:#dc2626,stroke-width:1.5px,color:#fecaca,font-weight:bold; ROOT[ScheduledThreadPoolExecutor\n继承 ThreadPoolExecutor] ROOT --> B1(1. 任务类型替换) B1 --> TASK["📦 ScheduledFutureTask\nextends FutureTask\n+ implements Delayed\n+ 三态 period 模型\n+ sequenceNumber 保序"] ROOT --> B2(2. 队列替换) B2 --> QUEUE["📥 DelayedWorkQueue\n自建二叉堆\n无界阻塞延迟队列\n扩展 leader/follower 模式"] ROOT --> B3(3. 调度方法替代) B3 --> SCHED["⚡ 三个入口方法"] SCHED --> S1["schedule()\n一次性延迟任务\nperiod = 0"] SCHED --> S2["scheduleAtFixedRate()\n固定速率\nperiod > 0"] SCHED --> S3["scheduleWithFixedDelay()\n固定延迟\nperiod < 0"] ROOT --> B4(4. 关闭后行为) B4 --> SHUT["🛑 两个布尔开关"] SHUT --> C1["continueExistingPeriodicTasksAfterShutdown\nshutdown 后是否继续执行周期任务"] SHUT --> C2["executeExistingDelayedTasksAfterShutdown\nshutdown 后是否执行已延迟的任务"] ROOT --> B5(5. 线程数策略) B5 --> SIZE["🔢 maximumPoolSize = Integer.MAX_VALUE\n队列无界,永远不需要额外线程\n仅核心线程数决定并发度"] class ROOT root; class B1,B2,B3,B4,B5 branch; class TASK,QUEUE,SCHED,SHUT,SIZE leaf; class S1,S2,S3,C1,C2 highlight; 五个改进点一句话总结 : ...

九月 4, 2022 · 12 分钟 · 2502 字 · 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

ThreadPoolExecutor 源码解析

ThreadPoolExecutor 源码解析:Worker 机制、生命周期、拒绝策略与动态线程池实践 🚀 道格·李为什么需要一个线程池 Java 1.0 就支持多线程,但管理线程生命周期这件事一直缺少标准方案。开发者每次需要异步执行时,要么 new Thread().start(),要么自己维护一个线程管理队列——前者浪费资源,后者极易出错。 一个线程的创建和销毁是有成本的。JVM 要为每个线程分配栈内存(默认约 1MB),操作系统要为每个线程维护内核线程表项和调度上下文。当并发请求量上来后,频繁创建/销毁线程会导致: 内存压力——大量线程的栈内存吃掉堆外空间 CPU 浪费在上下文切换——线程数远超 CPU 核心数时,CPU 的时间片都消耗在"换人"而不是"干活"上 线程数不可控——请求峰值时线程数无上限增长,最终 OOM 或系统不可用 道格·李在设计 JSR 166 时面对的核心问题是:如何让开发者既能享受多线程的并发收益,又不用直接管理线程的创建和销毁? 答案是把线程抽象为一种可复用的资源——线程池。 线程池的本质是一个"线程 + 任务队列"的组合:核心线程常驻,任务多时创建临时线程分担,任务少时回收空闲线程,任务太多时由拒绝策略兜底。从设计上看,ThreadPoolExecutor 把线程的创建策略(core/max)、存活策略(keepAliveTime)、排队策略(workQueue)和过载策略(rejectedExecutionHandler)全部暴露为可配置参数——这正是道格·李的设计风格:不替开发者做决定,而是把决策权交给调用方。 📐 七大核心参数 ThreadPoolExecutor 最完整的构造器接受 7 个参数: public ThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, TimeUnit unit, BlockingQueue<Runnable> workQueue, ThreadFactory threadFactory, RejectedExecutionHandler handler) ⚙️ 1. corePoolSize — 核心线程数 线程池中始终存活的线程数量(除非 allowCoreThreadTimeOut 设为 true)。即使这些线程当前空闲,也不会被回收。 关键行为:当提交任务时,即使有空闲的核心线程,只要当前线程数少于 corePoolSize,线程池也会继续创建新的线程——先凑够核心线程数量,再谈复用。这种"先扩容再复用"是出于设计上的简单性:判断是否达到核心线程数的开销远小于判断是否有空闲线程且空闲线程是否可用。 📐 2. maximumPoolSize — 最大线程数 线程池允许创建的最大线程数。只有当工作队列已满且当前线程数不足 maximumPoolSize 时,才会创建超出核心线程数的额外线程。 ...

九月 2, 2022 · 20 分钟 · 4131 字 · yaomingye

ForkJoinPool 源码深度解析

ForkJoinPool 源码深度解析:从分治思想到工作窃取的完整实现 🤔 一、道格·李为什么需要一个"能偷工作"的线程池 Java 5 的 ThreadPoolExecutor 解决了线程复用的问题,但它有一个结构性的局限:所有线程共享一个任务队列。当一个线程提交了子任务后阻塞等待子任务结果,而子任务又在同一个队列里等待被执行时,就会发生线程饥饿——等待的线程占着一个槽位但不干活,队列里的子任务没人执行,形成死锁。 这个问题在递归分治算法(把大问题拆成小问题递归求解)中尤为致命。分治算法天然适合并行——子问题之间互不依赖,可以同时计算。但如果每个线程都把子任务扔到共享队列然后等结果,队列很快就会堆满等待被执行的任务而所有线程都在等。 道格·李在 Java 7 中引入 ForkJoinPool 时,核心创新是工作窃取(Work-Stealing): 每个工作线程有自己的双端队列(Deque),线程从自己的队列头部取任务 当一个线程 fork 子任务时,子任务被 push 到该线程自己的队列 当线程自己的队列空了,它会从其他线程的队列尾部窃取任务来执行 这个设计解决了两个问题:① 递归 fork 的子任务不会堵塞共享队列;② 快线程不会空等——它会偷慢线程的活来干。ForkJoinPool 也是 Java 8 并行流(parallelStream())的底层引擎。 二、设计思想:分治算法 + 工作窃取 📌 2.1 分治算法(Divide-and-Conquer) ForkJoinPool 的设计基础是分治算法(Divide-and-Conquer),其核心过程为三个方面: flowchart TD %% 半暗底色 + 高亮描边:完美适配博客深色/浅色双主题 %% classDef process fill:#1e1e24,stroke:#6b7280,stroke-width:2px,color:#e5e7eb; classDef root fill:#0f172a,stroke:#3b82f6,stroke-width:2.5px,color:#bfdbfe,font-weight:bold; ROOT[根任务\n问题规模N] ROOT -->|拆分| L1[子任务\n规模N/2] ROOT -->|拆分| R1[子任务\n规模N/2] L1 -->|继续拆分| L2[子任务\n规模N/4] L1 -->|继续拆分| R2[子任务\n规模N/4] R1 -->|继续拆分| L3[子任务\n规模N/4] R1 -->|继续拆分| R3[子任务\n规模N/4] L2 -->|达到阈值\n直接计算| SOLVE1[原子任务] R2 -->|达到阈值\n直接计算| SOLVE2[原子任务] L3 -->|达到阈值\n直接计算| SOLVE3[原子任务] R3 -->|达到阈值\n直接计算| SOLVE4[原子任务] SOLVE1 -->|合并| MERGE1[汇总结果] SOLVE2 -->|合并| MERGE1 SOLVE3 -->|合并| MERGE2[汇总结果] SOLVE4 -->|合并| MERGE2 MERGE1 -->|最终合并| FINAL[最终结果] MERGE2 -->|最终合并| FINAL class FINAL,L1,L2,L3,MERGE1,MERGE2,R1,R2,R3,SOLVE1,SOLVE2,SOLVE3,SOLVE4 process; class ROOT root; 阶段 操作 说明 Divide(拆分) fork() 将大任务递归拆分为小任务,直到达到阈值 Conquer(求解) compute() 对原子任务执行实际计算 Combine(合并) join() 递归汇总所有子任务的结果 📌 2.2 工作窃取算法(Work-Stealing) 普通的 ThreadPoolExecutor 使用单一共享阻塞队列(BlockingQueue),所有线程竞争同一个队列的头元素,存在单点竞争瓶颈。ForkJoinPool 则采用完全不同的设计:每个工作线程维护自己的双端队列(WorkQueue)。 ...

九月 2, 2022 · 20 分钟 · 4180 字 · yaomingye

BlockingQueue 设计解析:四组方法语义、锁机制分化与生产者-消费者模型的工程实践

BlockingQueue 设计解析 🤔 道格·李为什么需要一个阻塞队列接口 生产者-消费者模式是多线程编程里最常见的协作模型——一个(或多个)线程生产数据,另一个(或多个)线程消费数据。在 JUC 出现之前,Java 开发者只能用 wait() / notify() 手写这个模型。 手写版本的典型代码如下: public synchronized void put(E e) throws InterruptedException { while (list.size() == capacity) { wait(); } list.addLast(e); notifyAll(); } 这段代码表面正确,但道格·李在分析并发程序的常见错误时发现了几个根深蒂固的问题: 生产者唤醒生产者:notifyAll() 唤醒等待队列里的所有线程——包括生产者和消费者。当队列满时,多个生产者同时被唤醒,只有第一个能成功插入,其余又回到 wait。这些"无效唤醒"不是 Bug,但大量浪费 CPU 无法区分等待原因:所有线程在同一个条件队列上等待,生产者因为"队列满"而等,消费者因为"队列空"而等。notifyAll() 叫醒所有人,但被叫醒的线程可能发现条件仍不满足,继续睡——这就是为什么 wait() 必须放在 while 循环里 没有标准接口:每个项目都在重新发明这个轮子,而且各自的行为语义不一致——有的用 null 表示失败,有的抛异常,有的阻塞等待 道格·李的解决方案是两层的:接口层——BlockingQueue 接口定义了四组标准方法(抛异常、返回特殊值、阻塞、超时),统一了所有阻塞队列的行为契约。实现层——用 ReentrantLock 的两个 Condition(notFull 和 notEmpty)精确分离生产者与消费者的等待条件,让"队列满"只唤醒消费者,“队列空"只唤醒生产者,消除无效唤醒。 🚧 BlockingQueue 接口设计:四组方法的语义定义 BlockingQueue 接口最核心的设计决策在于: 同一操作提供四种不同的线程协作策略 ,以方法名区分行为,以返回类型区分语义。 行为模式 插入 移除 检查 语义 抛异常 add(e) remove() element() 操作无法立即执行时抛出 IllegalStateException ,调用方需自行处理 返回特殊值 offer(e) poll() peek() 操作无法立即执行时返回 false 或 null ,调用方通过返回值判断是否成功 阻塞 put(e) take() — 操作无法立即执行时阻塞当前线程,直到条件满足,调用方被挂起 超时 offer(e, t, u) poll(t, u) — 操作无法立即执行时阻塞最多指定时长,超时返回 false 或 null 四组方法的核心设计哲学是"让调用方选择等待策略而非被动接受”。 同一个"放入元素"的需求,调用方可以根据业务场景选择: ...

八月 31, 2022 · 9 分钟 · 1713 字 · yaomingye
Cat Radio