0%

Java多线程等待的N种打开方式(第一篇)

单核CPU就是”一根筋”,只是简单的一条一条的串行执行机器指令,且每个时钟周期内只执行一条指令,如果一个进程不间断的运行,那么就只能运行一个任务;

将CPU的运行时间划分成一个个的时间段,这就是时间片(time slice),进程通过抢占时间片,以达到交替执行的效果,宏观上,多个任务同时运行,JVM正是采用抢占式调度的方式根据优先级来”有限”的控制线程的运行。

为了提高代码的执行效率,编译器会在编译时做指令重排等优化,这就造成了在多线程运行时,出现与单线程运行不一致的情况。再加上,为了优化多核执行效率,造成各个线程访问的变量值与内存主存的值不一致的情况,就出现内存可见性问题。

JMM通过内存屏障、CAS和mutex lock指令等手段,实现各种各样的锁,究竟Java提供了哪些方法可以实现线程的等待呢,接下来以一个样例场景打开”Java多线程等待的N中方式”。

题目与约定

假设有A、B、C三个线程,C必须等待A和B都结束时才能运行,

约定每个线程的耗时操作为方法doWork(String threadName)

接下来将以此样例分析实现方法和Java实现的原理:

  1. 第一篇: monitor
  • Thread.join等待线程完成
  • Object.wait-notify等待唤醒
  • LockSupport
  1. 第二篇: AQS
  • CountDownLatch
  • ReentrantLock
  • Semaphore
  • CyclicBarrier循环屏障
  1. 第三篇: 队列
  • Future
  • BlockingDeque
  • ListenableFuture

Thread.Join等待线程完成

1
public final void join(long millis,int nanos) throws InterruptedException

等待最多millis毫秒+nanos纳秒等待完成完成,通过循环调用this.wait等待this.isAlive条件完成,当调用this.notifyAll方法时,此线程终止。建议应用程序不要在线程实例上使用waitnotifynotifyAll

解题

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
final Thread threadA = new Thread(() -> {
doWork("A");
});
final Thread threadB = new Thread(() -> {
doWork("B");
});
final Thread threadC = new Thread(() -> {
try {
threadA.join(); // 等待A结束
threadB.join(); // 等待B结束
} catch (InterruptedException e) {
e.printStackTrace();
}
doWork("C");
});
// 启动A、B两个线程
threadA.start();
threadB.start();
// 启动C线程
threadC.start();

运行结果

1
2
3
4
5
6
Thread A starting 
Thread B starting
Thread B finished
Thread A finished
Thread C starting
Thread C finished

实现原理

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
public class Thread implements Runnable {
...
public final void join() throws InterruptedException {
join(0);
}
...
public final synchronized void join(long millis) throws InterruptedException {
long base = System.currentTimeMillis();
long now = 0;
if (millis < 0) {
throw new IllegalArgumentException("timeout value is negative");
}
if (millis == 0) { // 0表示没有设置超时时间
while (isAlive()) {//isAlive获取线程状态,无限等待直到previousThread线程结束
wait(0); //调用Object中的wait方法实现线程的阻塞
}
} else { //阻塞直到超时
while (isAlive()) {
long delay = millis - now;
if (delay <= 0) {
break;
}
wait(delay);
now = System.currentTimeMillis() - base;
}
}
}
/**
* 测试此线程是否存活,如果此线程started并且还没有死亡,则存活
* @return <code>true</code> 如果此线程存活;
* <code>false</code> 否则.
*/
public final native boolean isAlive();
...
}
public class Object {
/**
* 导致当前线程等待,直到另一个线程调用此对象的notify()方法或notifyAll()方法
* 当前线程必须拥有此对象的监视器。当前线程释放此监视器的所有权,并等待,直到另一个线程通过调用notify方法或notifyAll方法唤通知等待此对象的监视器的线程唤醒。然后,线程等待,直到它可以重新获得监视器的所有权并继续执行
*/
public final native void wait(long timeout) throws InterruptedException;
}
1
2
3
4
5
6
7
void JavaThread::exit(bool destroy_vm, ExitType exit_type) {
assert(this == JavaThread::current(), "thread consistency check");
...
// 在线程对象上通知等待者。必须在线程上调用exit()之后执行此操作(如果该线程是守护线程组中的最后一个线程,则通知等待对象之前,线程组应设置已破坏的位)。
ensure_join(this);
assert(!this->has_pending_exception(), "ensure_join should have cleared");
...

join方法本质是调用的Object中的wait方法实现线程的阻塞,调用wait方法必须要获取锁,所以join方法是被synchronized修饰的,synchronized修饰在方法层面相当于synchronized(this),this就是previousThread本身的实例,两次调用join的previousThread分别是A和B两个线程,此时previousThread线程对象的监视器也就只有一个C线程在等待,因此当previousThread执行完的时候,通过notifynotifyAll会立即唤醒(只有一个线程嘛)。

此例子中,有两把锁存在,分别等待threadAthreadB:

1
2
threadA.join(); // 调用threadA.wait(),threadA运行完后,默认会唤醒threadA对象上等待的threadC
threadB.join(); // 调用threadB.wait(),threadB运行完后,默认会唤醒threadB对象上等待的threadC

注意:建议调用Thread.join的时候,不要在线程实例上同时调用waitnofitynotifyAll方法

Object.wait-notify等待唤醒

wait-nofity是Object中的等待唤醒组合,实现等待两个线程完成,可以使用一把锁,也可以使用两把锁。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
// 一把锁的解法
// finishedCount是一个数组对象,可以通过这个对象实现C等待A和B完成,C循环扫描完成个数
int[] finishedCount = {0};
Thread aThread = new Thread(() -> {
doWork("A");
synchronized (finishedCount) {
// 完成一个,加一
finishedCount[0]++;
// 唤醒C
finishedCount.notify();
}
});
Thread bThread = new Thread(() -> {
doWork("B");
synchronized (finishedCount) {
// 完成一个,加一
finishedCount[0]++;
// 唤醒C
finishedCount.notify();
}
});
Thread cThread = new Thread(() -> {
// 循环扫描完成个数,只到两个都完成
while (finishedCount[0] < 2) {
synchronized (finishedCount) {
try {
finishedCount.wait();
} catch (InterruptedException e) {
}
}
}
doWork("C");
});
aThread.start();
bThread.start();
cThread.start();

C在对象finishedCount上等待,A和B完成时,首先把finishedCount加一,然后唤醒C

C被唤醒时,判断A和B已经完成(finishedCount是否为2),完成则开始执行C的任务

Object.wait-notify实现原理

在HotSpot虚拟机中,monitor采用ObjectMonitor实现。

ObjectMonitor对象中有两个队列,都用来保存ObjectWaiter对象,分别是WaitSetEntrySet_owner用来指向获得ObjectMonitor对象的线程

ObjectWaiter对象是双向链表结构,保存了_thread(当前线程)以及当前的状态TState等数据, 每个等待锁的线程都会被封装成ObjectWaiter对象。
_WaitSet :处于wait状态的线程,会被加入到wait set;

_EntrySet:处于等待锁block状态的线程,会被加入到entry set;

wait方法实现

lock.wait()方法最终通过ObjectMonitorwait(jlong millis, bool interruptable, TRAPS)实现

  1. 将当前线程封装成ObjectWaiter对象node

  2. 通过ObjectMonitor::AddWaiter方法将node添加到_WaitSet列表中

  3. 通过ObjectMonitor::exit方法释放当前的ObjectMonitor对象,这样其它竞争线程就可以获取该ObjectMonitor对象

  4. 最终底层的park方法会挂起线程

ObjectSynchorizer::wait方法通过Object对象找到ObjectMonitor对象来调用方法 ObjectMonitor::wait(),通过调用ObjectMonitor::AddWaiter()可以把新建的ObjectWaiter对象,放入到_WaitSet队列的末尾,然后在ObjectMonitor::exit释放锁,接着通过执行thread_ParkEvent->park来挂起线程,也就是进行wait

notify方法实现

lock.notify()方法最终通过ObjectMonitor的void notify(TRAPS)实现:
1、如果当前_WaitSet为空,即没有正在等待的线程,则直接返回;
2、通过ObjectMonitor::DequeueWaiter方法,获取_WaitSet列表中的第一个ObjectWaiter节点,实现也很简单。

这里需要注意的是,在jdk的notify方法注释是随机唤醒一个线程,其实是第一个ObjectWaiter节点

3、根据不同的策略,将取出来的ObjectWaiter节点,加入到_EntryList或则通过Atomic::cmpxchg_ptr指令进行自旋操作cxq,具体代码实现有点长,这里就不贴了,有兴趣的同学可以看objectMonitor::notify方法;

notifyAll方法实现

lock.notifyAll()方法最终通过ObjectMonitorvoid notifyAll(TRAPS)实现:
通过for循环取出_WaitSetObjectWaiter节点,并根据不同策略,加入到_EntryList或则进行自旋操作。

从JVM的方法实现中,可以发现:notifynotifyAll并不会释放所占有的ObjectMonitor对象,其实真正释放ObjectMonitor对象的时间点是在执行monitorexit指令,一旦释放ObjectMonitor对象了,entry set中ObjectWaiter节点所保存的线程就可以开始竞争ObjectMonitor对象进行加锁操作了。

参考文献

  1. https://docs.oracle.com/javase/8/docs/