AQS是JUC中很多同步组件的构建基础,简单来讲,它内部实现主要是状态变量state和一个FIFO队列来完成,同步队列的头结点是当前获取到同步状态的结点,获取同步状态state失败的线程,会被构造成一个结点(或共享式或独占式)加入到同步队列尾部(采用自旋CAS来保证此操作的线程安全),随后线程会阻塞;释放时唤醒头结点的后继结点,使其加入对同步状态的争夺中。
AQS为我们定义好了顶层的处理实现逻辑,我们在使用AQS构建符合我们需求的同步组件时,只需重写tryAcquire,tryAcquireShared,tryRelease,tryReleaseShared几个方法,来决定同步状态的释放和获取即可,至于背后复杂的线程排队,线程阻塞/唤醒,如何保证线程安全,都由AQS为我们完成了,这也是非常典型的模板方法的应用。AQS定义好顶级逻辑的骨架,并提取出公用的线程入队列/出队列,阻塞/唤醒等一系列复杂逻辑的实现,将部分简单的可由使用者决定的操作逻辑延迟到子类中去实现。
1package com.abstractqueuesynchronizer; 2 3import java.util.concurrent.locks.AbstractQueuedSynchronizer; 4 5public class SelfAbstractQueueSynchronizer { 6 7 8 //继承AbstractQueuedSynchronizer类 9 private static class Syn extends AbstractQueuedSynchronizer { 10 11 private static final long serialVersionUID = 1L; 12 13 //是否拥有锁 14 protected boolean isHeldExclusively() { 15 return getState() == 1; 16 } 17 18 //获取锁 19 public boolean tryAcquire(int acquires) { 20 if(compareAndSetState(0, 1)) { 21 setExclusiveOwnerThread(Thread.currentThread()); 22 return true; 23 } 24 return false; 25 } 26 27 //释放所 28 protected boolean tryRelease(int releases) { 29 if(getState() == 0) 30 throw new IllegalMonitorStateException(); 31 setExclusiveOwnerThread(null); 32 setState(0); 33 return true; 34 } 35 } 36 37 private final Syn syn = new Syn(); 38 39 public void lock() { 40 syn.acquire(1); 41 } 42 43 public boolean tryLock() { 44 return syn.tryAcquire(1); 45 } 46 47 public void unlock() { 48 syn.release(1); 49 } 50 51 public boolean isLocked() { 52 return syn.isHeldExclusively(); 53 } 54}
测试
1package com.abstractqueuesynchronizer; 2 3import java.util.concurrent.BrokenBarrierException; 4import java.util.concurrent.CyclicBarrier; 5 6public class Main { 7 8 private static CyclicBarrier barrier = new CyclicBarrier(31); 9 private static int a = 0; 10 private static SelfAbstractQueueSynchronizer test = new SelfAbstractQueueSynchronizer(); 11 12 public static void main(String[] args) throws InterruptedException, BrokenBarrierException { 13 14 for(int i = 0; i < 30; i++) { 15 Thread t = new Thread(new Runnable(){ 16 17 @Override 18 public void run() { 19 for(int i = 0; i < 1000; i++) { 20 unlockIncrement(); 21 } 22 try { 23 barrier.await(); 24 } catch (Exception e) { 25 e.printStackTrace();; 26 } 27 } 28 29 }); 30 t.start(); 31 } 32 33 barrier.await(); 34 System.out.println("unlock model a= " + a); 35 36 System.out.println("##########################"); 37 38 barrier.reset(); 39 a = 0; 40 for(int i = 0; i < 30; i++) { 41 Thread t = new Thread(new Runnable(){ 42 @Override 43 public void run() { 44 for(int i = 0; i < 1000; i++ ) { 45 lockIncrement(); 46 } 47 48 try { 49 barrier.await(); 50 } catch (InterruptedException e) { 51 e.printStackTrace(); 52 } catch (BrokenBarrierException e) { 53 e.printStackTrace(); 54 } 55 } 56 57 }); 58 t.start(); 59 } 60 61 barrier.await(); 62 System.out.println("lock model a= " + a); 63 64 } 65 66 public static void unlockIncrement() { 67 a++; 68 } 69 70 public static void lockIncrement() { 71 test.lock(); 72 a++; 73 test.unlock(); 74 } 75 76}