并发 11 · 同步工具类
秒杀活动零点开始之前系统通常要做一轮预热库存服务要把商品库存加载进 Redis优惠券服务要初始化好发券规则风控服务要拉取好黑名单——主线程必须等这几个子系统都准备好才能把允许下单的开关打开。这是一个典型的等待多个任务全部完成的协作场景靠前面几篇讲的锁或者BlockingQueue都不太顺手因为这里的诉求不是互斥、也不是传递数据而是线程之间的节奏协调。JDK 在java.util.concurrent包里专门准备了一组同步工具类来解决这类节奏协调问题CountDownLatch、CyclicBarrier、Semaphore、Phaser、Exchanger。它们和第 5、6 篇讲的锁不是一回事——锁解决的是同一时刻谁能进临界区这几个工具解决的是线程之间谁等谁、等到什么时候一起走。这一篇把它们挨个拆开看重点关注每一个工具适合什么场景、和相邻工具的边界在哪里。目录CountDownLatch一次性的门闩CyclicBarrier可重复使用的栅栏两者的边界门闩等外部事件栅栏等彼此到达Semaphore控制并发访问数量的信号量Phaser更灵活的多阶段协作简要了解Exchanger两个线程之间的数据交换实战秒杀系统预热分片处理限流的组合选型小结一、CountDownLatch一次性的门闩CountDownLatch倒计时门闩的用法非常直白用一个计数器初始化多个线程各自完成任务后调用countDown()让计数器减一等待的线程调用await()阻塞直到计数器归零才被放行。publicclassSeckillWarmup{publicstaticvoidmain(String[]args)throwsInterruptedException{CountDownLatchlatchnewCountDownLatch(3);// 3 个子系统要初始化AtomicReferenceThrowablefailurenewAtomicReference();newThread(()-{try{loadStockToRedis();}catch(Throwablet){failure.compareAndSet(null,t);}finally{latch.countDown();// 成功失败都必须减一否则等待方可能永久阻塞}}).start();newThread(()-{try{initCouponRules();}catch(Throwablet){failure.compareAndSet(null,t);}finally{latch.countDown();}}).start();newThread(()-{try{loadRiskBlacklist();}catch(Throwablet){failure.compareAndSet(null,t);}finally{latch.countDown();}}).start();latch.await();// 主线程阻塞直到计数器归零if(failure.get()!null){thrownewIllegalStateException(预热失败,failure.get());}System.out.println(预热完成开闸放行);openSeckillGate();}}底层实现直接复用了第 5 篇讲过的AQS 共享模式内部维护一个Sync extends AbstractQueuedSynchronizerstate就是那个计数器。countDown()本质是releaseShared()每次把state减一减到 0 时唤醒所有等待线程await()本质是acquireSharedInterruptibly()state不为 0 就把当前线程挂进 AQS 的等待队列。这也解释了为什么它天然支持多个线程同时await()——共享模式下一个唤醒动作可以级联唤醒队列里所有等待的节点而不像独占模式一次只放一个。CountDownLatch有个容易被忽略但很关键的限制计数器只能减、不能重置减到 0 之后这个门闩就永久保持开放后续await()会直接通过想开始一轮新的计数必须重新new一个。这决定了它适合的是一次性的等待场景——像上面这种应用启动预热天然只需要发生一次。示例里还有一条生产代码必须守住的规则countDown()要放在finally中同时用单独的失败通道把异常传回等待方。只放finally能避免永久等待但如果不记录失败主线程仍可能误把三个任务都结束了当成三个任务都成功了。二、CyclicBarrier可重复使用的栅栏如果预热逻辑不是启动一次而是要反复出现——比如秒杀系统把海量的扣库存请求分成一批一批处理每一批里的所有线程都必须先处理完当前这批才能一起进入下一批常见于分片计算、批量任务分阶段执行——这种周期性汇合点就该用CyclicBarrier循环栅栏了。publicclassBatchStockDeduction{publicstaticvoidmain(String[]args){intworkerCount4;CyclicBarrierbarriernewCyclicBarrier(workerCount,()-System.out.println(本批次全部处理完毕进入下一批));// 栅栏放开时执行一次AtomicBooleancancellednewAtomicBoolean();for(inti0;iworkerCount;i){intworkerIdi;newThread(()-{for(intbatch0;batch10!cancelled.get();batch){try{processBatch(workerId,batch);// 处理当前批次分配到的任务if(cancelled.get())break;barrier.await();// 等其余3个线程也处理完这一批}catch(InterruptedExceptione){cancelled.set(true);barrier.reset();Thread.currentThread().interrupt();break;}catch(BrokenBarrierExceptione){cancelled.set(true);break;// 栅栏已经损坏继续 await 只会再次失败}catch(RuntimeExceptione){cancelled.set(true);barrier.reset();// 自己提前退出时唤醒正在等这一批的同伴throwe;}}}).start();}}}关键区别在实现上CyclicBarrier不是基于 AQS而是靠一把ReentrantLock配合Condition手写的回顾第 6 篇。每个线程调用await()时先lock()把计数减一如果减到 0 就执行可选的Runnable就是上面例子里那个进入下一批的打印然后condition.signalAll()唤醒所有等待者同时把计数器重置回初始值——这就是循环Cyclic的含义用完一次自动复位可以无限次重复使用。如果计数没减到 0当前线程就condition.await()挂起等待。三、两者的边界门闩等外部事件栅栏等彼此到达CountDownLatch和CyclicBarrier长得很像——都是等到某个数量才放行——但语义上有一条清晰的分界线容易在面试或者选型时搞混维度CountDownLatchCyclicBarrier等待的对象一个或多个外部事件完成一组彼此协作的线程互相到达同一个点谁调用哪个方法完成任务的线程调countDown()等待的线程调await()——两者可以是完全不同的线程集合所有参与协作的线程都调用同一个await()——等待者和被等待者是同一批线程能否重复使用不能减到 0 即失效能放行后自动重置计数到达后的动作无内置回调可传入一个Runnable最后一个到达的线程触发在放行前执行用一句话记门闩Latch适合1 个或多个线程等另外几个线程把活干完如主线程等子系统初始化栅栏Barrier适合一组同伴互相等对方走到同一步再一起往下走如批量任务分阶段推进。秒杀系统的启动预热天然是前者分片批处理天然是后者选错了工具虽然能通过一些变通用法凑合实现但语义会很扭曲。四、Semaphore控制并发访问数量的信号量前两个工具解决的是等待Semaphore信号量解决的是另一类问题限制同一时刻能有多少个线程同时访问某个资源。它的模型很直观维护一定数量的许可permit线程访问资源前先acquire()拿一个许可没有许可就阻塞等待用完之后release()归还许可。publicclassSeckillRateLimiter{// 假设数据库连接池只能扛住 20 个并发扣库存操作用信号量限流保护它privatefinalSemaphoresemaphorenewSemaphore(20);publicvoiddeductStock(LongproductId){try{semaphore.acquire();// 拿不到许可就在这里排队等}catch(InterruptedExceptione){Thread.currentThread().interrupt();return;// 没拿到许可不能执行 release()}try{doDeductStock(productId);// 真正操作数据库}finally{semaphore.release();// 只归还已经成功拿到的许可}}}底层同样是基于 AQS 的共享模式和CountDownLatch用的是同一套家族state表示当前剩余的许可数量。acquire()是tryAcquireShared()失败则入队等待release()是增加许可后唤醒等待队列。release()必须和一次已经成功完成的acquire()配对。若线程在等待许可时被中断acquire()会直接抛异常此时不能进入释放逻辑否则会凭空增加许可、悄悄放宽并发上限。一个常被拿来类比又容易搞混的点如果把许可数设成 1Semaphore看起来就像一把互斥锁但和ReentrantLock有本质区别——Semaphore没有持有者概念不具备可重入性A 线程acquire()拿到的许可可以由 B 线程release()归还这在锁的语义里是不允许的锁必须由持有它的线程自己释放。这个特性反过来也是它的价值所在Semaphore表达的从来不是互斥关系而是资源配额——无论是限制数据库连接数、限制同时下载的文件数还是给下游接口做限流保护本质都是这批资源总量有限谁先申请到谁先用用完还回去这跟谁拿到锁就必须谁释放是两种不同的心智模型。五、Phaser更灵活的多阶段协作简要了解JDK 7 引入的Phaser可以看作CyclicBarrier的加强版解决的是同一类多阶段协作问题但补上了CyclicBarrier的两个短板参与者数量可以动态增减。CyclicBarrier的参与线程数在构造时就固定死了Phaser支持运行期通过register()/deregister()动态调整参与者数量——比如一个批处理任务中途可能有新的工作线程加入或者提前退出。支持分层多级 Phaser 树。大量参与者时单个Phaser内部同步的开销会变高可以把参与者组织成一棵树每层内部先同步再逐层向上汇聚降低整体的同步竞争。代价是 API 比CyclicBarrier复杂不少arrive()、arriveAndAwaitAdvance()、awaitAdvance()等一堆方法还要理解阶段号这个概念。实际业务代码里Phaser出现的频率并不高——多数分阶段同步的需求用CyclicBarrier就够了只有真的遇到参与者数量会动态变化的场景才值得多学一层Phaser的复杂度。这里知道它存在、明白它解决什么问题即可不做源码级展开。六、Exchanger两个线程之间的数据交换Exchanger是这几个工具里最小众的一个只解决一件很窄的事让恰好两个线程在同一个点上各自把手里的数据交给对方。ExchangerListOrderexchangernewExchanger();// 线程A填充数据的一方newThread(()-{ListOrderbufferfillOrderBuffer();try{ListOrderpartnerBufferexchanger.exchange(buffer);// 等待线程B交换后拿到对方的数据process(partnerBuffer);}catch(InterruptedExceptione){Thread.currentThread().interrupt();}}).start();// 线程B另一方newThread(()-{ListOrderbufferfillAnotherBuffer();try{ListOrderpartnerBufferexchanger.exchange(buffer);// 等待线程A交换后拿到对方的数据process(partnerBuffer);}catch(InterruptedExceptione){Thread.currentThread().interrupt();}}).start();exchange()是双向阻塞的——先到的线程会一直等直到另一个线程也调用exchange()两者才同时完成交换并各自继续往下走这跟第 10 篇讲的SynchronousQueue的面对面交接神似实际上SynchronousQueue内部的等待/传递机制和Exchanger的实现思路一脉相承。它适合的是那种两个角色分工明确、需要周期性互换缓冲区的场景比如一个线程专门负责填充数据、另一个线程专门处理数据用交换而不是加锁复制的方式减少数据拷贝开销。这个场景在日常业务代码里确实不算常见了解思路即可。七、实战秒杀系统预热分片处理限流的组合把本篇讲的三个主力工具CountDownLatch、CyclicBarrier、Semaphore串起来看一个更完整的秒杀系统骨架publicclassSeckillSystem{privatefinalSemaphoredbSemaphorenewSemaphore(20);// 保护数据库的并发限流publicvoidstart()throwsInterruptedException{// 第一阶段CountDownLatch —— 等待所有子系统初始化完成CountDownLatchwarmupLatchnewCountDownLatch(3);AtomicReferenceThrowablewarmupFailurenewAtomicReference();ExecutorServicewarmupPoolExecutors.newFixedThreadPool(3);submitWarmup(warmupPool,warmupLatch,warmupFailure,this::loadStockToRedis);submitWarmup(warmupPool,warmupLatch,warmupFailure,this::initCouponRules);submitWarmup(warmupPool,warmupLatch,warmupFailure,this::loadRiskBlacklist);try{warmupLatch.await();}catch(InterruptedExceptione){warmupPool.shutdownNow();throwe;}finally{warmupPool.shutdown();}if(warmupFailure.get()!null){thrownewIllegalStateException(系统预热失败,warmupFailure.get());}System.out.println(预热完成开始处理请求);// 第二阶段CyclicBarrier —— 分批处理请求每批同步一次intworkerCount4;CyclicBarrierbatchBarriernewCyclicBarrier(workerCount,()-System.out.println(本批处理完毕统计当前剩余库存));AtomicBooleancancellednewAtomicBoolean();for(inti0;iworkerCount;i){newThread(()-{for(intbatch0;batch100!cancelled.get();batch){// 第三阶段Semaphore —— 每个具体的扣库存操作前先限流try{dbSemaphore.acquire();try{deductStockForBatch(batch);}finally{dbSemaphore.release();// 只在 acquire 成功后进入这层 finally}if(cancelled.get())break;batchBarrier.await();// 等其他worker也处理完这一批再一起进下一批}catch(InterruptedExceptione){cancelled.set(true);batchBarrier.reset();// 若中断发生在 acquire 阶段也要放开其他等待者Thread.currentThread().interrupt();break;}catch(BrokenBarrierExceptione){cancelled.set(true);break;}catch(RuntimeExceptione){cancelled.set(true);batchBarrier.reset();// 业务异常退出时避免同伴永久卡在栅栏throwe;}}}).start();}}privatevoidsubmitWarmup(ExecutorServicepool,CountDownLatchlatch,AtomicReferenceThrowablefailure,Runnabletask){pool.execute(()-{try{task.run();}catch(Throwablet){failure.compareAndSet(null,t);}finally{latch.countDown();}});}}三个工具在这里各司其职CountDownLatch管启动前的一次性等待CyclicBarrier管处理过程中反复出现的阶段汇合Semaphore管任意时刻别让并发量把数据库压垮——它们互不冲突可以在同一个系统里叠着用各自解决自己那一层的协调问题。八、选型小结场景工具关键理由等待多个一次性任务全部完成再继续CountDownLatch一次性门闩基于AQS共享模式减到0永久放行一组线程反复地互相等对方到达再一起走CyclicBarrier可重复使用锁Condition实现支持到达后回调限制同一时刻并发访问某资源的线程数量Semaphore许可无持有者概念表达资源配额而非互斥参与者数量会动态变化的多阶段协作Phaser支持register/deregister动态增减分层降低竞争恰好两个线程之间周期性交换数据Exchanger双向阻塞的面对面数据交换减少复制开销一条主线贯穿全篇这几个工具解决的都不是互斥问题那是锁的地盘而是线程之间怎么排节奏——谁等谁、等到什么条件、要不要重复等——分辨清楚这一点选型时就不会纠结。带走两句话① 分不清用CountDownLatch还是CyclicBarrier时问自己等待者和被等待者是不是同一批线程——不是同一批用门闩是同一批互相等用栅栏。②Semaphore不是锁的替代品它表达的是资源配额release() 不要求必须是 acquire() 的那个线程这个特性用对了是灵活性用错了是隐藏的 bug 来源。下一篇我们进入ThreadLocal——它和这一篇的思路完全相反不是让多个线程协作而是让每个线程拥有自己独立的一份数据彼此互不干扰。会讲清楚它的实现原理、经典的内存泄漏问题从哪来以及InheritableThreadLocal、阿里开源的TransmittableThreadLocal分别解决了什么场景。