
开发一款企业级软件,不管是手机APP,还是WEB产品,就像雷神山医院的建设一样,是一个纷繁复杂的构建过程。相较于三维世界的宏观物理工程,软件工程显得更细小,更奇幻。Java的世界也是如此,sun公司的先驱们将物理世界的处事逻辑浓缩到了JRE的库文件中,以供大家观摩、使用,自己也从中获取名和利。作为一种共享内存的编程语言,Java多线程编程让程序员多长了很多根白头发(包括我),花了数万块钱买了一台服务器,出于资本家思维,我们希望能充分利用好服务器的硬件资源,单机编程面向多线程开发成了不二首选。为了能精准地控制好线程及线程间信息交换,sun公司的大拿们为我们引入了线程池、同步锁。为了能让急吼吼的线程冷却下来,大拿们又为我们提供了像CountDownLatch这样的类库。这就是Java的世界!庞大的社区让你在现实中能碰到的问题都已然开发成依赖库以供使用,你只要具备快速搜索信息的能力就可以了,这也激发了程序员的快速集成能力。下面我们就通透地走读CountDownLatch源码,试图窥视大拿们开发现场的情感思维。
1 CountDownLath的基础使用
这里简单地使用CountDownLatch控制主子线程的输出顺序的例子描述CountDownLath的使用。代码片段如下
public class CountDownLathExample {
public static class Seller implements Runnable{
private CountDownLatch latch;
public Seller(CountDownLatch latch) { this.latch = latch; }
public void run() {
try {
try {
TimeUnit.SECONDS.sleep(3);
} catch (InterruptedException e) {}
System.out.println(Thread.currentThread().getName() + " working...");
} finally {
latch.countDown();
}
}
}
public static void main(String[] args) {
CountDownLatch latch = new CountDownLatch(2);
for(int i = 0;i < latch.getCount();i++)
new Thread(new Seller(latch)).start();
try {
latch.await();
} catch (InterruptedException e) {}
System.out.println(Thread.currentThread().getName()+" done!");
}
}
如果不使用CountDownLatch控制,依照main函数的逻辑,输出如下
main done!
Thread-1 working...
Thread-0 working...
主线程的逻辑会首先输出,子线程会睡3秒然后输出。使用CountDownLatch控制之后,输出如下
Thread-0 working...
Thread-1 working...
main done!
主线会阻塞着等待子线程执行完成之后才自己执行。这也是我们想看到的结果。
2 CountDownLatch源码分析
例子中,我首先使用了CountDownLatch一个参数的构造器来初始化CountDownLatch实例,参数指定了抢夺临界区资源的线程数量。进入到源码中,
public CountDownLatch(int count) {
if (count < 0) throw new IllegalArgumentException("count < 0");
this.sync = new Sync(count);
}
将线程数量的参数委托给了Sync类,
private static final class Sync extends AbstractQueuedSynchronizer {
...
Sync(int count) {
setState(count);
}
...
}
Sync类为AQS的子实现类,构造器中将参数传递给了AQS的state变量,state变量使用volatile关键字修饰保证了线程可见性。
接着,调用countDown()接口来释放资源。实现逻辑为
public void countDown() {
sync.releaseShared(1);
}
...
public final boolean releaseShared(int arg) {
if (tryReleaseShared(arg)) {
doReleaseShared();
return true;
}
return false;
}
...
protected boolean tryReleaseShared(int releases) {
// Decrement count; signal when transition to zero
for (;;) {
int c = getState();
if (c == 0)
return false;
int nextc = c-1;
if (compareAndSetState(c, nextc))
return nextc == 0;
}
}
...
private void doReleaseShared() {
for (;;) {
Node h = head;
if (h != null && h != tail) {
int ws = h.waitStatus;
if (ws == Node.SIGNAL) {
if (!compareAndSetWaitStatus(h, Node.SIGNAL, 0))
continue; // loop to recheck cases
unparkSuccessor(h);
}
else if (ws == 0 &&
!compareAndSetWaitStatus(h, 0, Node.PROPAGATE))
continue; // loop on failed CAS
}
if (h == head) // loop if head changed
break;
}
}
...
private void unparkSuccessor(Node node) {
int ws = node.waitStatus;
if (ws < 0)
compareAndSetWaitStatus(node, ws, 0);
Node s = node.next;
if (s == null || s.waitStatus > 0) {
s = null;
for (Node t = tail; t != null && t != node; t = t.prev)
if (t.waitStatus <= 0)
s = t;
}
if (s != null)
LockSupport.unpark(s.thread);
}
CountDownLatch将释放资源的处理委托给了AQS的releaseShared()方法。在这个方法中,首先调用Sync子类的tryReleaseShared()方法来探测是否满足释放资源的条件。Sync在tryReleaseShared()方法中自旋CAS操作保证volatile变量state减一操作的原子性。如果state的值为0则调用doReleaseShared()方法释放共享资源。我们知道AQS将对资源的访问模式抽象成了共享模式和独占模式,并且维护着一个等待获取资源线程的双向队列,以便线程按顺序访问临界区数据,队列的结点中维护着一个waitStatus等待状态。等待状态可取值CANCELLED,值为1,代表结点维护的线程等待超时或被中断,需要从队列中取消等待;SIGNAL,值为-1,代表后续结点处于等待状态,当前结点的线程释放了资源或取消,将会通知后续结点,使后续结点的线程得以运行;CONDITION,值为-2,代表结点在条件队列中,当其他线程调用Condition的signal()方法之后,该结点会从条件队列转移到同步队列中;PROPAGATE,值为-3,代表下一次共享模式的共享资源获取将会无条件传播下去。这里,AQS的state属性初始值为2,当线程1进来,state会减一,并不会释放资源,当线程2进来,state也会减一,释放资源,并调用unparkSuccessor()方法唤醒(LockSupport.unpark)主线程。
最后,主线程中调用await()方法等待子线程执行完成,代码片如下
public void await() throws InterruptedException {
sync.acquireSharedInterruptibly(1);
}
...
public final void acquireSharedInterruptibly(int arg)
throws InterruptedException {
if (Thread.interrupted())
throw new InterruptedException();
if (tryAcquireShared(arg) < 0)
doAcquireSharedInterruptibly(arg);
}
...
protected int tryAcquireShared(int acquires) {
return (getState() == 0) ? 1 : -1;
}
...
private void doAcquireSharedInterruptibly(int arg)
throws InterruptedException {
// 增加共享模式的结点到队尾
final Node node = addWaiter(Node.SHARED);
boolean failed = true; // 标记是否成功获取资源
try {
for (;;) {
final Node p = node.predecessor(); // 获取前驱结点
if (p == head) { // 如果前驱结点是头结点
int r = tryAcquireShared(arg); // 尝试获取共享资源
if (r >= 0) {
setHeadAndPropagate(node, r);
p.next = null; // help GC
failed = false;
return;
}
}
// 如果当前线程可以挂起了,就进入到waiting状态,直到被unpark()
if (shouldParkAfterFailedAcquire(p, node) &&
parkAndCheckInterrupt())
throw new InterruptedException();
}
} finally {
if (failed)
cancelAcquire(node);
}
}
...
private Node addWaiter(Node mode) {
Node node = new Node(Thread.currentThread(), mode);
Node pred = tail;
if (pred != null) {
node.prev = pred;
if (compareAndSetTail(pred, node)) {
pred.next = node;
return node;
}
}
enq(node); // 将node加入到队尾
return node;
}
...
private Node enq(final Node node) {
// CAS 自旋,直到成功入队
for (;;) {
Node t = tail;
if (t == null) { // 队列为空,创建一个空结点作为head结点,并将tail也指向这个结点
if (compareAndSetHead(new Node()))
tail = head;
} else { // node正常入队
node.prev = t;
if (compareAndSetTail(t, node)) {
t.next = node;
return t;
}
}
}
}
...
// 检查状态,判断当前线程是否可以进入waiting状态
private static boolean shouldParkAfterFailedAcquire(Node pred, Node node) {
int ws = pred.waitStatus; // 获取前驱结点的等待状态
if (ws == Node.SIGNAL) // 前驱结点执行完成之后会通知自己
return true;
if (ws > 0) {
// 如果前驱结点取消了,就一直往前找,直到找到最近一个正常等待状态的结点,并排在它后面
do {
node.prev = pred = pred.prev;
} while (pred.waitStatus > 0);
pred.next = node;
} else {
// 如果前驱结点正常,那就把前驱结点的状态设置为SIGNAL,告诉它执行结束通知自己一下
compareAndSetWaitStatus(pred, ws, Node.SIGNAL);
}
return false;
}
...
private final boolean parkAndCheckInterrupt() {
LockSupport.park(this);// 让自己进入到waiting状态
return Thread.interrupted(); // 如果被唤醒,查看自己是不是被中断
}
这里在主线程中调用await()方法,会让自己进行到waiting状态,直到子线程执行完成,并调用LockSupport.unpark()唤醒。
3 AQS实现分析
从CountDownLatch的源码我们知道:它把线程对临界区资源的访问控制全部委托给了AQS的子类Sync类完成,并且对资源的访问属于共享模式。在Sync类中,仅仅在构造器中初始化了AQS的state属性(代表共享资源),覆盖了尝试以共享模式获取资源和尝试以共享模式释放资源的方法,增加了一个返回共享资源的方法,此外,并没有过多扩展,几乎所有对资源访问控制的核心工作都由AQS完成。那么,AQS到底是怎么设计的呢?从上面的分析我们可以看出,AQS把对共享资源的访问分了两种模式:共享模式(Share)和独占模式(Exclusive);此外,它还维护了一个代表共享资源的volatile修饰的整型变量state和一个先进先出的双向队列(线程抢占资源被阻塞时会进入此队列),至此,AQS的轮廓基本上就出来了:线程获取共享资源,如果共享资源已经被占用,线程进行排队;否则,占用共享资源。从上述需求可以分析出:AQS应该具备尝试获取资源、尝试释放资源、线程排队的能力和持有共享资源被占用的标志。由于AQS将资源访问划分为共享模式和独占模式,那么,尝试获取资源也就可以分为尝试以共享模式获取共享资源和尝试以独占方式获取共享资源;释放资源也自然可以分为尝试以共享模式释放共享资源和尝试以独占模式释放共享资源。针对于线程排队,队列可使用链表实现,原生Thread类不具备排队的能力,因此,可将Thread包装成为链表的结点。AQS抽取的结点信息如下
waitStatus - 整型,代表线程等待状态,使用volatile修饰,确保线程可见
prev - 前驱结点引用,使用volatile修饰
next - 后继结点引用,使用volatile修饰
thread - Thread类型,使用volatile修饰
nextWaiter - Node类型,代表当前线程对共享资源的访问模式
有了结点抽象,为了实现链表,自然需要持有头结点引用(head)和尾结点引用(tail)。最后,使用整型变量state代表共享资源占用标志。接下来就是各个方法的具体实现及凭借技术手段保障实现的正确性即可。
上面分析了尝试以共享模式获取资源(tryAcquireShared)和尝试以共享模式释放资源(tryReleaseShared)在AQS中的实现,出于篇幅考虑,后面我们再分析尝试以独占模式后去资源()和尝试以独占模式释放资源的实现。




