0%

<<Java高级软件工程师知识结构

Java多线程是Java基础的重要的一部分,支持多线程是Java的重要特性之一. 主要包括如下内容:
  1. Java多线程1: 线程生命周期和多线程基础
  2. Java多线程2: Lock、信号量、原子量与队列
  3. Java多线程3: volatile
  4. Java多线程4: 同步锁与Java线程同步方法比较
  5. Java多线程5: 线程池
  6. Java多线程6: Java阻塞队列与生产者消费者模式

更详细内容请参考博文 Java并发编程:Lock

Sun在Java5中, 对Java线程的类库做了大量的扩展

java.util.concurrent.locks 包下面, 里面有三个重要的接口ConditionLockReadWriteLock.

Condition

Condition 将 Object 监视器方法(waitnotifynotifyAll)分解成截然不同的对象, 以便通过将这些对象与任意 Lock 实现组合使用, 为每个对象提供多个等待 setwait-set).

Lock

Lock实现提供了比使用 synchronized 方法和语句可获得的更广泛的锁定操作.

ReadWriteLock

ReadWriteLock维护了一对相关的锁定, 一个用于只读操作, 另一个用于写入操作.

Lock 与 synchronized的区别

见文献 深入研究 Java Synchronize 和 Lock 的区别与用法

  1. 用法上的区别
    synchronized可以加在方法或代码块上,而Lock必须指定起始位置,一般使用ReentrantLock类做为锁,多个线程中必须要使用一个ReentrantLock类做为对象才能保证锁的生效。且在加锁和解锁处需要通过lock()unlock()显式指出。所以一般会在finally块中写unlock()以防死锁。

    1. ReentrantLock非阻塞
    2. ReentrantLock CAS实现无锁
    3. 可以指定等待时间,可以中断
    4. 公平锁:按照申请顺序
    5. 可以指定条件
    6. await必须与while一起使用
    7. Lock可以知道线程有没有成功获取到锁。这个是synchronized无法办到的
    8. synchronized在发生异常时,会自动释放线程占有的锁,不会导致死锁现象发生;
      而Lock在发生异常时,如果没有主动通过unLock()去释放锁,则很可能造成死锁现象,因此使用Lock时需要在finally块中释放锁;
  2. 性能上的区别
    synchronized是托管给JVM执行的,而lock是java写的控制锁的代码。synchronized是悲观锁,线程获取到的是独占锁。独占锁意味着其他线程只能依靠阻塞来等待线程释放锁。多个线程竞争资源,CPU频繁切换效率变低。

    Lock使用的是乐观锁,乐观锁就是CAS,调用CPU的指令,效率比较高。是非阻塞算法。 每次不加锁而是假设没有冲突而去完成某项操作,如果因为冲突失败就重试,直到成功为止。乐观锁实现的机制就是CAS操作(Compare and Swap)。获得锁的一个方法是compareAndSetState. CPU提供了指令,可以自动更新共享数据,而且能够检测到其他线程的干扰,而 compareAndSet() 就用这些代替了锁定。这个算法称作非阻塞算法,意思是一个线程的失败或者挂起不应该影响其他线程的失败或挂起的算法。

  3. 用途上的区别
    高并发情况下,比较适合使用Lock,特别是在下面的情况下:

    • 某个线程在等待一个锁的控制权的这段时间需要中断
    • 需要分开处理一些wait-notifyReentrantLock里面的Condition应用,能够控制notify哪个线程
    • 具有公平锁功能,每个到来的线程都将排队等候

ReentrantLock

ReentrantLock是唯一实现了Lock接口的类

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
//创建并发访问的账户
MyCount myCount = new MyCount("95599200901215522", 10000);
//创建一个锁对象
Lock lock = new ReentrantLock();
//创建一个线程池
ExecutorService pool = Executors.newCachedThreadPool();
//创建一些并发访问用户, 一个信用卡, 存的存, 取的取, 好热闹啊
User u1 = new User("张三", myCount, -4000, lock);
User u2 = new User("张三他爹", myCount, 6000, lock);
User u3 = new User("张三他弟", myCount, -8000, lock);
User u4 = new User("张三", myCount, 800, lock);
//在线程池中执行各个用户的操作
pool.execute(u1);
pool.execute(u2);
pool.execute(u3);
pool.execute(u4);
//关闭线程池
pool.shutdown();
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
/**
* 信用卡的用户
*/
class User implements Runnable {
private String name; //用户名
private MyCount myCount; //所要操作的账户
private int iocash; //操作的金额, 当然有正负之分了
private Lock myLock; //执行操作所需的锁对象
User(String name, MyCount myCount, int iocash, Lock myLock) {
this.name = name;
this.myCount = myCount;
this.iocash = iocash;
this.myLock = myLock;
}
public void run() {
//获取锁
myLock.lock();
//执行现金业务
System.out.println(name + "正在操作" + myCount + "账户, 金额为" + iocash + ", 当前金额为" + myCount.getCash());
myCount.setCash(myCount.getCash() + iocash);
System.out.println(name + "操作" + myCount + "账户成功, 金额为" + iocash + ", 当前金额为" + myCount.getCash());
//释放锁, 否则别的线程没有机会执行了
myLock.unlock();
}
}
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
/**
* 信用卡账户, 可随意透支
*/
class MyCount {
private String oid; //账号
private int cash; //账户余额
MyCount(String oid, int cash) {
this.oid = oid;
this.cash = cash;
}
public String getOid() {
return oid;
}
public void setOid(String oid) {
this.oid = oid;
}
public int getCash() {
return cash;
}
public void setCash(int cash) {
this.cash = cash;
}
@Override
public String toString() {
return "MyCount{" +
"oid='" + oid + '\'' +
", cash=" + cash +
'}';
}
}

线程的中断

Java线程中的中断,提供了三个方法:

  • interrupt()
    通知线程中断, 但是线程不会立即中断, 只会将interrupt标志位设置为 true, 当在线程中获取到这个标志位时自行判断
  • interrupted()
    判断是否中断, 如果未中断, 立即设置为中断
  • isInterrupted()
    判断是否中断

这三个方法都不能使线程中断执行。而ReetrantLock提供了响应的方案。

ReentrantLock的lock机制有2种,忽略中断锁和响应中断锁,这给我们带来了很大的灵活性。

比如:如果A、B2个线程去竞争锁,A线程得到了锁,B线程等待,但是A线程这个时候实在有太多事情要处理,就是一直不返回,B线程可能就会等不及了,想中断自己,不再等待这个锁了,转而处理其他事情。这个时候ReentrantLock就提供了2种机制,

  • 第一,B线程中断自己(或者别的线程中断它),但是ReentrantLock不去响应,继续让B线程等待,你再怎么中断,我全当耳边风(synchronized原语就是如此);Lock.lock()也不会响应中断操作
  • 第二,B线程中断自己(或者别的线程中断它),ReentrantLock处理了这个中断,并且不再等待这个锁的到来,完全放弃。
    在Thread类中使用lock.lockInterruptibly();锁定时可以响应中断操作,当调用Thread.interupt()方法,该lock会放弃锁定,抛出InterruptedException异常
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
volatile boolean isProcess = false;
ReentrantLock lock = new ReentrantLock();
Condition processReady = lock.newCondition();
thread: run() {
lock.lock();
isProcess = true;
try {
while(!isProcessReady) {
//isProcessReady 是另外一个线程的控制变量
processReady.await();
//释放了lock,在此等待signal
}catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
lock.unlock();
isProcess = false;
}
}
}
}
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
43
44
45
46
47
public class Test {
private Lock lock = new ReentrantLock();
public static void main(String[] args) {
Test test = new Test();
MyThread thread1 = new MyThread(test);
MyThread thread2 = new MyThread(test);
thread1.start();
thread2.start();
try {
Thread.sleep(2000);
} catch (InterruptedException e) {
e.printStackTrace();
}
thread2.interrupt();
}
public void insert(Thread thread) throws InterruptedException{
lock.lockInterruptibly();
//注意,如果需要正确中断等待锁的线程,必须将获取锁放在外面,然后将InterruptedException抛出
try {
System.out.println(thread.getName()+"得到了锁");
long startTime = System.currentTimeMillis();
for( ; ;) {
if(System.currentTimeMillis() - startTime >= Integer.MAX_VALUE)
break;
//插入数据
}
}finally {
System.out.println(Thread.currentThread().getName()+"执行finally");
lock.unlock();
System.out.println(thread.getName()+"释放了锁");
}
}
}
class MyThread extends Thread {
private Test test = null;
public MyThread(Test test) {
this.test = test;
}
@Override
public void run() {
try {
test.insert(Thread.currentThread());
} catch (InterruptedException e) {
System.out.println(Thread.currentThread().getName()+"被中断");
}
}
}

读写锁 ReadWriteLock

为了提高性能, Java提供了读写锁, 在读的地方使用读锁, 在写的地方使用写锁, 灵活控制, 在一定程度上提高了程序的执行效率.
Java中读写锁有个接口java.util.concurrent.locks.ReadWriteLock, 也有具体的实现ReentrantReadWriteLock

ReadWriteLock也是一个接口,在它里面只定义了两个方法:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
public interface ReadWriteLock {
/**
* Returns the lock used for reading.
*
* @return the lock used for reading.
*/
Lock readLock();
/**
* Returns the lock used for writing.
*
* @return the lock used for writing.
*/
Lock writeLock();
}

ReentrantReadWriteLock里面提供了很多丰富的方法,不过最主要的有两个方法:readLock()和writeLock()用来获取读锁和写锁。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
//创建并发访问的账户
MyCount myCount = new MyCount("95599200901215522", 10000);
//创建一个锁对象
ReadWriteLock lock = new ReentrantReadWriteLock(false);
//创建一个线程池
ExecutorService pool = Executors.newFixedThreadPool(2);
//创建一些并发访问用户, 一个信用卡, 存的存, 取的取, 好热闹啊
User u1 = new User("张三", myCount, -4000, lock, false);
User u2 = new User("张三他爹", myCount, 6000, lock, false);
User u3 = new User("张三他弟", myCount, -8000, lock, false);
User u4 = new User("张三", myCount, 800, lock, false);
User u5 = new User("张三他爹", myCount, 0, lock, true);
//在线程池中执行各个用户的操作
pool.execute(u1);
pool.execute(u2);
pool.execute(u3);
pool.execute(u4);
pool.execute(u5);
//关闭线程池
pool.shutdown();
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
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
/**
* 信用卡的用户
*/
class User implements Runnable {
private String name; //用户名
private MyCount myCount; //所要操作的账户
private int iocash; //操作的金额, 当然有正负之分了
private ReadWriteLock myLock; //执行操作所需的锁对象
private boolean ischeck; //是否查询
User(String name, MyCount myCount, int iocash, ReadWriteLock myLock, boolean ischeck) {
this.name = name;
this.myCount = myCount;
this.iocash = iocash;
this.myLock = myLock;
this.ischeck = ischeck;
}
public void run() {
if (ischeck) {
//获取读锁
myLock.readLock().lock();
System.out.println("读:" + name + "正在查询" + myCount + "账户, 当前金额为" + myCount.getCash());
//释放读锁
myLock.readLock().unlock();
} else {
//获取写锁
myLock.writeLock().lock();
//执行现金业务
System.out.println("写:" + name + "正在操作" + myCount + "账户, 金额为" + iocash + ", 当前金额为" + myCount.getCash());
myCount.setCash(myCount.getCash() + iocash);
System.out.println("写:" + name + "操作" + myCount + "账户成功, 金额为" + iocash + ", 当前金额为" + myCount.getCash());
//释放写锁
myLock.writeLock().unlock();
}
}
}

/**
* 信用卡账户, 可随意透支
*/
class MyCount {
private String oid; //账号
private int cash; //账户余额
MyCount(String oid, int cash) {
this.oid = oid;
this.cash = cash;
}
public String getOid() {
return oid;
}
public void setOid(String oid) {
this.oid = oid;
}
public int getCash() {
return cash;
}
public void setCash(int cash) {
this.cash = cash;
}
@Override
public String toString() {
return "MyCount{" +
"oid='" + oid + '\'' +
", cash=" + cash +
'}';
}
}

条件变量 Condition

条件变量就是表示条件的一种变量. 但是必须说明, 这里的条件是没有实际含义的, 仅仅是个标记而已, 并且条件的含义往往通过代码来赋予其含义.

条件变量都实现了java.util.concurrent.locks.Condition接口, 条件变量的实例化是通过一个Lock对象上调用newCondition()方法来获取的, 这样, 条件就和一个锁对象绑定起来了. 因此, Java中的条件变量只能和锁配合使用, 来控制并发程序访问竞争资源的安全.

在Java5中, 一个锁可以有多个条件, 每个条件上可以有多个线程等待, 通过调用await()方法, 可以让线程在该条件下等待. 当调用signalAll()方法, 又可以唤醒该条件下的等待的线程.

实例

有一个账户, 多个用户(线程)在同时操作这个账户, 有的存款有的取款, 存款随便存, 取款有限制, 不能透支, 任何试图透支的操作都将等待里面有足够存款才执行操作.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
//创建并发访问的账户
MyCount myCount = new MyCount("95599200901215522", 10000);
//创建一个线程池
ExecutorService pool = Executors.newFixedThreadPool(2);
Thread t1 = new SaveThread("张三", myCount, 2000);
Thread t2 = new SaveThread("李四", myCount, 3600);
Thread t3 = new DrawThread("王五", myCount, 2700);
Thread t4 = new SaveThread("老张", myCount, 600);
Thread t5 = new DrawThread("老牛", myCount, 1300);
Thread t6 = new DrawThread("胖子", myCount, 800);
//执行各个线程
pool.execute(t1);
pool.execute(t2);
pool.execute(t3);
pool.execute(t4);
pool.execute(t5);
pool.execute(t6);
//关闭线程池
pool.shutdown();
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
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
/**
* 存款线程类
*/
class SaveThread extends Thread {
private String name; //操作人
private MyCount myCount; //账户
private int x; //存款金额
SaveThread(String name, MyCount myCount, int x) {
this.name = name;
this.myCount = myCount;
this.x = x;
}
public void run() {
myCount.saving(x, name);
}
}

/**
* 取款线程类
*/
class DrawThread extends Thread {
private String name; //操作人
private MyCount myCount; //账户
private int x; //存款金额
DrawThread(String name, MyCount myCount, int x) {
this.name = name;
this.myCount = myCount;
this.x = x;
}
public void run() {
myCount.drawing(x, name);
}
}
/**
* 普通银行账户, 不可透支
*/
class MyCount {
private String oid; //账号
private int cash; //账户余额
private Lock lock = new ReentrantLock(); //账户锁
private Condition _save = lock.newCondition(); //存款条件
private Condition _draw = lock.newCondition(); //取款条件
MyCount(String oid, int cash) {
this.oid = oid;
this.cash = cash;
}
/**
* 存款
*
* @param x 操作金额
* @param name 操作人
*/
public void saving(int x, String name) {
lock.lock(); //获取锁
if (x > 0) {
cash += x; //存款
System.out.println(name + "存款" + x + ", 当前余额为" + cash);
}
_draw.signalAll(); //唤醒所有等待线程.
lock.unlock(); //释放锁
}
/**
* 取款
*
* @param x 操作金额
* @param name 操作人
*/
public void drawing(int x, String name) {
lock.lock(); //获取锁
try {
if (cash - x < 0) {
_draw.await(); //阻塞取款操作
} else {
cash -= x; //取款
System.out.println(name + "取款" + x + ", 当前余额为" + cash);
}
_save.signalAll(); //唤醒所有存款操作
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
lock.unlock(); //释放锁
}
}
}

CountDownLatch

参考来源:未引入样例

一个同步辅助类,在完成一组正在其他线程中执行的操作之前,它允许一个或多个线程一直等待。用给定的计数初始化 CountDownLatch。由于调用了 countDown() 方法,所以在当前计数到达零之前,await 方法会一直受阻塞。之后,会释放所有等待的线程,await 的所有后续调用都将立即返回。这种现象只出现一次——计数无法被重置。 一个线程(或者多个), 等待另外N个线程完成某个事情之后才能执行

在一些应用场合中,需要等待某个条件达到要求后才能做后面的事情;同时当线程都完成后也会触发事件,以便进行后面的操作。 这个时候就可以使用CountDownLatch。CountDownLatch最重要的方法是countDown()和await(),前者主要是倒数一次,后者是等待倒数到0,如果没有到达0,就只有阻塞等待了。

主要方法

countDown

1
2
3
4
5
6
public void countDown()

递减锁存器的计数,如果计数到达零,则释放所有等待的线程。如果当前计数大于零,则将计数减少。
如果新的计数为零,出于线程调度目的,将重新启用所有的等待线程。

如果当前计数等于零,则不发生任何操作。

await

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
public boolean await(long timeout,
TimeUnit unit)
throws InterruptedException

使当前线程在锁存器倒计数至零之前一直等待,除非线程被中断或超出了指定的等待时间。
如果当前计数为零,则此方法立刻返回 true 值。
如果当前计数大于零,则出于线程调度目的,将禁用当前线程,
且在发生以下三种情况之一前,该线程将一直处于休眠状态:

由于调用 countDown() 方法,计数到达零;或者其他某个线程中断当前线程;
或者已超出指定的等待时间。

如果计数到达零,则该方法返回 true 值。

如果当前线程:

在进入此方法时已经设置了该线程的中断状态;或者在等待时被中断,
则抛出 InterruptedException,并且清除当前线程的已中断状态。
如果超出了指定的等待时间,则返回值为 false。如果该时间小于等于零,则此方法根本不会等待。

参数:
timeout - 要等待的最长时间
unit - timeout 参数的时间单位。
返回:
如果计数到达零,则返回 true;如果在计数到达零之前超过了等待时间,则返回 false
抛出:
InterruptedException - 如果当前线程在等待时被中断

在控制访问量上的应用

信号量 Semaphore

Java的信号量实际上是一个功能完备的计数器, 对控制一定资源的消费与回收有着很重要的意义, 信号量常常用于多线程的代码中, 并能监控有多少数目的线程等待获取资源, 并且通过信号量可以得知可用资源的数目等等, 这里总是在强调“数目”二字, 但不能指出来有哪些在等待, 哪些资源可用.

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
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Semaphore;

public class Test {
public static void main(String[] args) {
MyPool myPool = new MyPool(20);
//创建线程池
ExecutorService threadPool = Executors.newFixedThreadPool(2);
MyThread t1 = new MyThread("任务A", myPool, 3);
MyThread t2 = new MyThread("任务B", myPool, 12);
MyThread t3 = new MyThread("任务C", myPool, 7);
//在线程池中执行任务
threadPool.execute(t1);
threadPool.execute(t2);
threadPool.execute(t3);
//关闭池
threadPool.shutdown();
}
}

/**
* 一个池
*/
class MyPool {
private Semaphore sp; //池相关的信号量
/**
* 池的大小, 这个大小会传递给信号量
* @param size 池的大小
*/
MyPool(int size) {
this.sp = new Semaphore(size);
}
public Semaphore getSp() {
return sp;
}
public void setSp(Semaphore sp) {
this.sp = sp;
}
}

class MyThread extends Thread {
private String threadname; //线程的名称
private MyPool pool; //自定义池
private int x; //申请信号量的大小
MyThread(String threadname, MyPool pool, int x) {
this.threadname = threadname;
this.pool = pool;
this.x = x;
}
public void run() {
try {
//从此信号量获取给定数目的许可
pool.getSp().acquire(x);
//TODO 也许这里可以做更复杂的业务
System.out.println(threadname + "成功获取了" + x + "个许可!");
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
//释放给定数目的许可, 将其返回到信号量.
pool.getSp().release(x);
System.out.println(threadname + "释放了" + x + "个许可!");
}
}
}

CyclicBarrier

原子量

所谓的原子量即操作变量的操作是“原子的”, 该操作不可再分, 因此是线程安全的.
为何要使用原子变量呢, 原因是多个线程对单个变量操作也会引起一些问题.
Java5之后, 专门提供了用来进行单变量多线程并发安全访问的工具包java.util.concurrent.atomic, 其中的类也很简单.

1
2
3
4
5
//原子量, 每个线程都可以自由操作
private static AtomicLong aLong = new AtomicLong(10000);
lock.lock();
System.out.println(name + "执行了" + x + ", 当前余额:" + aLong.addAndGet(x));
lock.unlock();

可重入锁

如果锁具备可重入性,则称作为可重入锁。像synchronized和ReentrantLock都是可重入锁,
可重入性 实际上表明了锁的分配机制:基于线程的分配,而不是基于方法调用的分配。
举个简单的例子,当一个线程执行到某个synchronized方法时,比如说method1,而在method1中会调用另外一个synchronized方法method2,
此时线程不必重新去申请锁,而是可以直接执行方法method2。

看下面这段代码就明白了:

1
2
3
4
5
6
7
class MyClass {
public synchronized void method1() {
method2();
}
public synchronized void method2() {
}
}

两个方法method1和method2都用synchronized修饰了,假如某一时刻,线程A执行到了method1,此时线程A获取了这个对象的锁,
而由于method2也是synchronized方法,假如synchronized不具备可重入性,此时线程A需要重新申请锁。但是这就会造成一个问题,
因为线程A已经持有了该对象的锁,而又在申请获取该对象的锁,这样就会线程A一直等待永远不会获取到的锁。

由于synchronized和Lock都具备可重入性,所以不会发生上述现象。

可中断锁

可中断锁:顾名思义,就是可以相应中断的锁。

在Java中,synchronized就不是可中断锁,而Lock是可中断锁。
如果某一线程A正在执行锁中的代码,另一线程B正在等待获取该锁,可能由于等待时间过长,线程B不想等待了,想先处理其他事情,
可以让它中断自己或者在别的线程中中断它,这种就是可中断锁。
lockInterruptibly()的用法时已经体现了Lock的可中断性。

公平锁

公平锁即尽量以请求锁的顺序来获取锁。比如同是有多个线程在等待一个锁,当这个锁被释放时,等待时间最久的线程(最先请求的线程)会获得该所,这种就是公平锁。

非公平锁即无法保证锁的获取是按照请求锁的顺序进行的。这样就可能导致某个或者一些线程永远获取不到锁。
在Java中,synchronized就是非公平锁,它无法保证等待的线程获取锁的顺序。

而对于ReentrantLock和ReentrantReadWriteLock,它默认情况下是非公平锁,但是可以设置为公平锁。

看一下这2个类的源代码就清楚了:

在ReentrantLock中定义了2个静态内部类,一个是NotFairSync,一个是FairSync,分别用来实现非公平锁和公平锁。

我们可以在创建ReentrantLock对象时,通过以下方式来设置锁的公平性:

1
ReentrantLock lock = new ReentrantLock(true);

如果参数为true表示为公平锁,为false为非公平锁。默认情况下,如果使用无参构造器,则是非公平锁。

另外在ReentrantLock类中定义了很多方法,比如:

1
2
3
4
isFair() //判断锁是否是公平锁
isLocked() //判断锁是否被任何线程获取了
isHeldByCurrentThread() //判断锁是否被当前线程获取了
hasQueuedThreads() //判断是否有线程在等待该锁

在ReentrantReadWriteLock中也有类似的方法,同样也可以设置为公平锁和非公平锁。不过要记住,ReentrantReadWriteLock并未实现Lock接口,它实现的是ReadWriteLock接口。

读写锁

读写锁将对一个资源(比如文件)的访问分成了2个锁,一个读锁和一个写锁。
正因为有了读写锁,才使得多个线程之间的读操作不会发生冲突。
ReadWriteLock就是读写锁,它是一个接口,ReentrantReadWriteLock实现了这个接口。
可以通过readLock()获取读锁,通过writeLock()获取写锁。

减少锁带来的开销

使用Java API中的并发类库

可以采用java.util.concurrent等包下面的并发类,通常它们已经经过了充分的优化,能有效地支持高并发环境下的操作,并发类中大量采用了非阻塞算法,有些利用了CAS实现无锁。这里有一个小提示:使用并发哈希表时应优先采用ConcurrentHashMap而不是Hashtable,前者通过分解锁的方法使得效率更高。

用CAS代替锁

和第一点中提到的一样,CAS可以减小锁的开销,但是CAS本身是基于轮询的操作,实际使用中反而可能增加开销,这一点需要实验来测试。

减小锁的粒度

尽可能缩小锁定的范围,可以从两个方面入手。第一是只锁定所需的对象,少用synchronized(this)。第二是尽可能缩小锁定的方法块,缩小临界区大小,避免将不必要的操作也归入临界区。

拆分锁

将普通的对象锁、互斥锁按照场景拆分为读写锁或像ConcurrentHashMap一样拆分为若干把锁。

利用写时复制

对于读操作次数远远大于写操作的场景,可以在读操作时不加锁,写操作时利用写时复制来完成,但是内存占用会相应上升。

其次,系统中的线程数最好不是固定的,而是按CPU数来计算,这样当CPU增加时,相应的系统会自动增加线程提高并发率。
最后,如果对系统的CPU使用率还不满意,应当考虑分解一些单线程任务,改为多线程并发执行,以提高效率。
增加内存时,需要进行如下调整:
首先类似于线程数按CPU数计算,将Cache大小按内存大小计算,扩展内存后,Cache可以自动增长大小。
其次,可以分配更大的JVM堆内存给虚拟机,能减少OOM发生的几率。


【参考文献】:

Sun在Java5中, 对Java线程的类库做了大量的扩展, 其中 线程池 就是Java5的新特征之一

线程池使用

  • 线程池的作用:
    线程池作用就是限制系统中执行线程的数量. 根据系统的环境情况, 可以自动或手动设置线程数量, 达到运行的最佳效果;
    少了浪费系统资源, 多了造成系统拥挤效率不高. 用线程池控制线程数量, 其他线程排队等候.
    一个任务执行完毕, 再从队列的中取最前面的任务开始执行. 若队列中没有等待进程, 线程池的这一资源处于等待.
    当一个新任务需要运行时, 如果线程池中有等待的工作线程, 就可以开始运行了; 否则进入等待队列.
  • 为什么要用线程池:
    1. 减少了创建和销毁线程的次数, 每个工作线程都可以被重复利用, 可执行多个任务.
    2. 可以根据系统的承受能力, 调整线程池中工作线线程的数目, 防止因为消耗过多的内存, 而把服务器累趴下(每个线程需要大约1MB内存, 线程开的越多, 消耗的内存也就越大, 最后死机).

在Java5中, 需要了解java.util.concurrent.Executors类的API, 这个类提供大量创建线程池的静态方法, 是必须掌握的.
Java里面线程池的顶级接口是Executor, 但是严格意义上讲 Executor 并不是一个线程池, 而只是一个执行线程的工具. 真正的线程池接口是 ExecutorService.

比较重要的几个类:

  • ExecutorService 真正的线程池接口.
  • ScheduledExecutorService 能和Timer/TimerTask类似, 解决那些需要任务重复执行的问题.
  • ThreadPoolExecutor ExecutorService的默认实现.
  • ScheduledThreadPoolExecutor 继承ThreadPoolExecutorScheduledExecutorService接口实现, 周期性任务调度的类实现.

要配置一个线程池是比较复杂的, 在Executors类里面提供了一些 静态工厂 用于简便的生成一些常用的线程池.

  1. newSingleThreadExecutor
    创建一个单线程的线程池. 这个线程池只有一个线程在工作, 也就是相当于单线程串行执行所有任务. 如果这个唯一的线程因为异常结束, 那么会有一个新的线程来替代它. 此线程池保证所有任务的执行顺序按照任务的提交顺序执行.
  2. newFixedThreadPool
    创建固定大小的线程池. 每次提交一个任务就创建一个线程, 直到线程达到线程池的最大大小. 线程池的大小一旦达到最大值就会保持不变, 如果某个线程因为执行异常而结束, 那么线程池会补充一个新线程.
  3. newCachedThreadPool
    创建一个可缓存的线程池. 如果线程池的大小超过了处理任务所需要的线程, 那么就会回收部分空闲(60秒不执行任务)的线程, 当任务数增加时, 此线程池又可以智能的添加新线程来处理任务. 此线程池不会对线程池大小做限制, 线程池大小完全依赖于操作系统(或者说JVM)能够创建的最大线程大小.
  4. newScheduledThreadPool
    创建一个大小无限的线程池. 此线程池支持定时以及周期性执行任务的需求.

可重用固定线程数的线程池

Executors.newFixedThreadPool

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
//创建一个可重用固定线程数的线程池
ExecutorService pool = Executors.newFixedThreadPool(2);
//创建实现了Runnable接口对象, Thread对象当然也实现了Runnable接口
Thread t1 = new MyThread();
Thread t2 = new MyThread();
Thread t3 = new MyThread();
Thread t4 = new MyThread();
Thread t5 = new MyThread();
//将线程放入池中进行执行
pool.execute(t1);
pool.execute(t2);
pool.execute(t3);
pool.execute(t4);
pool.execute(t5);
//虽然放入了5个Runnable, 但必须在两个线程中循环执行.
//关闭线程池
pool.shutdown();

单个worker线程的Executor

Executors.newSingleThreadExecutor

创建一个使用单个 worker 线程的 Executor, 以无界队列方式来运行该线程.

1
ExecutorService pool = Executors.newSingleThreadExecutor();

以上两种连接池, 大小都是固定的, 当要加入的池的线程(或者任务)超过池最大尺寸时候, 则入此线程池需要排队等待.
一旦池中有线程完毕, 则排队等待的某个线程会入池执行. 如果当前线程在执行任务时突然中断, 则会创建一个新的线程替代它继续执行任务

可变尺寸的线程池

Executors.newCachedThreadPool

创建一个可根据需要创建新线程的线程池, 但是在以前构造的线程可用时将重用它们.

1
ExecutorService pool = Executors.newCachedThreadPool();

延迟连接池

Executors.newScheduledThreadPool

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
//创建一个线程池, 它可安排在给定延迟后运行命令或者定期地执行.
ScheduledExecutorService pool = Executors.newScheduledThreadPool(2);
Thread t1 = new MyThread();
Thread t2 = new MyThread();
Thread t3 = new MyThread();
Thread t4 = new MyThread();
Thread t5 = new MyThread();
//将线程放入池中进行执行
pool.execute(t1);
pool.execute(t2);
pool.execute(t3);

//使用延迟执行风格的方法
pool.schedule(t4, 10, TimeUnit.MILLISECONDS);
pool.schedule(t5, 10, TimeUnit.MILLISECONDS);

//关闭线程池
pool.shutdown();

单任务延迟连接池

Executors.newSingleThreadScheduledExecutor

1
2
//创建一个单线程执行程序, 它可安排在给定延迟后运行命令或者定期地执行.
ScheduledExecutorService pool = Executors.newSingleThreadScheduledExecutor();

线程池原理

在实现上, 线程池包含两个部分: worker线程组和等待队列. 当没有等待队列时, 线程都不能等待或缓冲, 所有的线程在放入的时候, 即决定开始执行或不执行. 其他情况下, 均是worker线程组从等待队列中取”任务”.

ThreadPoolExecutor

1
2
3
4
//创建等待队列
BlockingQueue<Runnable> bqueue = new ArrayBlockingQueue<Runnable>(20);
//创建一个单线程执行程序, 它可安排在给定延迟后运行命令或者定期地执行.
ThreadPoolExecutor pool = new ThreadPoolExecutor(2,3,2,TimeUnit.MILLISECONDS,bqueue);

创建自定义线程池的构造方法很多, 本例中参数的含义如下:

1
2
3
4
5
6
ThreadPoolExecutor
public ThreadPoolExecutor(int corePoolSize,
int maximumPoolSize,
long keepAliveTime,
TimeUnit unit,
BlockingQueue<Runnable> workQueue)

用给定的初始参数和默认的线程工厂及处理程序创建新的 ThreadPoolExecutor.
使用 Executors 工厂方法之一比使用此通用构造方法方便得多.

  • 参数:
    • corePoolSize - 池中所保存的线程数, 包括空闲线程.
    • maximumPoolSize - 池中允许的最大线程数.
    • keepAliveTime - 当线程数大于核心时, 此为终止前多余的空闲线程等待新任务的最长时间.
    • unit - keepAliveTime参数的时间单位.
    • workQueue - 执行前用于保持任务的队列. 此队列仅保持由 execute 方法提交的 Runnable 任务.
  • 抛出:
    • IllegalArgumentException - 如果 corePoolSize 或 keepAliveTime 小于零, 或者 maximumPoolSize 小于或等于零, 或者 corePoolSize 大于 maximumPoolSize.
    • NullPointerException - 如果 workQueue 为 null

自定义连接池稍微麻烦些, 不过通过创建的ThreadPoolExecutor线程池对象, 可以获取到当前线程池的尺寸、正在执行任务的线程数、工作队列等等

源码

下面介绍一下几个类的源码:

  1. ExecutorService newFixedThreadPool(int nThreads): 固定大小线程池.
    可以看到, corePoolSizemaximumPoolSize的大小是一样的(实际上, 后面会介绍, 如果使用无界queue的话maximumPoolSize参数是没有意义的), keepAliveTimeunit的设值表明什么?-就是该实现不想keep alive!最后的BlockingQueue选择了LinkedBlockingQueue, 该queue有一个特点, 他是无界的.
1
2
3
4
5
public static ExecutorService newFixedThreadPool(int nThreads) {
return new ThreadPoolExecutor(nThreads, nThreads,
0L, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<Runnable>());
}
  1. ExecutorService newSingleThreadExecutor(): 单线程
1
2
3
4
5
6
public static ExecutorService newSingleThreadExecutor() {
return new FinalizableDelegatedExecutorService
(new ThreadPoolExecutor(1, 1,
0L, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<Runnable>()));
}
  1. ExecutorService newCachedThreadPool(): 无界线程池, 可以进行自动线程回收
    这个实现就有意思了. 首先是无界的线程池, 所以我们可以发现maximumPoolSize为big big. 其次BlockingQueue的选择上使用SynchronousQueue. 可能对于该BlockingQueue有些陌生, 简单说: 该QUEUE中, 每个插入操作必须等待另一个线程的对应移除操作.
1
2
3
4
5
public static ExecutorService newCachedThreadPool() {
return new ThreadPoolExecutor(0, Integer.MAX_VALUE,
60L, TimeUnit.SECONDS,
new SynchronousQueue<Runnable>());
}

阻塞队列

先从BlockingQueue<Runnable> workQueue这个入参开始说起. 在JDK中, 其实已经说得很清楚了, 一共有三种类型的queue.

所有 BlockingQueue 都可用于传输和保持提交的任务. 可以使用此队列与池大小进行交互:

如果运行的线程少于 corePoolSize, 则 Executor始终首选添加新的线程, 而不进行排队. (如果当前运行的线程小于corePoolSize, 则任务根本不会存放, 添加到queue中, 而是直接抄家伙(thread)开始运行)

如果运行的线程等于或多于 corePoolSize, 则 Executor 始终首选将请求加入队列, 而不添加新的线程.
如果无法将请求加入队列, 则创建新的线程, 除非创建此线程超出 maximumPoolSize, 在这种情况下, 任务将被拒绝.
queue上的三种类型.

排队有三种通用策略:

  • 直接提交. 工作队列的默认选项是 SynchronousQueue, 它将任务直接提交给线程而不保持它们. 在此, 如果不存在可用于立即运行任务的线程, 则试图把任务加入队列将失败, 因此会构造一个新的线程. 此策略可以避免在处理可能具有内部依赖性的请求集时出现锁. 直接提交通常要求无界 maximumPoolSizes 以避免拒绝新提交的任务. 当命令以超过队列所能处理的平均数连续到达时, 此策略允许无界线程具有增长的可能性.
  • 无界队列. 使用无界队列(例如, 不具有预定义容量的 LinkedBlockingQueue)将导致在所有 corePoolSize 线程都忙时新任务在队列中等待. 这样, 创建的线程就不会超过 corePoolSize. (因此, maximumPoolSize的值也就无效了. )当每个任务完全独立于其他任务, 即任务执行互不影响时, 适合于使用无界队列; 例如, 在 Web页服务器中. 这种排队可用于处理瞬态突发请求, 当命令以超过队列所能处理的平均数连续到达时, 此策略允许无界线程具有增长的可能性.
  • 有界队列. 当使用有限的 maximumPoolSizes时, 有界队列(如 ArrayBlockingQueue)有助于防止资源耗尽, 但是可能较难调整和控制. 队列大小和最大池大小可能需要相互折衷: 使用大型队列和小型池可以最大限度地降低 CPU 使用率、操作系统资源和上下文切换开销, 但是可能导致人工降低吞吐量. 如果任务频繁阻塞(例如, 如果它们是 I/O边界), 则系统可能为超过您许可的更多线程安排时间. 使用小型队列通常要求较大的池大小, CPU使用率较高, 但是可能遇到不可接受的调度开销, 这样也会降低吞吐量.

BlockingQueue的选择.

  • 例子一: 使用直接提交策略, 也即SynchronousQueue.
    首先SynchronousQueue是无界的, 也就是说他存数任务的能力是没有限制的, 但是由于该Queue本身的特性, 在某次添加元素后必须等待其他线程取走后才能继续添加. 在这里不是核心线程便是新创建的线程, 但是我们试想一样下, 下面的场景.
    我们使用一下参数构造ThreadPoolExecutor:
1
2
3
4
5
new ThreadPoolExecutor(
2, 3, 30, TimeUnit.SECONDS,
new SynchronousQueue<Runnable>(),
new RecorderThreadFactory("CookieRecorderPool"),
new ThreadPoolExecutor.CallerRunsPolicy());

当核心线程已经有2个正在运行.

  • 此时继续来了一个任务(A), 根据前面介绍的“如果运行的线程等于或多于 corePoolSize, 则 Executor始终首选将请求加入队列, 而不添加新的线程. ”,所以A被添加到queue中.
  • 又来了一个任务(B), 且核心2个线程还没有忙完, OK, 接下来首先尝试1中描述, 但是由于使用的SynchronousQueue, 所以一定无法加入进去.
    • 此时便满足了上面提到的“如果无法将请求加入队列, 则创建新的线程, 除非创建此线程超出maximumPoolSize, 在这种情况下, 任务将被拒绝. ”, 所以必然会新建一个线程来运行这个任务.
    • 暂时还可以, 但是如果这三个任务都还没完成, 连续来了两个任务, 第一个添加入queue中, 后一个呢?queue中无法插入, 而线程数达到了maximumPoolSize, 所以只好执行异常策略了.

所以在使用SynchronousQueue通常要求maximumPoolSize是无界的, 这样就可以避免上述情况发生(如果希望限制就直接使用有界队列). 对于使用SynchronousQueue的作用jdk中写的很清楚: 此策略可以避免在处理可能具有内部依赖性的请求集时出现锁.
什么意思?如果你的任务A1, A2有内部关联, A1需要先运行, 那么先提交A1, 再提交A2, 当使用SynchronousQueue我们可以保证, A1必定先被执行, 在A1么有被执行前, A2不可能添加入queue中.

  • 例子二: 使用无界队列策略, 即LinkedBlockingQueue
    这个就拿newFixedThreadPool来说, 根据前文提到的规则:
    如果运行的线程少于 corePoolSize, 则 Executor 始终首选添加新的线程, 而不进行排队. 那么当任务继续增加, 会发生什么呢?
    如果运行的线程等于或多于 corePoolSize, 则 Executor 始终首选将请求加入队列, 而不添加新的线程. OK, 此时任务变加入队列之中了, 那什么时候才会添加新线程呢?
    如果无法将请求加入队列, 则创建新的线程, 除非创建此线程超出 maximumPoolSize, 在这种情况下, 任务将被拒绝. 这里就很有意思了, 可能会出现无法加入队列吗?不像SynchronousQueue那样有其自身的特点, 对于无界队列来说, 总是可以加入的(资源耗尽, 当然另当别论). 换句说, 永远也不会触发产生新的线程!corePoolSize大小的线程数会一直运行, 忙完当前的, 就从队列中拿任务开始运行. 所以要防止任务疯长, 比如任务运行的实行比较长, 而添加任务的速度远远超过处理任务的时间, 而且还不断增加, 不一会儿就爆了.

  • 例子三: 有界队列, 使用ArrayBlockingQueue.
    这个是最为复杂的使用, 所以JDK不推荐使用也有些道理. 与上面的相比, 最大的特点便是可以防止资源耗尽的情况发生.
    举例来说, 请看如下构造方法:

1
2
3
4
5
new ThreadPoolExecutor(
2, 4, 30, TimeUnit.SECONDS,
new ArrayBlockingQueue<Runnable>(2),
new RecorderThreadFactory("CookieRecorderPool"),
new ThreadPoolExecutor.CallerRunsPolicy());

假设, 所有的任务都永远无法执行完.
对于首先来的A,B来说直接运行, 接下来, 如果来了C,D, 他们会被放到queue中, 如果接下来再来E,F, 则增加线程运行E, F. 但是如果再来任务, 队列无法再接受了, 线程数也到达最大的限制了, 所以就会使用拒绝策略来处理.

keepAliveTime

jdk中的解释是: 当线程数大于核心时, 此为终止前多余的空闲线程等待新任务的最长时间.
有点拗口, 其实这个不难理解, 在使用了“池”的应用中, 大多都有类似的参数需要配置. 比如数据库连接池, DBCP中的maxIdle, minIdle参数.
什么意思?接着上面的解释, 后来向老板派来的工人始终是“借来的”, 俗话说“有借就有还”, 但这里的问题就是什么时候还了, 如果借来的工人刚完成一个任务就还回去, 后来发现任务还有, 那岂不是又要去借?这一来一往, 老板肯定头也大死了.

合理的策略: 既然借了, 那就多借一会儿. 直到“某一段”时间后, 发现再也用不到这些工人时, 便可以还回去了. 这里的某一段时间便是keepAliveTime的含义, TimeUnit为keepAliveTime值的度量.

RejectedExecutionHandler
另一种情况便是, 即使向老板借了工人, 但是任务还是继续过来, 还是忙不过来, 这时整个队伍只好拒绝接受了.
RejectedExecutionHandler接口提供了对于拒绝任务的处理的自定方法的机会. 在ThreadPoolExecutor中已经默认包含了4中策略, 因为源码非常简单, 这里直接贴出来.
CallerRunsPolicy: 线程调用运行该任务的 execute 本身. 此策略提供简单的反馈控制机制, 能够减缓新任务的提交速度.

1
2
3
4
5
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
if (!e.isShutdown()) {
r.run();
}
}

这个策略显然不想放弃执行任务. 但是由于池中已经没有任何资源了, 那么就直接使用调用该execute的线程本身来执行.

AbortPolicy:
处理程序遭到拒绝将抛出运行时RejectedExecutionException

1
2
3
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
throw new RejectedExecutionException();
}

这种策略直接抛出异常, 丢弃任务.

DiscardPolicy: 不能执行的任务将被删除

1
2
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
}

这种策略和AbortPolicy几乎一样, 也是丢弃任务, 只不过他不抛出异常.
DiscardOldestPolicy: 如果执行程序尚未关闭, 则位于工作队列头部的任务将被删除, 然后重试执行程序(如果再次失败, 则重复此过程)

1
2
3
4
5
6
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
if (!e.isShutdown()) {
e.getQueue().poll();
e.execute(r);
}
}

该策略就稍微复杂一些, 在pool没有关闭的前提下首先丢掉缓存在队列中的最早的任务, 然后重新尝试运行该任务. 这个策略需要适当小心.
设想:如果其他线程都还在运行, 那么新来任务踢掉旧任务, 缓存在queue中, 再来一个任务又会踢掉queue中最老任务.

总结:
keepAliveTime和maximumPoolSize及BlockingQueue的类型均有关系. 如果BlockingQueue是无界的, 那么永远不会触发maximumPoolSize, 自然keepAliveTime也就没有了意义.
反之, 如果核心数较小, 有界BlockingQueue数值又较小, 同时keepAliveTime又设的很小, 如果任务频繁, 那么系统就会频繁的申请回收线程.

1
2
3
4
5
public static ExecutorService newFixedThreadPool(int nThreads) {
return new ThreadPoolExecutor(nThreads, nThreads,
0L, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<Runnable>());
}

有返回值的线程池

可返回值的任务必须实现Callable接口, 类似的, 无返回值的任务必须Runnable接口.
执行Callable取到Callable任务返回的Object了.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
//创建一个线程池
ExecutorService pool = Executors.newFixedThreadPool(2);
//创建两个有返回值的任务
Callable c1 = new MyCallable("A");
Callable c2 = new MyCallable("B");
//执行任务并获取Future对象
Future f1 = pool.submit(c1);
Future f2 = pool.submit(c2);
//从Future对象上获取任务的返回值, 并输出到控制台
System.out.println(">>>"+f1.get().toString());
System.out.println(">>>"+f2.get().toString());
//关闭线程池
pool.shutdown();

class MyCallable implements Callable{

}

【参考文献】:

  1. Java线程池使用说明

2个线程轮流输出

A、B两个线程轮流输出1~100的数字

A线程输出: 1、3、5…47、49、 52、54…98、100

B线程输出: 2、4、6…48、50、51、53、55…97、99 两个线程输出个数相同

分析:

  1. 每个线程持有一把锁,运行时,首先获取自己的锁,运行完释放另一个线程的锁
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
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
import java.util.concurrent.Semaphore;

/**
* A、B两个线程轮流输出1~100的数字<br/>
* A线程输出: 1、3、5...47、49、 52、54...98、100 <br/>
* B线程输出: 2、4、6...48、50、51、53、55...97、99 两个线程输出个数相同
*/
public class RunningInTurn {
/**
* 每个线程
*/
static class Worker extends Thread {
private final String name;
private final Semaphore thisSemaphore;
private final Semaphore nextSemaphore;
private int value;

public Worker(String name, Semaphore thisSemaphore, Semaphore nextSemaphore, int initialValue) {
this.name = name;
this.thisSemaphore = thisSemaphore;
this.nextSemaphore = nextSemaphore;
this.value = initialValue;
}

@Override
public void run() {
int cnt = 0;
while (value <= 100) {
try {
thisSemaphore.acquire();
System.out.println(name + ":\t" + value);
cnt++;
if (value == 50) {
value = 51;
System.out.println(name + ":\t" + value);
cnt++;
value += 2;
} else if (value == 49) {
value = 52;
} else {
value += 2;
}
nextSemaphore.release();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
System.out.println(name + " cnt : " + cnt);
}
}

public static void main(String[] args) throws InterruptedException {

Semaphore aSemaphore = new Semaphore(1);
Semaphore bSemaphore = new Semaphore(1);
Worker workerA = new Worker("a", aSemaphore, bSemaphore, 1);
Worker workerB = new Worker("b", bSemaphore, aSemaphore, 2);
bSemaphore.acquire();
workerA.start();
workerB.start();
}

}

理解Java多线程, 需要深入理解`线程状态和锁`

[TOC]

线程

由于每个时钟周期内, CPU 实际上只能执行一条指令. CPU每一个时刻只能做一件事, 多线程是通过任务调度给CPU分配任务实现的, 多线程的目的是为了最大限度的利用CPU资源.

操作系统负责管理进程和线程, 轮流(没有固定的顺序)分配每个进程很短的时间(不一定是均分), 然后在每个线程内部, 程序代码自己处理该进程内部线程的时间分配, 多个线程之间相互的切换去执行, 这个切换时间也是非常短的.

程序、进程、线程之间的关系

  • 程序是一段静态的代码,是应用软件执行的蓝本。
  • 进程是程序一次动态执行的过程,它对应了从代码加载、执行完毕的一个完整过程,这也是进程开始到消亡的过程。
  • 线程是进程中独立、可调度的执行单元,是执行中最小单位。
  • 一个程序一般是一个进程,但一个程序中也可以有多个进程。
  • 一个进程中可以有多个线程,但只有一个主线程。
  • Java应用程序中默认的主线程是main方法,如果main方法中创建了其他线程,JVM就会执行其他的线程。

Java 进程

Java编写的程序是运行在JVM中的, 启动一个Java应用程序, 就会启动一个JVM进程. 在同一个JVM进程中, 有且只有一个进程, 就是它自己. 因此, 所有的程序代码的运行都是以线程运行的. 同一个进程中的所有线程共享一块内存块, 因此线程间通信很容易且速度很快.

  • Java 中的线程是一个对象, 与其他 Java 中的对象一样, 具有变量和方法, 生死于堆上.

调用栈

  • Java 中的每个线程都有一个调用栈, 即使不创建任何新的线程, 线程也在后台运行着.
  • 一旦创建一个新的线程, 就产生一个新的调用栈.

主线程

在JVM上运行一个应用程序时, JVM首先寻找程序入口的main()方法, 然后运行main()方法, 此时就产生了一个Java线程, 这个线程就是主线程. 当main方法结束后, 主线程运行完成, 如果不存在额外的线程运行, JVM进程随即退出.

调度的方式有两种:分时调度和抢占式调度, Java中采用的是抢占式调度

线程优先级

1
2
3
4
5
6
static int MAX_PRIORITY
线程可以具有的最高优先级.
static int MIN_PRIORITY
线程可以具有的最低优先级.
static int NORM_PRIORITY
分配给线程的默认优先级.

当线程池中线程都具有相同的优先级, 调度程序的JVM实现自由选择它喜欢的线程. 这时候调度程序的操作有两种可能:
一是选择一个线程运行, 直到它阻塞或者运行完成为止.
二是时间分片, 为池内的每个线程提供均等的运行机会.

1~10之间的值是没有保证的. 一些JVM可能不能识别10个不同的值, 而将这些优先级进行每两个或多个合并, 变成少于10个的优先级, 则两个或多个优先级的线程可能被映射为一个优先级.

与线程休眠类似, 线程的优先级仍然无法保障线程的执行次序. 只不过, 优先级高的线程获取CPU资源的概率较大, 优先级低的并非没机会执行.

线程的生命周期

创建–运行–中断–死亡

  • 创建:线程构造
  • 运行:调用start()方法,进入run()方法
  • 中断:sleep()、wait()
  • 死亡:执行完run()方法或强制run()方法结束,线程死亡

线程的状态与转换

线程生命周期与状态转换

new与新建状态

新建线程有两个方法: 继承Thread类 和 实现Runnable接口.

Runnable 接口的定义如下:

1
2
3
public interface Runnable {
public void run();
}

实现Runnable接口就需要实现run方法. run方法的内容就是线程要执行的任务.
Thread类, 是实现了Runable接口的类, 因此在Thread类中也存在run方法.

新建线程有若干种重载方法:

1
2
3
4
5
6
7
8
public Thread(Runnable target) {
//
}
public Thread(String name) {
//
}
//由于 Thread 实现了 Runnable 接口,
//这种形式也可以以 Thread 对象为参数

start 启动与就绪状态(即可运行状态)

调用Thread.start()方法, 该线程进入Runable(可运行)状态, 等待分配 CPU 资源, 当抢占到CPU资源时,
该线程启动,开始执行run方法.

Thread.start()是唯一可以新建线程的方法. 执行 Thread.run() 和 Runable.run() 只会执行run方法, 不会启动新的线程.

一旦线程启动, 它就永远不能再重新启动. 只有一个新的线程可以被启动, 并且只能一次. 一个可运行的线程或死线程可以被重新启动.

线程的调度是JVM的一部分, 在一个CPU的机器上, 实际上一次只能运行一个线程. 一次只有一个线程栈执行. JVM线程调度程序决定实际运行哪个处于可运行状态的线程. 众多可运行线程中的某一个会被选中作为当前线程. 可运行线程被选择运行的顺序是没有保障的. 尽管通常采用队列形式, 但这是没有保障的. 队列形式是指当一个线程完成“一轮”时, 它移到可运行队列的尾部等待, 直到它最终排队到该队列的前端为止, 它才能被再次选中. 事实上, 我们把它称为可运行池而不是一个可运行队列, 目的是帮助认识线程并不都是以某种有保障的顺序排列成一个队列的事实.

Running 运行

运行状态, 执行run()方法的内容.

当 Java 虚拟机继续执行线程, 直到下面任一情况出现为止:

  • 调用 Runtime的 exit 方法 System.exit()
  • 非守护线程全部停止运行, 无论是从 run 方法返回还是通过抛出一个传播到 run 方法之外的异常.

几种特殊情况可能使线程离开运行状态:

  1. 线程的run()方法完成.
  2. 在对象上调用wait()方法(不是在线程上调用).
  3. 线程不能在对象上获得锁定, 它正试图运行该对象的方法代码.
  4. 线程调度程序可以决定将当前运行状态移动到可运行状态, 以便让另一个线程获得运行机会, 而不需要任何理由.

sleep() 休眠

1
2
3
4
5
6
try{
Thread.sleep();
}catch(Exception e){
}
//原始定义
public static native void sleep(long millis) throws InterruptedException;

调用 Thread 的静态方法可以进入休眠, sleep()作用:

  • 执行带锁的代码时, 不会释放锁
  • 进入Blocked状态
  • 线程在mills时间内不会醒来
  • mills时间到了, 线程变为Runnable状态
  • 再次进入运行状态, 继续执行 sleep() 后面的代码

Thread.sleep(long millis)Thread.sleep(long millis, int nanos)静态方法强制当前正在执行的线程休眠(暂停执行), 以“减慢线程”.
当线程睡眠时, 它入睡在某个地方, 在苏醒之前不会返回到可运行状态. 当睡眠时间到期, 则返回到可运行状态.
线程睡眠的原因:线程执行太快, 或者需要强制进入下一轮, 因为Java规范不保证合理的轮换.

当休眠一定时间后, 线程会苏醒, 进入准备状态等待执行.

Blocked 阻塞状态

阻塞状态是线程因为某种原因放弃CPU使用权, 暂时停止运行. 直到线程进入就绪状态, 才有机会转到运行状态. 阻塞的情况分三种:

  • 等待阻塞:运行的线程执行wait()方法, JVM会把该线程放入等待池中.
  • 同步阻塞:运行的线程在获取对象的同步锁时, 若该同步锁被别的线程占用, 则JVM会把该线程放入锁池中.
  • 其他阻塞:运行的线程执行sleep()join()方法, 或者发出了I/O请求时, JVM会把该线程置为阻塞状态.
    sleep()状态超时、join()等待线程终止或者超时、或者I/O处理完毕时, 线程重新转入就绪状态.

Synchronized锁与同步

monitor与synchronized

在 Java 中每个对象都有一个锁,(问题来了: Java对象锁信息保存在哪里?) 并且对象的锁同时只能被一个线程使用, 因此当某个线程得到对象的锁时, 其他线程也就没办法获得锁.
利用对象的锁, 可以实现只允许一个线程访问, 即同步.

当线程运行到 synchronized 时, 首先检测是否可以获得对象的锁, 如果可以获得,则马上获取锁. 如果不能获取锁, 线程阻塞, 开始等待其他线程释放锁.

当同步锁被释放时, 线程重新进入 Runnable 可运行状态.

需要同步时,一定要搞清楚加锁的对象是什么

关于锁和同步, 有一下几个要点:

  • 只能同步方法, 而不能同步变量和类;
  • 每个对象只有一个锁;当提到同步时, 应该清楚在什么上同步?也就是说, 在哪个对象上同步?
  • 不必同步类中所有的方法, 类可以同时拥有同步和非同步方法.
  • 线程睡眠(执行sleep)时, 它所持的任何锁都不会释放.
  • 线程可以获得多个锁. 比如, 在一个对象的同步方法里面调用另外一个对象的同步方法, 则获取了两个对象的同步锁.

同步方法

在方法名称前面增加 synchronized 关键字, 此时相当于以类的对象(即this)作为对象锁.

1
2
public synchronized void function(){
}

同步代码块

1
2
synchronized(object1){
}

object1作为对象锁

同步静态方法

  1. 静态方法同步是以方法所在的class对象作为锁的.
    要同步静态方法, 需要一个用于整个类对象的锁, 这个对象是就是这个类(XXX.class).
    例如:

    1
    2
    3
    4
    5
    6
    7
    8
    9
    public static synchronized int setName(String name){
    Xxx.name = name;
    }
    // 等价于
    public static int setName(String name){
    synchronized(Xxx.class){
    Xxx.name = name;
    }
    }

    实质上, 线程进入该对象的的一种池中, 必须在那里等待, 直到其锁被释放

  2. 调用同一个类中的静态同步方法的线程将彼此阻塞, 它们都是锁定在相同的Class对象上.

  3. 静态同步方法和非静态同步方法将永远不会彼此阻塞, 因为静态方法锁定在Class对象上, 非静态方法锁定在该类的对象上.
    对于非静态字段中可更改的数据, 通常使用非静态方法访问.
    对于静态字段中可更改的数据, 通常使用静态方法访问.

减少锁定时间

线程同步的目的是为了保护多个线程访问一个资源时对资源的破坏. 线程同步方法是通过锁来实现, 每个对象都有且仅有一个锁, 这个锁与一个特定的对象关联, 线程一旦获取了对象锁, 其他访问该对象的线程就无法再访问该对象的其他同步方法.

synchronized中慎用sleep和yield

在使用synchronized关键字时候, 应该尽可能避免在synchronized方法或synchronized块中使用sleep或者yield方法:
因为synchronized程序块占有着对象锁, 你休息那么其他的线程只能一边等着你醒来执行完了才能执行. 不但严重影响效率, 也不合逻辑. 同样, 在同步程序块内调用yield方法让出CPU资源也没有意义, 因为你占用着锁, 其他互斥线程还是无法访问同步程序块. 当然与同步程序块无关的线程可以获得更多的执行时间.

wait 等待 与 notify/notifyAll 通知

wait 让本线程等待, notify 通知某个线程不再等待

wait()作用主要有:

  • 释放锁
  • 不继续执行 wait 后面的代码
  • 本线程进入等待阻塞

wait必须与while一起使用:(避免假唤醒)

A thread can also wake up without being notified, interrupted, or timing out, a so-called spurious wakeup. While this will rarely occur in practice, applications must guard against it by testing for the condition that should have caused the thread to be awakened, and continuing to wait if the condition is not satisfied. In other words, waits should always occur in loops, like this one:
线程也可以在没有被通知,中断或超时的情况下唤醒,即所谓的虚假唤醒。 虽然这在实践中很少发生,但应用程序必须通过测试应该导致线程被唤醒的条件来防范它,并且如果条件不满足则继续等待。 换句话说,等待应该总是出现在循环中

1
2
3
while (!condition) {
obj.wait();
}

notify()与 notifyAll()
针对同一个对象锁上的线程, 主要的作用是:

  • 不再等待, 从等待 Blocked 中转出
  • 进入 Runnable 状态, 等待获取 CPU 资源. ??

线程唤醒 notify

Object类中的notify()方法, 唤醒在此对象监视器上等待的单个线程. 如果所有线程都在此对象上等待, 则会选择唤醒其中一个线程. 选择是任意性的, 并在对实现做出决定时发生. 线程通过调用其中一个 wait 方法, 在对象的监视器上等待. 直到当前的线程放弃此对象上的锁定, 才能继续执行被唤醒的线程. 被唤醒的线程将以常规方式与在该对象上主动同步的其他所有线程进行竞争;例如, 唤醒的线程在作为锁定此对象的下一个线程方面没有可靠的特权或劣势. 类似的方法还有一个notifyAll(), 唤醒在此对象监视器上等待的所有线程.

wait()notify()notifyAll()都是Object的实例方法. 与每个对象具有锁一样, 每个对象可以有一个线程列表, 他们等待来自该信号(通知). 线程通过执行对象上的wait()方法获得这个等待列表.

这3个方法必须处于synchronized代码块或者synchronized方法中,否则就会抛出IllegalMonitorStateException异常,这是因为这几个方法必须拿到当前对象的监视器monitor对象,也就是说notify/notifyAll和wait方法依赖于monitor对象,

千万注意

当在对象上调用wait()方法时, 执行该代码的线程立即放弃它在对象上的锁. 然而调用notify()时, 并不意味着这时线程会放弃其锁. 如果线程仍然在完成同步代码, 则线程在移出之前不会放弃锁. 因此, 只要调用notify()并不意味着这时该锁变得可用.

notifyAll() 方法, 起到的是一个通知作用,不释放锁, 也不获取锁. 只是告诉该对象上等待的线程“可以竞争执行了, 都醒来去执行吧”

yield() 让步

  • 当前线程让出, 进入Runnable可执行状态
  • 但是继续占着锁
  • 同级别或较高级别的开始竞争 CPU 资源

Thread.yield()方法作用是:暂停当前正在执行的线程对象, 并执行其他线程.
yield()应该做的是让当前运行线程回到可运行状态, 以允许具有相同优先级的其他线程获得运行机会.

yield()从未导致线程转到等待/睡眠/阻塞状态. 在大多数情况下, yield()将导致线程从运行状态转到可运行状态, 但有可能没有效果.

// TODO 使用场景

join() 合并

假设在 A 线程中,执行B.join() , B 线程放到 A 线程前面执行. A 线程转入阻塞状态首先执行 B 线程,
直到执行完 B 线程后, A 线程转入可运行状态就绪, 获取到 CPU 资源后再继续执行 join 后面的代码

还有 join() 的重载形式:

1
B.join(1000);  

首先执行 B 1000毫秒, 1000毫秒后, 即使是没有执行完, 也会停止执行 B 线程, 开始执行join 语句后面的代码

请看join的原始定义

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
43
44
45
/**
* Waits at most {@code millis} milliseconds for this thread to
* die. A timeout of {@code 0} means to wait forever.
*
* <p> This implementation uses a loop of {@code this.wait} calls
* conditioned on {@code this.isAlive}. As a thread terminates the
* {@code this.notifyAll} method is invoked. It is recommended that
* applications not use {@code wait}, {@code notify}, or
* {@code notifyAll} on {@code Thread} instances.
*
* @param millis
* the time to wait in milliseconds
*
* @throws IllegalArgumentException
* if the value of {@code millis} is negative
*
* @throws InterruptedException
* if any thread has interrupted the current thread. The
* <i>interrupted status</i> of the current thread is
* cleared when this exception is thrown.
*/
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) {
while (isAlive()) {
wait(0);
}
} else {
while (isAlive()) {
long delay = millis - now;
if (delay <= 0) {
break;
}
wait(delay);
now = System.currentTimeMillis() - base;
}
}
}

其实Join 方法实现是通过wait实现的.
当main 线程调用t.join 时候, main 线程会获得线程对象t 的锁 (wait 意味着拿到该对象的锁), 调用该对象的wait( 等待时间) ,直到该对象唤醒main 线程, 比如退出后.

stop 停止

避免使用stop(),是因为它不安全。它会解除由线程获取的所有锁定,而且如果对象处于一种不连贯状态,那么其他线程能在那种状态下检查和修改它们。结果很难检查出真正的问题所在。

suspend 暂停

suspend()方法容易发生死锁。调用suspend()的时候,目标线程会停下来,但却仍然持有在这之前获得的锁定。此时,其他任何线程都不能访问锁定的资源,除非被”挂起”的线程恢复运行。对任何线程来说,如果它们想恢复目标线程,同时又试图使用任何一个锁定的资源,就会造成死锁。所以不应该使用suspend(),而应在自己的Thread类中置入一个标志,指出线程应该活动还是挂起。若标志指出线程应该挂起,便用wait()命其进入等待状态。若标志指出线程应当恢复,则用一个notify()重新启动线程。

Thread

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
static Thread currentThread()
返回对当前正在执行的线程对象的引用.
ClassLoader getContextClassLoader()
返回该线程的上下文 ClassLoader.
Thread.State getState()
返回该线程的状态.
ThreadGroup getThreadGroup()
返回该线程所属的线程组.
static boolean holdsLock(Object obj)
当且仅当当前线程在指定的对象上保持监视器锁时, 才返回 true.
static void setDefaultUncaughtExceptionHandler(Thread.UncaughtExceptionHandler eh)
设置当线程”由于未捕获到异常而突然终止, 并且没有为该线程定义其他处理程序时”所调用的默认处理程序.
void stop()
已过时.
stop 的许多使用都应由只修改某些变量以指示目标线程应该停止运行的代码来取代.
目标线程应定期检查该变量, 并且如果该变量指示它要停止运行, 则从其运行方法依次返回.
如果目标线程等待很长时间(例如基于一个条件变量), 则应使用 interrupt 方法来中断该等待.

线程中断与synchronized

线程中断

正如中断二字所表达的意义,在线程运行(run方法)中间打断它,在Java中,提供了以下3个有关线程中断的方法

1
2
3
4
5
6
7
8
//中断线程(实例方法)
public void Thread.interrupt();

//判断线程是否被中断(实例方法)
public boolean Thread.isInterrupted();

//判断是否被中断并清除当前中断状态(静态方法)
public static boolean Thread.interrupted();

当一个线程处于被阻塞状态或者试图执行一个阻塞操作时,使用Thread.interrupt()方式中断该线程,注意此时将会抛出一个InterruptedException的异常,同时中断状态将会被复位(由中断状态改为非中断状态),如下代码将演示该过程:

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
public class InterruptSleepThread3 {
public static void main(String[] args) throws InterruptedException {
Thread t1 = new Thread() {
@Override
public void run() {
//while在try中,通过异常中断就可以退出run循环
try {
while (true) {
//当前线程处于阻塞状态,异常必须捕捉处理,无法往外抛出
TimeUnit.SECONDS.sleep(2);
}
} catch (InterruptedException e) {
System.out.println("Interruted When Sleep");
boolean interrupt = this.isInterrupted();
//中断状态被复位
System.out.println("interrupt:"+interrupt);
}
}
};
t1.start();
TimeUnit.SECONDS.sleep(2);
//中断处于阻塞状态的线程
t1.interrupt();

/**
* 输出结果:
Interruted When Sleep
interrupt:false
*/
}
}

如上述代码所示,我们创建一个线程,并在线程中调用了sleep方法从而使线程进入阻塞状态,启动线程后,调用线程实例对象的interrupt方法中断阻塞异常,并抛出InterruptedException异常,此时中断状态也将被复位。除了阻塞中断的情景,我们还可能会遇到处于运行期且非阻塞的状态的线程,这种情况下,直接调用Thread.interrupt()中断线程是不会得到任响应的,如下代码,将无法中断非阻塞状态下的线程:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
public class InterruputThread {
public static void main(String[] args) throws InterruptedException {
Thread t1=new Thread(){
@Override
public void run(){
while(true){
System.out.println("未被中断");
}
}
};
t1.start();
TimeUnit.SECONDS.sleep(2);
t1.interrupt();

/**
* 输出结果(无限执行):
未被中断
未被中断
未被中断
......
*/
}
}

虽然我们调用了interrupt方法,但线程t1并未被中断,因为处于非阻塞状态的线程需要我们手动进行中断检测并结束程序,改进后代码如下:

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
public class InterruputThread {
public static void main(String[] args) throws InterruptedException {
Thread t1=new Thread(){
@Override
public void run(){
while(true){
//判断当前线程是否被中断
if (this.isInterrupted()){
System.out.println("线程中断");
break;
}
}

System.out.println("已跳出循环,线程中断!");
}
};
t1.start();
TimeUnit.SECONDS.sleep(2);
t1.interrupt();
/**
* 输出结果:
线程中断
已跳出循环,线程中断!
*/
}
}
  • 一种是当线程处于阻塞状态或者试图执行一个阻塞操作时,我们可以使用实例方法interrupt()进行线程中断,执行中断操作后将会抛出interruptException异常(该异常必须捕捉无法向外抛出)并将中断状态复位,
  • 另外一种是当线程处于运行状态时,我们也可调用实例方法interrupt()进行线程中断,但同时必须手动判断中断状态,并编写中断线程的代码(其实就是结束run方法体的代码)。有时我们在编码时可能需要兼顾以上两种情况,那么就可以如下编写:
1
2
3
4
5
6
7
8
9
10
public void run(){
try {
//判断当前线程是否已中断,注意interrupted方法是静态的,执行后会对中断状态进行复位
while (!Thread.interrupted()) {
TimeUnit.SECONDS.sleep(2);
}
} catch (InterruptedException e) {

}
}

事实上线程的中断操作对于正在等待获取的锁对象的synchronized方法或者代码块并不起作用,也就是对于synchronized来说,如果一个线程在等待锁,那么结果只有两种,要么它获得这把锁继续执行,要么它就保存等待,即使调用中断线程的方法,也不会生效。演示代码如下

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
43
44
45
46
47
/**
* Created by zejian on 2017/6/2.
* Blog : http://blog.csdn.net/javazejian [原文地址,请尊重原创]
*/
public class SynchronizedBlocked implements Runnable{

public synchronized void f() {
System.out.println("Trying to call f()");
while(true){ // Never releases lock
Thread.yield();
}
}

/**
* 在构造器中创建新线程并启动获取对象锁
*/
public SynchronizedBlocked() {
//该线程已持有当前实例锁
new Thread() {
public void run() {
f(); // Lock acquired by this thread
}
}.start();
}
public void run() {
//中断判断
while (true) {
if (Thread.interrupted()) {
System.out.println("中断线程!!");
break;
} else {
f();
}
}
}


public static void main(String[] args) throws InterruptedException {
SynchronizedBlocked sync = new SynchronizedBlocked();
Thread t = new Thread(sync);
//启动后调用f()方法,无法获取当前实例锁处于等待状态
t.start();
TimeUnit.SECONDS.sleep(1);
//中断线程,无法生效
t.interrupt();
}
}

我们在SynchronizedBlocked构造函数中创建一个新线程并启动获取调用f()获取到当前实例锁,由于SynchronizedBlocked自身也是线程,启动后在其run方法中也调用了f(),但由于对象锁被其他线程占用,导致t线程只能等到锁,此时我们调用了t.interrupt();但并不能中断线程。

线程死锁

线程A当前持有互斥所锁 lock1 ,线程B当前持有互斥锁 lock2 。 接下来,当线程A仍然持有lock1时,它试图获取lock2,因为线程B正持有lock2,因此线程A会阻塞等待线程B对lock2的释放。如果此时线程B在持有lock2的时候,也在试图获取lock1,因为线程A正持有lock1,因此线程B会阻塞等待A对lock1的释放。二者都在等待对方所持有锁的释放,而二者却又都没释放自己所持有的锁,这时二者便会一直阻塞下去。这种情形称为死锁。

规避死锁:

  1. 只在必要的最短时间内持有锁,考虑使用同步语句块代替整个同步方法;
  2. 尽量编写不在同一时刻需要持有多个锁的代码,如果不可避免,则确保线程持有第二个锁的时间尽量短暂;
  3. 创建和使用一个大锁来代替若干小锁,并把这个锁用于互斥,而不是用作单个对象的对象级别锁;

同步方案

在需要同步的时候,第一选择应该是synchronized关键字,这是最安全的方式,尝试其他任何方式都是有风险的。尤其在jdK1.5之后,对synchronized同步机制做了很多优化,如:自适应的自旋锁、锁粗化、锁消除、轻量级锁等,使得它的性能明显有了很大的提升。

[volatile][a-volatile] 也是确保可见性的方法之一,但是不能实现原子性

互斥锁

如果同一个方法内同时有两个或更多线程,则每个线程有自己的局部变量拷贝, 采用 synchronized 修饰符实现的同步机制叫做互斥锁机制,它所获得的锁叫做互斥锁。 每个对象都有一个monitor(锁标记),当线程拥有这个锁标记时才能访问这个资源,没有锁标记便进入锁池。 任何一个对象系统都会为其创建一个互斥锁,这个锁是为了分配给线程的,防止打断原子操作。每个对象的锁只能分配给一个线程,因此叫做互斥锁。

互斥是实现同步的一种手段,临界区、互斥量和信号量都是主要的互斥实现方式。synchronized 关键字经过编译后,会在同步块的前后分别形成 monitorentermonitorexit 这两个字节码指令。根据虚拟机规范的要求,在执行monitorenter 指令时,首先要尝试获取对象的锁,如果获得了锁,把锁的计数器加1,相应地,在执行 monitorexit 指令时会将锁计数器减1,当计数器为0时,锁便被释放了。由于synchronized同步块对同一个线程是可重入的,因此一个线程可以多次获得同一个对象的互斥锁,同样,要释放相应次数的该互斥锁,才能最终释放掉该锁。

当一个线程进入一个对象的一个synchronized方法后,其它线程是否可进入此对象的其它方法? 分几种情况:

  1. 其他方法前是否加了synchronized关键字,如果没加,则能。
  2. 如果这个方法内部调用了wait,则可以进入其他synchronized方法。
  3. 如果其他个方法都加了synchronized关键字,并且内部没有调用wait,则不能。
  4. 如果其他方法是static,它用的同步锁是当前类的字节码,与非静态的方法不能同步,因为非静态的方法用的是this。

线程安全类

请查看 Java集合

守护线程

线程总体分两类:用户线程和守候线程.

当所有用户线程执行完毕的时候, JVM自动关闭. 但是守候线程却不独立于JVM, 守候线程一般是由操作系统或者用户自己创建的.

举例来说, JVM的垃圾回收、内存管理等线程都是守护线程.
还有就是在做数据库应用时候, 使用的数据库连接池, 连接池本身也包含着很多后台线程, 监控连接个数、超时时间、状态等等.

Timer

有一个后台运行的线程, 按照指定的时间执行定时任务. Timer.schedual()方法向后台线程添加定时任务,
后台线程按照的既定的时间执行定时任务.

调用Timer.cancel()取消所有已安排的定时任务, 正在的执行的任务不会被取消.

调用构造方法, 后台线程就已经启动.


[参考文献]:

AQS

CountDownLatch

CountDownLatch 的作用:当一个线程需要另外一个或多个线程完成后,再开始执行

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
final CountDownLatch countDownLatch = new CountDownLatch(2); //等待2次countdown
final Thread ta = new Thread(() -> {
doWork("A");
countDownLatch.countDown();
});
final Thread tb = new Thread(() -> {
doWork("B");
countDownLatch.countDown();
});
final Thread tc = new Thread(() -> {
try {
countDownLatch.await();
} catch (final InterruptedException e1) {
e1.printStackTrace();
}
doWork("C");
});
ta.start();
tb.start();
tc.start();

CountDownLatch实现原理

CountDownLatch是通过AQS实现的。 AQS 全称 AbstractQueuedSynchronizer,是 java.util.concurrent 中提供的一种高效且可扩展的同步机制。它可以用来实现依赖 int 状态(state)的同步器,除了CountDownLatchReentrantLockSemaphore 等功能实现都使用了它。

在调用 awit()countDown()的时候,发生了几个关键的调用关系

img

首先在 CountDownLatch 类内部定义了一个 Sync 内部类,这个内部类就是继承自 AbstractQueuedSynchronizer 的,并且重写了方法 tryAcquireSharedtryReleaseShared。当调用 awit()方法时,CountDownLatch 会调用内部类SyncacquireSharedInterruptibly() 方法,然后在这个方法中会调用 tryAcquireShared 方法,这个方法就是 CountDownLatch 的内部类 Sync 里重写的 AbstractQueuedSynchronizer 的方法。调用 countDown() 方法同理。

AQS的使用方法

使用 AbstractQueuedSynchronizer 的标准化方式,大致分为两步:

  1. 内部持有继承自 AbstractQueuedSynchronizer 的对象 Sync

  2. 并在 Sync 内重写 AbstractQueuedSynchronizerprotected 部分或全部方法,这些方法包括如下几个:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
protected boolean tryAcquire(int arg) {
throw new UnsupportedOperationException();
}

protected boolean tryRelease(int arg) {
throw new UnsupportedOperationException();
}

protected int tryAcquireShared(int arg) {
throw new UnsupportedOperationException();
}

protected boolean tryReleaseShared(int arg) {
throw new UnsupportedOperationException();
}
protected boolean isHeldExclusively() {
throw new UnsupportedOperationException();
}

使用者可以通过重写这些方法,加入自己的判断逻辑,例如 CountDownLatchtryAcquireShared中加入了判断,判断 state 是否不为0,如果不为0,才符合调用条件。

  • tryAcquiretryRelease是对应的,前者是独占模式获取,后者是独占模式释放。

  • tryAcquireSharedtryReleaseShared是对应的,前者是共享模式获取,后者是共享模式释放。

CountDownLatch 重写的方法 tryAcquireShared 实现如下:

1
2
3
protected int tryAcquireShared(int acquires) {
return (getState() == 0) ? 1 : -1;
}

判断 state 值是否为0,为0 返回1,否则返回 -1。state 值是 AbstractQueuedSynchronizer 类中的一个 volatile 变量。

1
private volatile int state;

CountDownLatch 中这个 state 值就是计数器,在调用 await 方法的时候,将值赋给 state

等待线程入队

调用 await() 方法时,先去获取 state 的值,当计数器不为0的时候,说明还有需要等待的线程在运行,则调用 doAcquireSharedInterruptibly 方法,尝试加入等待队列 ,即调用 addWaiter()方法, 源码如下:

AQS 的核心部分: AQS 用内部的一个 Node 类维护一个 CHL Node FIFO 队列。将当前线程加入等待队列,并通过 parkAndCheckInterrupt()方法实现当前线程的阻塞。下面一大部分都是在说明 CHL 队列的实现,里面用 CAS 实现队列出入不会发生阻塞。

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
private void doAcquireSharedInterruptibly(int arg)
throws InterruptedException {
//加入等待队列
final Node node = addWaiter(Node.SHARED);
boolean failed = true;
// 进入 CAS 循环
try {
for (;;) {
//当一个节点(关联一个线程)进入等待队列后, 获取此节点的 prev 节点
final Node p = node.predecessor();
// 如果获取到的 prev 是 head,也就是队列中第一个等待线程
if (p == head) {
// 再次尝试申请 反应到 CountDownLatch 就是查看是否还有线程需要等待(state是否为0)
int r = tryAcquireShared(arg);
// 如果 r >=0 说明 没有线程需要等待了 state==0
if (r >= 0) {
//尝试将第一个线程关联的节点设置为 head
setHeadAndPropagate(node, r);
p.next = null; // help GC
failed = false;
return;
}
}
//经过自旋tryAcquireShared后,state还不为0,就会到这里,第一次的时候,waitStatus是0,那么node的waitStatus就会被置为SIGNAL,第二次再走到这里,就会用LockSupport的park方法把当前线程阻塞住
if (shouldParkAfterFailedAcquire(p, node) &&
parkAndCheckInterrupt())
throw new InterruptedException();
}
} finally {
if (failed)
cancelAcquire(node);
}
}

我看看到上面先执行了 addWaiter() 方法,就是将当前线程加入等待队列,源码如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
/** Marker to indicate a node is waiting in shared mode */
static final Node SHARED = new Node();
/** Marker to indicate a node is waiting in exclusive mode */
static final Node EXCLUSIVE = null;

private Node addWaiter(Node mode) {
Node node = new Node(Thread.currentThread(), mode);
// 尝试快速入队操作,因为大多数时候尾节点不为 null
Node pred = tail;
if (pred != null) {
node.prev = pred;
if (compareAndSetTail(pred, node)) {
pred.next = node;
return node;
}
}
//如果尾节点为空(也就是队列为空) 或者尝试CAS入队失败(由于并发原因),进入enq方法
enq(node);
return node;
}

上面是向等待队列中添加等待者(waiter)的方法。首先构造一个 Node 实体,参数为当前线程和一个mode,这个mode有两种形式,一个是 SHARED ,一个是 EXCLUSIVE,请看上面的代码。然后执行下面的入队操作 addWaiter,和 enq() 方法的 else 分支操作是一样的,这里的操作如果成功了,就不用再进到 enq() 方法的循环中去了,可以提高性能。如果没有成功,再调用 enq() 方法。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
private Node enq(final Node node) {
// 死循环+CAS保证所有节点都入队
for (;;) {
Node t = tail;
// 如果队列为空 设置一个空节点作为 head
if (t == null) { // Must initialize
if (compareAndSetHead(new Node()))
tail = head;
} else {
//加入队尾
node.prev = t;
if (compareAndSetTail(t, node)) {
t.next = node;
return t;
}
}
}
}

说明:循环加 CAS 操作是实现乐观锁的标准方式,CAS 是为了实现原子操作而出现的,所谓的原子操作指操作执行期间,不会受其他线程的干扰。Java 实现的 CAS 是调用 unsafe 类提供的方法,底层是调用 c++ 方法,直接操作内存,在 cpu 层面加锁,直接对内存进行操作。

上面是 AQS 等待队列入队方法,操作在无限循环中进行,如果入队成功则返回新的队尾节点,否则一直自旋,直到入队成功。假设入队的节点为 node ,上来直接进入循环,在循环中,先拿到尾节点。

1、if 分支,如果尾节点为 null,说明现在队列中还没有等待线程,则尝试 CAS 操作将头节点初始化,然后将尾节点也设置为头节点,因为初始化的时候头尾是同一个,这和 AQS 的设计实现有关, AQS 默认要有一个虚拟节点。此时,尾节点不在为空,循环继续,进入 else 分支;

2、else 分支,如果尾节点不为 null, node.prev = t ,也就是将当前尾节点设置为待入队节点的前置节点。然后又是利用 CAS 操作,将待入队的节点设置为队列的尾节点,如果 CAS 返回 false,表示未设置成功,继续循环设置,直到设置成功,接着将之前的尾节点(也就是倒数第二个节点)的 next 属性设置为当前尾节点,对应 t.next = node 语句,然后返回当前尾节点,退出循环。

setHeadAndPropagate 方法负责将自旋等待或被 LockSupport 阻塞的线程唤醒。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
private void setHeadAndPropagate(Node node, int propagate) {
//备份现在的 head
Node h = head;
//抢到锁的线程被唤醒 将这个节点设置为head
setHead(node)
// propagate 一般都会大于0 或者存在可被唤醒的线程
if (propagate > 0 || h == null || h.waitStatus < 0 ||
(h = head) == null || h.waitStatus < 0) {
Node s = node.next;
// 只有一个节点 或者是共享模式 释放所有等待线程 各自尝试抢占锁
if (s == null || s.isShared())
doReleaseShared();
}
}

Node 对象中有一个属性是 waitStatus ,它有四种状态,分别是:

1
2
3
4
5
6
7
8
//线程已被 cancelled ,这种状态的节点将会被忽略,并移出队列
static final int CANCELLED = 1;
// 表示当前线程已被挂起,并且后继节点可以尝试抢占锁
static final int SIGNAL = -1;
//线程正在等待某些条件
static final int CONDITION = -2;
//共享模式下 无条件所有等待线程尝试抢占锁
static final int PROPAGATE = -3;

等待线程被唤醒

当执行 CountDownLatch 的 countDown()方法,将计数器减一,也就是state减一,当减到0的时候,等待队列中的线程被释放。是调用 AQS 的 releaseShared 方法来实现的,下面代码中的方法是按顺序调用的,摘到了一起,方便查看:

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
43
44
45
46
47
// AQS类
public final boolean releaseShared(int arg) {
// arg 为固定值 1
// 如果计数器state 为0 返回true,前提是调用 countDown() 之前不能已经为0
if (tryReleaseShared(arg)) {
// 唤醒等待队列的线程
doReleaseShared();
return true;
}
return false;
}

// CountDownLatch 重写的方法
protected boolean tryReleaseShared(int releases) {
// Decrement count; signal when transition to zero
// 依然是循环+CAS配合 实现计数器减1
for (;;) {
int c = getState();
if (c == 0)
return false;
int nextc = c-1;
if (compareAndSetState(c, nextc))
return nextc == 0;
}
}

/// AQS类
private void doReleaseShared() {
for (;;) {
Node h = head;
if (h != null && h != tail) {
int ws = h.waitStatus;
// 如果节点状态为SIGNAL,则他的next节点也可以尝试被唤醒
if (ws == Node.SIGNAL) {
if (!compareAndSetWaitStatus(h, Node.SIGNAL, 0))
continue; // loop to recheck cases
unparkSuccessor(h);
}
// 将节点状态设置为PROPAGATE,表示要向下传播,依次唤醒
else if (ws == 0 &&
!compareAndSetWaitStatus(h, 0, Node.PROPAGATE))
continue; // loop on failed CAS
}
if (h == head) // loop if head changed
break;
}
}

因为这是共享型的,当计数器为 0 后,会唤醒等待队列里的所有线程,所有调用了 await() 方法的线程都被唤醒,并发执行。这种情况对应到的场景是,有多个线程需要等待一些动作完成,比如一个线程完成初始化动作,其他5个线程都需要用到初始化的结果,那么在初始化线程调用 countDown 之前,其他5个线程都处在等待状态。一旦初始化线程调用了 countDown ,其他5个线程都被唤醒,开始执行。

Lock.condition

CyclicBarrier循环屏障

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
final CyclicBarrier cyclicBarrier = new CyclicBarrier(2, () -> doWork("C"));
final Thread ta = new Thread(() -> {
doWork("A");
try {
cyclicBarrier.await();
} catch (final InterruptedException | BrokenBarrierException e) {
e.printStackTrace();
}
});
final Thread tb = new Thread(() -> {
doWork("B");
try {
cyclicBarrier.await();
} catch (InterruptedException | BrokenBarrierException e) {
e.printStackTrace();
}
});

ta.start();
tb.start();

Semaphore信号量

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
final Semaphore semaphore = new Semaphore(2);
Thread ta = new Thread(() -> {
try {
semaphore.acquire();
} catch (InterruptedException e1) {
e1.printStackTrace();
}
doWork("A");
semaphore.release();
});
Thread tb = new Thread(() -> {
try {
semaphore.acquire();
} catch (InterruptedException e1) {
e1.printStackTrace();
}
doWork("B");
semaphore.release();
});
Thread tc = new Thread(() -> {
try {
semaphore.acquire(2);
} catch (InterruptedException e1) {
e1.printStackTrace();
}
doWork("C");
});
ta.start();
tb.start();
tc.start();

Future

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
ThreadPoolExecutor executor = new ThreadPoolExecutor(3, 3, 1, TimeUnit.MINUTES, new LinkedBlockingDeque<>());
Future<Boolean> aFuture = executor.submit(() -> {
doWork("A");
return true;
});
Future<Boolean> bFuture = executor.submit(() -> {
doWork("B");
return true;
});
executor.execute(() -> {
try {
aFuture.get();
bFuture.get();
} catch (InterruptedException | ExecutionException e1) {
e1.printStackTrace();
}
doWork("C");
});

Queue

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
LinkedBlockingDeque<Integer> queue = new LinkedBlockingDeque<>(2);
Thread ta = new Thread(() -> {
doWork("A");
queue.add(1);
});
Thread tb = new Thread(() -> {
doWork("B");
queue.add(1);
});
Thread tc = new Thread(() -> {
try {
queue.take();
queue.take();
} catch (InterruptedException e) {
e.printStackTrace();
}
doWork("C");
});
ta.start();
tb.start();
tc.start();

LockSupport

ListenableFuture

参考文献

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

单核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/

MapJoin Optimization

主要的几个优化参数:

  1. smalltable的大小,超过这个值才会使用map join
  2. map join的合并
  3. local task内存: mapjoin的小表转换为hashtable后的大小阈值,超过了会报错
  4. 本地任务可以使用的最大内存比例
  5. join条件字段类型需要一致
  6. map节点的内存大小

在Hive中,common join是很慢的,如果我们是一张大表关联多张小表,可以使用mapjoin加快速度。

mapjoin主要有以下参数:

1
2
3
4
5
6
7
8
--是否自动转换为mapjoin
hive.auto.convert.join
--小表的最大文件大小,默认为25000000,即25M
hive.mapjoin.smalltable.filesize
--是否将多个mapjoin合并为一个
hive.auto.convert.join.noconditionaltask
--多个mapjoin转换为1个时,所有小表的文件大小总和的最大值
hive.auto.convert.join.noconditionaltask.size

例如,一个大表顺序关联3个小表a(10M), b(8M),c(12M),
如果hive.auto.convert.join.noconditionaltask.size的值:

  1. <=18M,则无法合并mapjoin,必须执行3个mapjoin;
  2. 18M <30M,则可以合并a和b表的mapjoin,所以只需要执行2个mapjoin;

  3. 30M,则可以将3个mapjoin都合并为1个。

合并mapjoin

合并mapjoin有啥好处呢?
因为每个mapjoin都要执行一次map,需要读写一次数据,所以多个mapjoin就要做多次的数据读写,合并mapjoin后只用读写一次,自然能大大加快速度。
但是执行map是内存大小是有限制的,在一次map里对多个小表做mapjoin就必须把多个小表都加入内存,为了防止内存溢出,所以加了hive.auto.convert.join.noconditionaltask.size参数来做限制。不过,这个值只是限制输入的表文件的大小,并不代表实际mapjoinhashtable的大小。

我们可以通过explain查看执行计划,来看看mapjoin是否生效。

具体示例

1
2
3
4
set hive.auto.convert.join=true;
set hive.mapjoin.smalltable.filesize=300000000;
set hive.auto.convert.join.noconditionaltask=true;
set hive.auto.convert.join.noconditionaltask.size=300000000;

使用mapjoin时,会先执行一个本地任务(mapreduce local task)将小表转成hashtable并序列化为文件再压缩,随后这些hashtable文件会被上传到hadoop缓存,提供给各个mapjoin使用。这里有三个参数我们需要注意:

local task memory: 小表转换成hashtable的内存阈值

1
2
3
4
5
6
--将小表转成hashtable的本地任务的最大内存使用率,默认0.9
hive.mapjoin.localtask.max.memory.usage
--如果mapjoin后面紧跟着一个group by任务,这种情况下 本地任务的最大内存使用率,默认是0.55
hive.mapjoin.followby.gby.localtask.max.memory.usage
--localtask每处理完多少行,就执行内存检查。默认为100000
hive.mapjoin.check.memory.rows

如果我们的localtask的内存使用超过阀值,任务会直接失败。

字段类型要一致

此外,使用mapjoin时还要注意,用作join的关联字段的字段类型最好要一致。

我就碰到一个诡异的问题,执行mapjoin 的local task时一直卡住,40万行的小表处理了好几个小时,正常情况下应该几秒钟就完成了。查了好久原因,结果原来是做join的关联字段的类型不一致,一边是int, 一边是string,hive解释计划里显示它们都会被转成double再来join。我把字段类型改为一致的,瞬间就快了。照理说就算转成double也不该这么慢,不知道是不是hive的bug。

1
2
3
4
5
6
7
8
--是否自动转换为mapjoin
set hive.auto.convert.join = true;
--小表的最大文件大小,默认为25000000,即25M
set hive.mapjoin.smalltable.filesize = 25000000;
--是否将多个mapjoin合并为一个
set hive.auto.convert.join.noconditionaltask = true;
--多个mapjoin转换为1个时,所有小表的文件大小总和的最大值。
set hive.auto.convert.join.noconditionaltask.size = 10000000;

hive的join 有一种优化的方式:map join

但是,使用这种优化的时候要小心一点,先说一下优化配置的参数:

localtask max memory usage: 本地任务内存百分比

1
2
3
4
set hive.optimize.correlation=true
set hive.auto.convert.join=true
set hive.mapjoin.localtask.max.memory.usage=0.99
hive.mapjoin.localtask.max.memory.usage

说明:本地任务可以使用内存的百分比 默认值: 0.90,如果你的localtask mapjoin 表很小可以试试,但彻底解决需要

set hive.auto.convert.join=false;关闭自动mapjoin 但这个参数用的时候一定要注意,

如果你的sql 很长join会常多,关闭mapjoin任务数会成10倍激增,contener满了任务同样会非常之慢,
set hive.auto.convert.join=false;一定要用在localtask级别这种超轻量及的job上。

local mem: map节点的内存大小

1
2
--设置本地memory的大小值,单位为M
hive.mapred.local.mem

<< Java高级软件工程师知识结构

  1. 掌握InputStream、OutputStream、Reader、Writer的继承体系。
  2. 掌握字节流(FileInputStream、DataInputStream、BufferedInputStream、FileOutputSteam、DataOutputStream、BufferedOutputStream)和 字符流(BufferedReader、InputStreamReader、FileReader、BufferedWriter、OutputStreamWriter、PrintWriter、FileWriter),并熟练运用。
  3. 掌握NIO实现原理及使用方法。

Java IO包括:

【案例】ZipOutputStream类
先看一下ZipOutputStream类的继承关系
java.lang.Object
java.io.OutputStream
java.io.FilterOutputStream
java.util.zip.DeflaterOutputStream
java.util.zip.ZipOutputStream

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
import java.io.File;
import java.io.FileInputStream;
import java.io.FileOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.util.zip.ZipEntry;
import java.util.zip.ZipOutputStream;

public class ZipOutputStreamDemo1{
public static void main(String[] args) throws IOException{
File file = new File("d:" + File.separator +"hello.txt");
File zipFile = new File("d:" + File.separator +"hello.zip");
InputStream input = new FileInputStream(file);
ZipOutputStream zipOut = new ZipOutputStream(new FileOutputStream(
zipFile));
zipOut.putNextEntry(new ZipEntry(file.getName()));
// 设置注释
zipOut.setComment("hello");
int temp = 0;
while((temp = input.read()) != -1){
zipOut.write(temp);
}
input.close();
zipOut.close();
}
}

【案例】ZipOutputStream类压缩多个文件

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
import java.io.File;
import java.io.FileInputStream;
import java.io.FileOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.util.zip.ZipEntry;
import java.util.zip.ZipOutputStream;

/**
* 一次性压缩多个文件
* */
public class ZipOutputStreamDemo2{
public static void main(String[] args) throws IOException{
// 要被压缩的文件夹
File file = new File("d:" + File.separator +"temp");
File zipFile = new File("d:" + File.separator + "zipFile.zip");
InputStream input = null;
ZipOutputStream zipOut = new ZipOutputStream(new FileOutputStream(
zipFile));
zipOut.setComment("hello");
if(file.isDirectory()){
File[] files = file.listFiles();
for(int i = 0; i < files.length; ++i){
input = newFileInputStream(files[i]);
zipOut.putNextEntry(newZipEntry(file.getName()
+ File.separator +files[i].getName()));
int temp = 0;
while((temp = input.read()) !=-1){
zipOut.write(temp);
}
input.close();
}
}
zipOut.close();
}
}

【案例】ZipFile类展示

1
2
3
4
5
6
7
8
9
10
11
12
13
14
import java.io.File;
import java.io.IOException;
import java.util.zip.ZipFile;

/**
*ZipFile演示
* */
public class ZipFileDemo{
public static void main(String[] args) throws IOException{
File file = new File("d:" + File.separator +"hello.zip");
ZipFile zipFile = new ZipFile(file);
System.out.println("压缩文件的名称为:" + zipFile.getName());
}
}

【案例】解压缩文件(压缩文件中只有一个文件的情况)

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
import java.io.File;
import java.io.FileOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import java.util.zip.ZipEntry;
import java.util.zip.ZipFile;

/**
* 解压缩文件(压缩文件中只有一个文件的情况)
* */
public class ZipFileDemo2{
public static void main(String[] args) throws IOException{
File file = new File("d:" + File.separator +"hello.zip");
File outFile = new File("d:" + File.separator +"unZipFile.txt");
ZipFile zipFile = new ZipFile(file);
ZipEntry entry =zipFile.getEntry("hello.txt");
InputStream input = zipFile.getInputStream(entry);
OutputStream output = new FileOutputStream(outFile);
int temp = 0;
while((temp = input.read()) != -1){
output.write(temp);
}
input.close();
output.close();
}
}

【案例】ZipInputStream类解压缩一个压缩文件中包含多个文件的情况

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
import java.io.File;
import java.io.FileInputStream;
import java.io.FileOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import java.util.zip.ZipEntry;
import java.util.zip.ZipFile;
import java.util.zip.ZipInputStream;

/**
* 解压缩一个压缩文件中包含多个文件的情况
* */
public class ZipFileDemo3{
public static void main(String[] args) throws IOException{
File file = new File("d:" +File.separator + "zipFile.zip");
File outFile = null;
ZipFile zipFile = new ZipFile(file);
ZipInputStream zipInput = new ZipInputStream(new FileInputStream(file));
ZipEntry entry = null;
InputStream input = null;
OutputStream output = null;
while((entry = zipInput.getNextEntry()) != null){
System.out.println("解压缩" + entry.getName() + "文件");
outFile = new File("d:" + File.separator + entry.getName());
if(!outFile.getParentFile().exists()){
outFile.getParentFile().mkdir();
}
if(!outFile.exists()){
outFile.createNewFile();
}
input = zipFile.getInputStream(entry);
output = new FileOutputStream(outFile);
int temp = 0;
while((temp = input.read()) != -1){
output.write(temp);
}
input.close();
output.close();
}
}
}

[参考文献]:

  1. Think in Java