list = Collections.synchronizedList(new ArrayList<>());
for (int i = 0; i < 30; i++) {
new Thread(() -> {
//向集合添加内容
list.add(UUID.randomUUID().toString().substring(0, 8));
//从集合获取内容
System.out.println(list);
}, String.valueOf(i)).start();
}
}
}
```
输出:
> [504a67ec, f0176711, 78b9d3a6, faf52eb9, 1a8620a1]
> [504a67ec, f0176711, 78b9d3a6, faf52eb9, 1a8620a1, 2f8b2783, 1d75042f, b2db93e2, a8c2132c, e7272b1f, abd7a302, 4769a65c]
> ...
> [504a67ec, f0176711, 78b9d3a6, faf52eb9, 1a8620a1, 2f8b2783, 1d75042f, b2db93e2, a8c2132c, e7272b1f, abd7a302, 4769a65c]
没有并发修改异常,查看方法源码 :
```java
/*
* The returned list will be serializable if the specified list is
* serializable.
*
* @param the class of the objects in the list
* @param list the list to be "wrapped" in a synchronized list.
* @return a synchronized view of the specified list.
*/
public static List synchronizedList(List list) {
return (list instanceof RandomAccess ?
new SynchronizedRandomAccessList<>(list) :
new SynchronizedList<>(list));
}
```
##### 5.1.2.3)解决方案3-CopyOnWriteArrayList
CopyOnWriteArrayLis相当于线程安全的ArrayList,和ArrayList一样,它是个可变数组;但是和 ArrayList不同的时,它具有以下特性:
1. 它最适合于具有以下特征的应用程序:List大小通常保持很小,只读操作远多于可变操作,需要在遍历期间防止线程间的冲突;
2. 它是线程安全的;
3. 因为通常需要复制整个基础数组,所以可变操作(add()、set() 和 remove() 等)的开销很大;
4. 迭代器支持 hasNext(), next()等不可变操作,但不支持可变remove()等操作。
5. 使用迭代器进行遍历的速度很快,并且不会与其他线程发生冲突。在构造迭代器时,迭代器依赖于不变的数组快照。
CopyOnWriteArrayList 的核心思想就是复制思想,即拷贝一份:
当我们往一个容器添加元素的时候,不直接往当前容器添加,而是先将当前容器进行 Copy,复制出一个新的容器,然后新的容器里添加元素,添加完元素之后,再将原容器的引用指向新的容器,这时候会抛出来一个新的问题,也就是数据不一致的问题。如果写线程还没来得及写会内存,其他的线程就会读到了脏数据。
CopyOnWriteArrayList的核心技术就是 **写时复制技术**,流程如下图:

读集合数据时是并发读,写集合数据时是独立写,这样就避免了在写数据时导致的并发修改异常问题
代码如下:
```java
// CopyOnWriteArrayList解决ArrayList集合线程不安全问题
public class CopyOnWriteArrayListDemo {
public static void main(String[] args) {
// CopyOnWriteArrayList解决
List list = new CopyOnWriteArrayList<>();
for (int i = 0; i < 30; i++) {
new Thread(() -> {
//向集合添加内容
list.add(UUID.randomUUID().toString().substring(0, 8));
//从集合获取内容
System.out.println(list);
}, String.valueOf(i)).start();
}
}
}
```
输出:
> [4578f50c, 6937bb8b, ce696b2a, d6b6f9da]
> [4578f50c, 6937bb8b, ce696b2a, d6b6f9da, 970e0a0f, 569d7b7a, 39ac8475, 2aadbc14, 45df403d]
> ...
> [4578f50c, 6937bb8b, ce696b2a, d6b6f9da, 970e0a0f, 569d7b7a, 39ac8475, 2aadbc14, 45df403d, 32c0db21, 62493f05, 4ac5126f, 507391bb, 0d3633fb, def2ef6d, 1f20ecd3, 73c1a75c, 7e506688, b07c30f2, e56914cd, c7e2e730]
没有并发修改异常原因分析:CopyOnWriteArrayList 的 add() 方法源码如下
```java
/**
* Appends the specified element to the end of this list.
*
* @param e element to be appended to this list
* @return {@code true} (as specified by {@link Collection#add})
*/
public boolean add(E e) {
final ReentrantLock lock = this.lock;
// 上锁
lock.lock();
try {
Object[] elements = getArray();
int len = elements.length;
Object[] newElements = Arrays.copyOf(elements, len + 1);
newElements[len] = e;
setArray(newElements);
return true;
} finally {
// 解锁
lock.unlock();
}
}
```
没有并发修改异常原因分析:动态数组与线程安全
下面从“动态数组”和“线程安全”两个方面进一步对 CopyOnWriteArrayList的原理进行说明
**“动态数组”机制**
- 它内部有个“volatile数组”(array)来保持数据;
- 在“添加/修改/删除”数据时,都会新建一个数组,并将更新后的数据拷贝到新建的数组中,最后再将该数组赋值给“volatile数组”, 这就是它叫做CopyOnWriteArrayList的原因
- 由于它在“添加/修改/删除”数据时,都会新建数组,所以涉及到修改数据的操作,CopyOnWriteArrayList效率很低;但是单单只是进行遍历查找的话, 效率比较高。
**“线程安全”机制**
- 通过volatile和互斥锁来实现的;
- 通过“volatile数组”来保存数据;一个线程读取volatile数组时,总能看到其它线程对该volatile变量最后的写入;这样通过volatile提供了“读取到的数据总是最新的”这个机制的保证。
- 通过互斥锁来保护数据。在“添加/修改/删除”数据时,会先“获取互斥锁”,再修改完毕之后,先将数据更新到“volatile数组”中,然后再“释放互斥锁”,就达到了保护数据的目的。
### 5.2)HashSet集合线程
#### 5.2.1)HashSet集合线程不安全案例展示
代码如下:
```java
// HashSet集合线程不安全案例展示
public class HashSetErrorDemo {
public static void main(String[] args) {
//演示Hashset
Set set = new HashSet<>();
for (int i = 0; i < 30; i++) {
new Thread(() -> {
//向集合添加内容
set.add(UUID.randomUUID().toString().substring(0, 8));
//从集合获取内容
System.out.println(set);
}, String.valueOf(i)).start();
}
}
}
```
输出:
> java.util.ConcurrentModificationException
> at java.util.HashMap$HashIterator.nextNode(HashMap.java:1445)
> at java.util.HashMap$KeyIterator.next(HashMap.java:1469)
> at java.util.AbstractCollection.toString(AbstractCollection.java:461)
> at java.lang.String.valueOf(String.java:2994)
> at java.io.PrintStream.println(PrintStream.java:821)
> at com.study.collection.HashSetErrorDemo.lambda$main$0(HashSetErrorDemo.java:25)
> at java.lang.Thread.run(Thread.java:748)
> java.util.ConcurrentModificationException
> at java.util.HashMap$HashIterator.nextNode(HashMap.java:1445)
> at java.util.HashMap$KeyIterator.next(HashMap.java:1469)
> at java.util.AbstractCollection.toString(AbstractCollection.java:461)
> at java.lang.String.valueOf(String.java:2994)
> at java.io.PrintStream.println(PrintStream.java:821)
> at com.study.collection.HashSetErrorDemo.lambda$main$0(HashSetErrorDemo.java:25)
> at java.lang.Thread.run(Thread.java:748)
发现报错,java.util.ConcurrentModificationException ——》并发修改异常
> 问题: 为什么会出现并发修改异常?
查看HashSet的add方法源码 ,如下:
```java
/**
* Adds the specified element to this set if it is not already present.
* More formally, adds the specified element e to this set if
* this set contains no element e2 such that
* (e==null ? e2==null : e.equals(e2)).
* If this set already contains the element, the call leaves the set
* unchanged and returns false.
*
* @param e element to be added to this set
* @return true if this set did not already contain the specified
* element
*/
public boolean add(E e) {
return map.put(e, PRESENT)==null;
}
```
add() 方法没有添加 synchronized 关键字,即没有对多线程情况进行处理
#### 5.2.2)HashSet集合线程不安全解决方案
##### 5.2.2.1)解决方案-CopyOnWriteArraySet
使用 CopyOnWriteArraySet解决HashSet集合线程不安全问题,代码如下:
```java
// CopyOnWriteArray解决HashSet集合线程不安全问题
public class CopyOnWriteArraySetDemo {
public static void main(String[] args) {
//演示Hashset
Set set = new CopyOnWriteArraySet<>();
for (int i = 0; i < 30; i++) {
new Thread(() -> {
//向集合添加内容
set.add(UUID.randomUUID().toString().substring(0, 8));
//从集合获取内容
System.out.println(set);
}, String.valueOf(i)).start();
}
}
}
```
输出:
> com.study.collection.HashSet.CopyOnWriteArraySetDemo
> [6cfeb375, cc511c07, 79eddeb2, 6ba3f255, edf8d95d, 32cb0924, 5a3189bf, 196e5a54, 9875e04b, e240ad35, 22debfb2, 435bae49, 968d06de, d7fdacae, a46e7686, 6ea6c822, b380ce03, 54ddfa83, 16d0a020]
> [6cfeb375, cc511c07, 79eddeb2, 6ba3f255, edf8d95d, 32cb0924, 5a3189bf, 196e5a54, 9875e04b, e240ad35, 22debfb2, 435bae49, 968d06de, d7fdacae, a46e7686, 6ea6c822, b380ce03, 54ddfa83, 16d0a020, 2f4befd6]
> ...
> [6cfeb375, cc511c07, 79eddeb2, 6ba3f255, edf8d95d, 32cb0924, 5a3189bf, 196e5a54, 9875e04b, e240ad35, 22debfb2, 435bae49, 968d06de, d7fdacae, a46e7686, 6ea6c822, b380ce03, 54ddfa83, 16d0a020, 2f4befd6, 8d3201e7, b4d2fb7a]
>
> Process finished with exit code 0
CopyOnWriteArraySet的add() 方法源码如下:addIfAbsent()中实现了lock
```java
final transient ReentrantLock lock = new ReentrantLock();
public boolean add(E e) {
return al.addIfAbsent(e);
}
/**
* Appends the element, if not present.
*
* @param e element to be added to this list, if absent
* @return {@code true} if the element was added
*/
public boolean addIfAbsent(E e) {
Object[] snapshot = getArray();
return indexOf(e, snapshot, 0, snapshot.length) >= 0 ? false :
addIfAbsent(e, snapshot);
}
/**
* A version of addIfAbsent using the strong hint that given
* recent snapshot does not contain e.
*/
private boolean addIfAbsent(E e, Object[] snapshot) {
final ReentrantLock lock = this.lock;
lock.lock();
try {
Object[] current = getArray();
int len = current.length;
if (snapshot != current) {
// Optimize for lost race to another addXXX operation
int common = Math.min(snapshot.length, len);
for (int i = 0; i < common; i++)
if (current[i] != snapshot[i] && eq(e, current[i]))
return false;
if (indexOf(e, current, common, len) >= 0)
return false;
}
Object[] newElements = Arrays.copyOf(current, len + 1);
newElements[len] = e;
setArray(newElements);
return true;
} finally {
lock.unlock();
}
}
```
### 5.3)HashMap集合线程
#### 5.3.1)HashMap集合线程不安全案例展示
代码如下:
```java
// HashMap集合线程不安全案例展示
public class HashMapErrorDemo {
public static void main(String[] args) {
//演示HashMap
Map map = new HashMap<>();
for (int i = 0; i <30; i++) {
String key = String.valueOf(i);
new Thread(()->{
//向集合添加内容
map.put(key,UUID.randomUUID().toString().substring(0,8));
//从集合获取内容
System.out.println(map);
},String.valueOf(i)).start();
}
}
}
```
输出:
> Exception in thread "17" Exception in thread "1" Exception in thread "16" Exception in thread "26" java.util.ConcurrentModificationException
> at java.util.HashMap$HashIterator.nextNode(HashMap.java:1445)
> at java.util.HashMap$EntryIterator.next(HashMap.java:1479)
> at java.util.HashMap$EntryIterator.next(HashMap.java:1477)
> at java.util.AbstractMap.toString(AbstractMap.java:554)
> at java.lang.String.valueOf(String.java:2994)
> at java.io.PrintStream.println(PrintStream.java:821)
> at com.study.collection.HashMap.HashMapErrorDemo.lambda$main$0(HashMapErrorDemo.java:27)
> at java.lang.Thread.run(Thread.java:748)
> java.util.ConcurrentModificationException
> at java.util.HashMap$HashIterator.nextNode(HashMap.java:1445)
> at java.util.HashMap$EntryIterator.next(HashMap.java:1479)
> at java.util.HashMap$EntryIterator.next(HashMap.java:1477)
> at java.util.AbstractMap.toString(AbstractMap.java:554)
> at java.lang.String.valueOf(String.java:2994)
> at java.io.PrintStream.println(PrintStream.java:821)
> at com.study.collection.HashMap.HashMapErrorDemo.lambda$main$0(HashMapErrorDemo.java:27)
> at java.lang.Thread.run(Thread.java:748)
发现报错,java.util.ConcurrentModificationException ——》并发修改异常
> 问题: 为什么会出现并发修改异常?
查看HashMap的put方法源码 ,如下:
```java
public V put(K key, V value) {
return putVal(hash(key), key, value, false, true);
}
/**
* Implements Map.put and related methods.
*
* @param hash hash for key
* @param key the key
* @param value the value to put
* @param onlyIfAbsent if true, don't change existing value
* @param evict if false, the table is in creation mode.
* @return previous value, or null if none
*/
final V putVal(int hash, K key, V value, boolean onlyIfAbsent,
boolean evict) {
Node[] tab; Node p; int n, i;
if ((tab = table) == null || (n = tab.length) == 0)
n = (tab = resize()).length;
if ((p = tab[i = (n - 1) & hash]) == null)
tab[i] = newNode(hash, key, value, null);
else {
Node e; K k;
if (p.hash == hash &&
((k = p.key) == key || (key != null && key.equals(k))))
e = p;
else if (p instanceof TreeNode)
e = ((TreeNode)p).putTreeVal(this, tab, hash, key, value);
else {
for (int binCount = 0; ; ++binCount) {
if ((e = p.next) == null) {
p.next = newNode(hash, key, value, null);
if (binCount >= TREEIFY_THRESHOLD - 1) // -1 for 1st
treeifyBin(tab, hash);
break;
}
if (e.hash == hash &&
((k = e.key) == key || (key != null && key.equals(k))))
break;
p = e;
}
}
if (e != null) { // existing mapping for key
V oldValue = e.value;
if (!onlyIfAbsent || oldValue == null)
e.value = value;
afterNodeAccess(e);
return oldValue;
}
}
++modCount;
if (++size > threshold)
resize();
afterNodeInsertion(evict);
return null;
}
```
put() 方法没有添加 synchronized 关键字,即没有对多线程情况进行处理
#### 5.3.2)HashMap集合线程不安全解决方案
##### 5.3.2.1)解决方案-ConcurrentHashMap
使用 ConcurrentHashMap 解决HashMap集合线程不安全问题,代码如下:
```java
// ConcurrentHashMap解决HashMap集合线程不安全问题
public class ConcurrentHashMapDemo {
public static void main(String[] args) {
//演示ConcurrentHashMap
Map map = new ConcurrentHashMap<>();
for (int i = 0; i <30; i++) {
String key = String.valueOf(i);
new Thread(()->{
//向集合添加内容
map.put(key, UUID.randomUUID().toString().substring(0,8));
//从集合获取内容
System.out.println(map);
},String.valueOf(i)).start();
}
}
}
```
输出:
```
{0=8d7844d0, 1=19cc45b9, 2=a6ece204, 3=3878495c, 4=5f6e1594, 5=a88e727a, 6=9cd46cbd, 7=2f507f9a, 8=3e0da179, 9=61c9b227}
{11=feff97c8, 0=8d7844d0, 1=19cc45b9, 2=a6ece204, 3=3878495c, 4=5f6e1594, 5=a88e727a, 6=9cd46cbd, 7=2f507f9a, 8=3e0da179, 9=61c9b227, 10=6b7f272f}
...
{0=8d7844d0, 1=19cc45b9, 2=a6ece204, 3=3878495c, 4=5f6e1594, 5=a88e727a, 6=9cd46cbd, 7=2f507f9a, 8=3e0da179, 9=61c9b227, 10=6b7f272f}
```
ConcurrentHashMap的put() 方法源码如下:putVal()中添加了synchronized
```java
public V put(K key, V value) {
return putVal(key, value, false);
}
/** Implementation for put and putIfAbsent */
final V putVal(K key, V value, boolean onlyIfAbsent) {
if (key == null || value == null) throw new NullPointerException();
int hash = spread(key.hashCode());
int binCount = 0;
for (Node[] tab = table;;) {
Node f; int n, i, fh;
if (tab == null || (n = tab.length) == 0)
tab = initTable();
else if ((f = tabAt(tab, i = (n - 1) & hash)) == null) {
if (casTabAt(tab, i, null,
new Node(hash, key, value, null)))
break; // no lock when adding to empty bin
}
else if ((fh = f.hash) == MOVED)
tab = helpTransfer(tab, f);
else {
V oldVal = null;
synchronized (f) {
if (tabAt(tab, i) == f) {
if (fh >= 0) {
binCount = 1;
for (Node e = f;; ++binCount) {
K ek;
if (e.hash == hash &&
((ek = e.key) == key ||
(ek != null && key.equals(ek)))) {
oldVal = e.val;
if (!onlyIfAbsent)
e.val = value;
break;
}
Node pred = e;
if ((e = e.next) == null) {
pred.next = new Node(hash, key,
value, null);
break;
}
}
}
else if (f instanceof TreeBin) {
Node p;
binCount = 2;
if ((p = ((TreeBin)f).putTreeVal(hash, key,
value)) != null) {
oldVal = p.val;
if (!onlyIfAbsent)
p.val = value;
}
}
}
}
if (binCount != 0) {
if (binCount >= TREEIFY_THRESHOLD)
treeifyBin(tab, i);
if (oldVal != null)
return oldVal;
break;
}
}
}
addCount(1L, binCount);
return null;
}
```
## 六、多线程锁
### 6.1)演示锁的八种情况
#### 6.1.1)标准访问
代码如下:
```java
class Phone {
// 打印短信
public synchronized void sendSMS() throws Exception {
System.out.println("------sendSMS");
}
// 打印邮件
public synchronized void sendEmail() throws Exception {
System.out.println("------sendEmail");
}
}
public class Lock_8 {
public static void main(String[] args) throws Exception {
Phone phone = new Phone();
// 创建线程AA
new Thread(() -> {
try {
phone.sendSMS();
} catch (Exception e) {
e.printStackTrace();
}
}, "AA").start();
Thread.sleep(100);
// 创建线程BB
new Thread(() -> {
try {
phone.sendEmail();
} catch (Exception e) {
e.printStackTrace();
}
}, "BB").start();
}
}
```
输出:标准访问,先打印短后打印邮件
> ------sendSMS
> ------sendEmail
#### 6.1.2)停4秒在短信方法内
代码如下:
```java
class Phone {
// 打印短信
public synchronized void sendSMS() throws Exception {
//停留4秒
TimeUnit.SECONDS.sleep(4);
System.out.println("------sendSMS");
}
// 打印邮件
public synchronized void sendEmail() throws Exception {
System.out.println("------sendEmail");
}
}
public class Lock_8 {
public static void main(String[] args) throws Exception {
Phone phone = new Phone();
// 创建线程AA
new Thread(() -> {
try {
phone.sendSMS();
} catch (Exception e) {
e.printStackTrace();
}
}, "AA").start();
Thread.sleep(100);
// 创建线程BB
new Thread(() -> {
try {
phone.sendEmail();
} catch (Exception e) {
e.printStackTrace();
}
}, "BB").start();
}
}
```
输出:停4秒在短信方法内,先打印短后打印邮件
> ------sendSMS
> ------sendEmail
说明:上述两种情况输出结果相同,表示添加 synchronized 关键字后加锁的是当前使用的对象(即phone),即先执行phone.sendSMS()再执行phone.sendEmail()。
#### 6.1.3)新增普通方法
代码如下:
```java
class Phone {
// 打印短信
public synchronized void sendSMS() throws Exception {
//停留4秒
TimeUnit.SECONDS.sleep(4);
System.out.println("------sendSMS");
}
// 打印邮件
public synchronized void sendEmail() throws Exception {
System.out.println("------sendEmail");
}
// 普通方法
public void getHello() {
System.out.println("------getHello");
}
}
public class Lock_8 {
public static void main(String[] args) throws Exception {
Phone phone = new Phone();
// 创建线程AA
new Thread(() -> {
try {
phone.sendSMS();
} catch (Exception e) {
e.printStackTrace();
}
}, "AA").start();
Thread.sleep(100);
// 创建线程BB
new Thread(() -> {
try {
phone.getHello();
} catch (Exception e) {
e.printStackTrace();
}
}, "BB").start();
}
}
```
输出:新增普通的hello方法,先执行普通hello方法再打印短信
> ------getHello
> ------sendSMS
说明:上述输出结果表示,不添加 synchronized关键字的普通方法先执行,添加synchronized关键字的方法后执行,即先执行phone.getHello()再执行phone.sendSMS()。
#### 6.1.4)两部手机
代码如下:
```java
class Phone {
// 打印短信
public synchronized void sendSMS() throws Exception {
//停留4秒
TimeUnit.SECONDS.sleep(4);
System.out.println("------sendSMS");
}
// 打印邮件
public synchronized void sendEmail() throws Exception {
System.out.println("------sendEmail");
}
}
public class Lock_8 {
public static void main(String[] args) throws Exception {
Phone phone = new Phone();
Phone phone2 = new Phone();
// 创建线程AA
new Thread(() -> {
try {
phone.sendSMS();
} catch (Exception e) {
e.printStackTrace();
}
}, "AA").start();
Thread.sleep(100);
// 创建线程BB
new Thread(() -> {
try {
phone2.sendEmail();
} catch (Exception e) {
e.printStackTrace();
}
}, "BB").start();
}
}
```
输出:现在有两部手机,先打印邮件再打印短信
> ------sendEmail
> ------sendSMS
说明:上述输出结果表示添加 synchronized 关键字后加锁的是当前使用的对象(即phone),现在出现了2个对象所以加锁失效,因为实际上是加了2把锁,并不是同一把锁,即先执行phone2.sendEmail()再执行phone.sendSMS()。
#### 6.1.5)两个静态同步方法,1部手机
代码如下:
```java
class Phone {
// 打印短信
public static synchronized void sendSMS() throws Exception {
//停留4秒
TimeUnit.SECONDS.sleep(4);
System.out.println("------sendSMS");
}
// 打印邮件
public static synchronized void sendEmail() throws Exception {
System.out.println("------sendEmail");
}
}
public class Lock_8 {
public static void main(String[] args) throws Exception {
Phone phone = new Phone();
// 创建线程AA
new Thread(() -> {
try {
phone.sendSMS();
} catch (Exception e) {
e.printStackTrace();
}
}, "AA").start();
Thread.sleep(100);
// 创建线程BB
new Thread(() -> {
try {
phone.sendEmail();
} catch (Exception e) {
e.printStackTrace();
}
}, "BB").start();
}
}
```
输出:两个静态同步方法,1部手机,先打印短信再打印邮件
> ------sendSMS
> ------sendEmail
#### 6.1.6)两个静态同步方法,2部手机
代码如下:
```java
class Phone {
// 打印短信
public static synchronized void sendSMS() throws Exception {
//停留4秒
TimeUnit.SECONDS.sleep(4);
System.out.println("------sendSMS");
}
// 打印邮件
public static synchronized void sendEmail() throws Exception {
System.out.println("------sendEmail");
}
}
public class Lock_8 {
public static void main(String[] args) throws Exception {
Phone phone = new Phone();
Phone phone2 = new Phone();
// 创建线程AA
new Thread(() -> {
try {
phone.sendSMS();
} catch (Exception e) {
e.printStackTrace();
}
}, "AA").start();
Thread.sleep(100);
// 创建线程BB
new Thread(() -> {
try {
phone2.sendEmail();
} catch (Exception e) {
e.printStackTrace();
}
}, "BB").start();
}
}
```
输出:两个静态同步方法,2部手机,先打印短信再打印邮件
> ------sendSMS
> ------sendEmail
说明:上述两种情况输出结果相同,表示添加 static synchronized 关键字后加锁的是当前的字节码对象(即 class Phone )而不是实例对象(即 Phone ),相当于全局加锁了,即先执行phone2.sendEmail()再执行phone.sendSMS()。
#### 6.1.7)1个静态同步方法,1个普通同步方法,1部手机
代码如下:
```java
class Phone {
// 打印短信
public static synchronized void sendSMS() throws Exception {
//停留4秒
TimeUnit.SECONDS.sleep(4);
System.out.println("------sendSMS");
}
// 打印邮件
public synchronized void sendEmail() throws Exception {
System.out.println("------sendEmail");
}
}
public class Lock_8 {
public static void main(String[] args) throws Exception {
Phone phone = new Phone();
// 创建线程AA
new Thread(() -> {
try {
phone.sendSMS();
} catch (Exception e) {
e.printStackTrace();
}
}, "AA").start();
Thread.sleep(100);
// 创建线程BB
new Thread(() -> {
try {
phone.sendEmail();
} catch (Exception e) {
e.printStackTrace();
}
}, "BB").start();
}
}
```
输出:1个静态同步方法,1个普通同步方法,1部手机,先打印邮件再打印短信
> ------sendEmail
> ------sendSMS
#### 6.1.8)1个静态同步方法,1个普通同步方法,2部手机
代码如下:
```java
class Phone {
// 打印短信
public static synchronized void sendSMS() throws Exception {
//停留4秒
TimeUnit.SECONDS.sleep(4);
System.out.println("------sendSMS");
}
// 打印邮件
public synchronized void sendEmail() throws Exception {
System.out.println("------sendEmail");
}
}
public class Lock_8 {
public static void main(String[] args) throws Exception {
Phone phone = new Phone();
Phone phone2 = new Phone();
// 创建线程AA
new Thread(() -> {
try {
phone.sendSMS();
} catch (Exception e) {
e.printStackTrace();
}
}, "AA").start();
Thread.sleep(100);
// 创建线程BB
new Thread(() -> {
try {
phone2.sendEmail();
} catch (Exception e) {
e.printStackTrace();
}
}, "BB").start();
}
}
```
输出:1个静态同步方法,1个普通同步方法,2部手机,先打印邮件再打印短信
> ------sendEmail
> ------sendSMS
说明:上述两种情况输出结果相同,表示静态同步方法和普通同步方法的加锁对象不同,还是2把锁,所以锁实际是失效的,即先执行phone2.sendEmail()再执行phone.sendSMS()。
#### 6.1.9)小结
一个对象里面如果有多个synchronized方法,某一个时刻内,只要一个线程去调用其中的一个synchronized方法了, 其它的线程都只能等待,就是说,某一个时刻内只能有唯一一个线程去访问这些 synchronized方法;
锁的是当前this对象,被锁定后,其它的线程都不能进入到当前对象的其它的 synchronized方法,加入普通方法后发现和同步锁无关;
换成两个对象后,不是同一把锁了,锁实际上是无效的。
synchronized实现同步的基础:Java中的每一个对象都可以作为锁,具体表现为以下3种形式:
1. 对于普通同步方法,锁是当前实例对象;
2. 对于静态同步方法,锁是当前类的Class对象;
3. 对于同步方法块,锁是Synchonized括号里配置的对象
当一个线程试图访问同步代码块时,它首先必须得到锁,退出或抛出异常时必须释放锁;也就是说如果一个实例对象的非静态同步方法获取锁后,该实例对象的其他非静态同步 法必须等待获取锁的方法释放锁后才能获取锁, 可是别的实例对象的非静态同步方法因为跟该实例对象的非静态同步方法用的是不同的锁, 所以毋须等待该实例对象已获取锁的非静态同步方法释放锁就可以获取他们自己的锁。
所有的静态同步方法用的也是同一把锁——类对象本身,这两把锁是两个不同的对象, 以静态同步方法与非静态同步方法之间是不会有竞态条件的。
但是一旦一个静态同步方法获取锁后,其他的静态同步方法都必须等待该方法释放锁后才能获取锁,而不管是同一个实例对象的静态同步方法之间还是不同的实例对象的静态同 步方法之间,只要它们同一个类的实例对象!
### 6.2)公平锁和非公平锁
#### 6.2.1)代码示例
##### 6.2.1.1)公平锁
以卖票为例进行举例,代码如下:
```java
//第一步 创建资源类,定义属性和和操作方法
class LTicket {
//票数量
private int number = 30;
// 公平锁
private final ReentrantLock lock = new ReentrantLock(true);
//卖票方法
public void sale() {
//上锁
lock.lock();
try {
//判断是否有票
if (number > 0) {
System.out.println(Thread.currentThread().getName() + " :卖出" + (number--) + " 剩余:" + number);
}
} finally {
//解锁
lock.unlock();
}
}
}
public class LSaleTicket {
//第二步 创建多个线程,调用资源类的操作方法
//创建三个线程
public static void main(String[] args) {
LTicket ticket = new LTicket();
// 创建线程 AA
new Thread(() -> {
for (int i = 0; i < 40; i++) {
ticket.sale();
}
}, "AA").start();
// 创建线程 BB
new Thread(() -> {
for (int i = 0; i < 40; i++) {
ticket.sale();
}
}, "BB").start();
// 创建线程 CC
new Thread(() -> {
for (int i = 0; i < 40; i++) {
ticket.sale();
}
}, "CC").start();
}
}
```
输出:
> AA :卖出30 剩余:29
> AA :卖出29 剩余:28
> AA :卖出28 剩余:27
> AA :卖出27 剩余:26
> AA :卖出26 剩余:25
> AA :卖出25 剩余:24
> AA :卖出24 剩余:23
> AA :卖出23 剩余:22
> AA :卖出22 剩余:21
> BB :卖出21 剩余:20
> CC :卖出20 剩余:19
> AA :卖出19 剩余:18
> BB :卖出18 剩余:17
> CC :卖出17 剩余:16
> AA :卖出16 剩余:15
> BB :卖出15 剩余:14
> CC :卖出14 剩余:13
> AA :卖出13 剩余:12
> BB :卖出12 剩余:11
> CC :卖出11 剩余:10
> AA :卖出10 剩余:9
> BB :卖出9 剩余:8
> CC :卖出8 剩余:7
> AA :卖出7 剩余:6
> BB :卖出6 剩余:5
> CC :卖出5 剩余:4
> AA :卖出4 剩余:3
> BB :卖出3 剩余:2
> CC :卖出2 剩余:1
> AA :卖出1 剩余:0
##### 6.2.1.2)非公平锁
以卖票为例进行举例,代码如下:
```java
//第一步 创建资源类,定义属性和和操作方法
class LTicket {
//票数量
private int number = 30;
// 非公平锁实现一
//private final ReentrantLock lock = new ReentrantLock(false);
// 非公平锁实现二
private final ReentrantLock lock = new ReentrantLock();
//卖票方法
public void sale() {
//上锁
lock.lock();
try {
//判断是否有票
if (number > 0) {
System.out.println(Thread.currentThread().getName() + " :卖出" + (number--) + " 剩余:" + number);
}
} finally {
//解锁
lock.unlock();
}
}
}
public class LSaleTicket {
//第二步 创建多个线程,调用资源类的操作方法
//创建三个线程
public static void main(String[] args) {
LTicket ticket = new LTicket();
// 创建线程 AA
new Thread(() -> {
for (int i = 0; i < 40; i++) {
ticket.sale();
}
}, "AA").start();
// 创建线程 BB
new Thread(() -> {
for (int i = 0; i < 40; i++) {
ticket.sale();
}
}, "BB").start();
// 创建线程 CC
new Thread(() -> {
for (int i = 0; i < 40; i++) {
ticket.sale();
}
}, "CC").start();
}
}
```
输出:
> AA :卖出30 剩余:29
> AA :卖出29 剩余:28
> AA :卖出28 剩余:27
> AA :卖出27 剩余:26
> AA :卖出26 剩余:25
> AA :卖出25 剩余:24
> AA :卖出24 剩余:23
> AA :卖出23 剩余:22
> AA :卖出22 剩余:21
> AA :卖出21 剩余:20
> AA :卖出20 剩余:19
> AA :卖出19 剩余:18
> AA :卖出18 剩余:17
> AA :卖出17 剩余:16
> AA :卖出16 剩余:15
> AA :卖出15 剩余:14
> AA :卖出14 剩余:13
> AA :卖出13 剩余:12
> AA :卖出12 剩余:11
> AA :卖出11 剩余:10
> AA :卖出10 剩余:9
> AA :卖出9 剩余:8
> AA :卖出8 剩余:7
> AA :卖出7 剩余:6
> AA :卖出6 剩余:5
> AA :卖出5 剩余:4
> AA :卖出4 剩余:3
> AA :卖出3 剩余:2
> AA :卖出2 剩余:1
> AA :卖出1 剩余:0
#### 6.2.2)说明
通过上述代码示例,可以发现下列特点:
##### 6.2.2.1)公平锁
对于线程来说,是阳光普照性质的,所以线程都有执行到的机会,不会发生线程饿死的情况,但执行效率相对较低
```java
// 公平锁
private final ReentrantLock lock = new ReentrantLock(true);
```
源码如下:
```java
/**
* Creates an instance of {@code ReentrantLock} with the
* given fairness policy.
*
* @param fair {@code true} if this lock should use a fair ordering policy
*/
public ReentrantLock(boolean fair) {
sync = fair ? new FairSync() : new NonfairSync();
}
```
##### 6.2.2.2)非公平锁
对于线程来说,是竞争性质的,会发生线程饿死的情况,但执行效率相对较高
```java
// 非公平锁实现一
//private final ReentrantLock lock = new ReentrantLock(false);
// 非公平锁实现二
private final ReentrantLock lock = new ReentrantLock();
```
源码如下:
```java
/**
* Creates an instance of {@code ReentrantLock}.
* This is equivalent to using {@code ReentrantLock(false)}.
*/
public ReentrantLock() {
sync = new NonfairSync();
}
/**
* Creates an instance of {@code ReentrantLock} with the
* given fairness policy.
*
* @param fair {@code true} if this lock should use a fair ordering policy
*/
public ReentrantLock(boolean fair) {
sync = fair ? new FairSync() : new NonfairSync();
}
```
### 6.3)可重入锁
**核心问题场景:**当一个线程获得当前实例的锁lock,并且进入了方法A,该线程在方法A没有释放该锁的时候,是否可以再次进入使用该锁的方法B?
不可重入锁:在方法A释放锁之前,不可以再次进入方法B
可重入锁:在方法A释放该锁之前,可以再次进入方法B
**不可重入锁:**
当线程在访问A方法的时候,获取的A方法的锁,在A方法锁释放之前不能够访问其他方法(如方法B)的锁。
不可重入锁模型:{}{}{}{}{}都是独立的访问每一个方法,加锁 - 释放;加锁 - 释放。。。
**可重入锁:**
当线程在访问A方法的时候,获取A方法的锁,然后访问B方法获取B方法的锁,并计数加1,以此类推可以访问完了以后依次解锁。
可重入锁模型:{{{{}}}} 每次都可访问另一个方法,且加锁计数器加1,完全释放锁为计数器等于0
可重入锁,是指同一个线程可以重入上锁的代码段,不同的线程进入则需要进行阻塞等待。
Java的可重入锁有:reentrantLock(显式的可重入锁)【手动上锁解锁】、synchronized(隐式的可重入锁)【自动上锁解锁】
可重入锁诞生的目的:**防止死锁**,导致同一个线程不可重入上锁代码段,目的就是让同一个线程可以重新进入上锁代码段
#### 6.3.1)synchronized实现可重入锁
synchronized实现可重入锁代码示例:
```java
//synchronized实现可重入锁代码示例
public class SyncLockDemo {
public static void main(String[] args) {
// synchronized实现可重入锁
Object o = new Object();
new Thread(() -> {
synchronized (o) {
System.out.println(Thread.currentThread().getName() + " 外层");
synchronized (o) {
System.out.println(Thread.currentThread().getName() + " 中层");
synchronized (o) {
System.out.println(Thread.currentThread().getName() + " 内层");
}
}
}
}, "t1").start();
}
}
```
输出:在外层的锁释放前,仍然可以进入上锁的中层和内层代码块
> t1 外层
> t1 中层
> t1 内层
可重入锁又称递归锁,递归锁代码示例:
```java
public class SyncLockDemo1 {
public synchronized void add() {
add();
}
public static void main(String[] args) {
new SyncLockDemo1().add();
}
}
```
输出:
> Exception in thread "main" java.lang.StackOverflowError
> at com.study.sync.SyncLockDemo1.add(SyncLockDemo1.java:12)
#### 6.3.2)Lock 实现可重入锁
Lock实现可重入锁代码示例:
```java
// Lock实现可重入锁代码示例
public class SyncLockDemo2 {
public static void main(String[] args) {
//Lock演示可重入锁
Lock lock = new ReentrantLock();
//创建线程
new Thread(() -> {
try {
//上锁
lock.lock();
System.out.println(Thread.currentThread().getName() + " 外层");
try {
//上锁
lock.lock();
System.out.println(Thread.currentThread().getName() + " 内层");
} finally {
//释放锁
lock.unlock();
}
} finally {
//释放做
lock.unlock();
}
}, "t1").start();
//创建新线程
new Thread(() -> {
lock.lock();
System.out.println("aaaa");
lock.unlock();
}, "aa").start();
}
}
```
输出:lock 需要手动上锁和解锁,如果不解锁,将会一直卡在该线程中,无法继续执行后面的线程创建操作
> t1 外层
> t1 内层
> aaaa
### 6.4)死锁
#### 6.4.1)什么是死锁
死锁:两个或者两个以上进程在执行过程中,因为争夺资源而造成一种互相等待的现象,如果没有外力干涉,它们就无法在执行下去,会卡在这里,如下图所示:

线程A持有锁A试图获取锁B,线程B持有锁B试图获取锁A,代码实现如下:
```java
/**
* 演示死锁
*/
public class DeadLock {
//创建两个对象
static Object a = new Object();
static Object b = new Object();
public static void main(String[] args) {
new Thread(()->{
synchronized (a) {
System.out.println(Thread.currentThread().getName()+" 持有锁a,试图获取锁b");
try {
TimeUnit.SECONDS.sleep(1);
} catch (InterruptedException e) {
e.printStackTrace();
}
synchronized (b) {
System.out.println(Thread.currentThread().getName()+" 获取锁b");
}
}
},"A").start();
new Thread(()->{
synchronized (b) {
System.out.println(Thread.currentThread().getName()+" 持有锁b,试图获取锁a");
try {
TimeUnit.SECONDS.sleep(1);
} catch (InterruptedException e) {
e.printStackTrace();
}
synchronized (a) {
System.out.println(Thread.currentThread().getName()+" 获取锁a");
}
}
},"B").start();
}
}
```
输出:
> A 持有锁a,试图获取锁b
> B 持有锁b,试图获取锁a
#### 6.4.2)验证是否死锁
##### 6.4.2.1)查看运行的Java进程
命令: jps -l
```
E:\StudyCode\JUC\JucStudy>jps -l
1556 sun.tools.jps.Jps
9992
11884 com.study.sync.DeadLock
13964 org.jetbrains.jps.cmdline.Launcher
8188 org.jetbrains.jps.cmdline.Launcher
```
##### 6.4.2.2)查看堆栈信息
查看DeadLock的堆栈信息,命令:jstack 进程号
```
E:\StudyCode\JUC\JucStudy>jstack 11884
2022-05-11 15:35:29
Full thread dump Java HotSpot(TM) 64-Bit Server VM (25.211-b12 mixed mode):
"DestroyJavaVM" #14 prio=5 os_prio=0 tid=0x0000000003063800 nid=0x984 waiting on condition [0x0000000000000000]
java.lang.Thread.State: RUNNABLE
"B" #13 prio=5 os_prio=0 tid=0x000000001ad5d800 nid=0x3b6c waiting for monitor entry [0x000000001b72f000]
java.lang.Thread.State: BLOCKED (on object monitor)
at com.study.sync.DeadLock.lambda$main$1(DeadLock.java:38)
- waiting to lock <0x00000000d63fde08> (a java.lang.Object)
- locked <0x00000000d63fde18> (a java.lang.Object)
at com.study.sync.DeadLock$$Lambda$2/1078694789.run(Unknown Source)
at java.lang.Thread.run(Thread.java:748)
"A" #12 prio=5 os_prio=0 tid=0x000000001ad5b800 nid=0x498c waiting for monitor entry [0x000000001b62f000]
java.lang.Thread.State: BLOCKED (on object monitor)
at com.study.sync.DeadLock.lambda$main$0(DeadLock.java:24)
- waiting to lock <0x00000000d63fde18> (a java.lang.Object)
- locked <0x00000000d63fde08> (a java.lang.Object)
at com.study.sync.DeadLock$$Lambda$1/1324119927.run(Unknown Source)
at java.lang.Thread.run(Thread.java:748)
"Service Thread" #11 daemon prio=9 os_prio=0 tid=0x0000000019f9a000 nid=0x3b00 runnable [0x0000000000000000]
java.lang.Thread.State: RUNNABLE
"C1 CompilerThread3" #10 daemon prio=9 os_prio=2 tid=0x0000000019f48800 nid=0x4988 waiting on condition [0x0000000000000000]
java.lang.Thread.State: RUNNABLE
"C2 CompilerThread2" #9 daemon prio=9 os_prio=2 tid=0x0000000019f10800 nid=0x6c0 waiting on condition [0x0000000000000000]
java.lang.Thread.State: RUNNABLE
"C2 CompilerThread1" #8 daemon prio=9 os_prio=2 tid=0x0000000019f0b000 nid=0x450c waiting on condition [0x0000000000000000]
java.lang.Thread.State: RUNNABLE
"C2 CompilerThread0" #7 daemon prio=9 os_prio=2 tid=0x0000000019f02800 nid=0x47a0 waiting on condition [0x0000000000000000]
java.lang.Thread.State: RUNNABLE
"Monitor Ctrl-Break" #6 daemon prio=5 os_prio=0 tid=0x0000000019f09000 nid=0x840 runnable [0x000000001a72f000]
java.lang.Thread.State: RUNNABLE
at java.net.SocketInputStream.socketRead0(Native Method)
at java.net.SocketInputStream.socketRead(SocketInputStream.java:116)
at java.net.SocketInputStream.read(SocketInputStream.java:171)
at java.net.SocketInputStream.read(SocketInputStream.java:141)
at sun.nio.cs.StreamDecoder.readBytes(StreamDecoder.java:284)
at sun.nio.cs.StreamDecoder.implRead(StreamDecoder.java:326)
at sun.nio.cs.StreamDecoder.read(StreamDecoder.java:178)
- locked <0x00000000d63883f8> (a java.io.InputStreamReader)
at java.io.InputStreamReader.read(InputStreamReader.java:184)
at java.io.BufferedReader.fill(BufferedReader.java:161)
at java.io.BufferedReader.readLine(BufferedReader.java:324)
- locked <0x00000000d63883f8> (a java.io.InputStreamReader)
at java.io.BufferedReader.readLine(BufferedReader.java:389)
at com.intellij.rt.execution.application.AppMainV2$1.run(AppMainV2.java:47)
"Attach Listener" #5 daemon prio=5 os_prio=2 tid=0x0000000019e6c800 nid=0x3a48 waiting on condition [0x0000000000000000]
java.lang.Thread.State: RUNNABLE
"Signal Dispatcher" #4 daemon prio=9 os_prio=2 tid=0x0000000019ec0000 nid=0x170 runnable [0x0000000000000000]
java.lang.Thread.State: RUNNABLE
"Finalizer" #3 daemon prio=8 os_prio=1 tid=0x0000000019e51800 nid=0x1ecc in Object.wait() [0x000000001a42f000]
java.lang.Thread.State: WAITING (on object monitor)
at java.lang.Object.wait(Native Method)
- waiting on <0x00000000d6208ed0> (a java.lang.ref.ReferenceQueue$Lock)
at java.lang.ref.ReferenceQueue.remove(ReferenceQueue.java:144)
- locked <0x00000000d6208ed0> (a java.lang.ref.ReferenceQueue$Lock)
at java.lang.ref.ReferenceQueue.remove(ReferenceQueue.java:165)
at java.lang.ref.Finalizer$FinalizerThread.run(Finalizer.java:216)
"Reference Handler" #2 daemon prio=10 os_prio=2 tid=0x0000000019e50800 nid=0x40ac in Object.wait() [0x000000001a32e000]
java.lang.Thread.State: WAITING (on object monitor)
at java.lang.Object.wait(Native Method)
- waiting on <0x00000000d6206bf8> (a java.lang.ref.Reference$Lock)
at java.lang.Object.wait(Object.java:502)
at java.lang.ref.Reference.tryHandlePending(Reference.java:191)
- locked <0x00000000d6206bf8> (a java.lang.ref.Reference$Lock)
at java.lang.ref.Reference$ReferenceHandler.run(Reference.java:153)
"VM Thread" os_prio=2 tid=0x0000000018059000 nid=0x4abc runnable
"GC task thread#0 (ParallelGC)" os_prio=0 tid=0x0000000003079000 nid=0x1110 runnable
"GC task thread#1 (ParallelGC)" os_prio=0 tid=0x000000000307a800 nid=0x320c runnable
"GC task thread#2 (ParallelGC)" os_prio=0 tid=0x000000000307c000 nid=0xdd8 runnable
"GC task thread#3 (ParallelGC)" os_prio=0 tid=0x000000000307e800 nid=0x2920 runnable
"GC task thread#4 (ParallelGC)" os_prio=0 tid=0x0000000003081000 nid=0x3b08 runnable
"GC task thread#5 (ParallelGC)" os_prio=0 tid=0x0000000003082000 nid=0x3e18 runnable
"GC task thread#6 (ParallelGC)" os_prio=0 tid=0x0000000003085000 nid=0x3ab0 runnable
"GC task thread#7 (ParallelGC)" os_prio=0 tid=0x0000000003086800 nid=0x4578 runnable
"VM Periodic Task Thread" os_prio=2 tid=0x0000000019f9b800 nid=0x4644 waiting on condition
JNI global references: 317
Found one Java-level deadlock:
=============================
"B":
waiting to lock monitor 0x000000000315bca8 (object 0x00000000d63fde08, a java.lang.Object),
which is held by "A"
"A":
waiting to lock monitor 0x000000000315e5e8 (object 0x00000000d63fde18, a java.lang.Object),
which is held by "B"
Java stack information for the threads listed above:
===================================================
"B":
at com.study.sync.DeadLock.lambda$main$1(DeadLock.java:38)
- waiting to lock <0x00000000d63fde08> (a java.lang.Object)
- locked <0x00000000d63fde18> (a java.lang.Object)
at com.study.sync.DeadLock$$Lambda$2/1078694789.run(Unknown Source)
at java.lang.Thread.run(Thread.java:748)
"A":
at com.study.sync.DeadLock.lambda$main$0(DeadLock.java:24)
- waiting to lock <0x00000000d63fde18> (a java.lang.Object)
- locked <0x00000000d63fde08> (a java.lang.Object)
at com.study.sync.DeadLock$$Lambda$1/1324119927.run(Unknown Source)
at java.lang.Thread.run(Thread.java:748)
Found 1 deadlock.
```
发现在堆栈信息中 Found 1 deadlock ,说明确实发生了死锁
#### 6.4.3)产生死锁的原因
产生死锁的原因有以下三点:
1. 系统资源不足;
2. 进程运行推进顺序不合适
3. 资源分配不当
## 七、Callable&Future接口
### 7.1)创建线程的方式
创建线程的方式有以下4种:
1. 继承Thread类;
2. 实现Runnable接口;
3. 实现Callable接口;
4. 调用线程池
Runnable接口缺少的一项功能是,当线程终止时(即run() 完成时),无法使线程返回结果,为了支持此功能, Java中提供了Callable接口。
### 7.2)Callable接口
Callable接口和Runnable接口相比较有以下特点:
1. 为了实现Runnable需要实现不返回任何内容的run()方法,而对于 Callable,需要实现在完成时返回结果的call() 方法;
2. call() 方法可以抛出异常,而run() 则不能抛出异常;
3. 为实现Callable而必须重写call方法,不能直接替换Runnable,因为Thread类的构造方法根本没有Callable
【解决该问题需要引入FutureTask,它和Runnable、Callable都有关联,Runnable接口有实现类FutureTask,FutureTask构造函数可以传递Callable】
Callable接口和Runnable接口代码如下:
```java
//比较两个接口
//实现Runnable接口
class MyThread1 implements Runnable {
@Override
public void run() {
// ...
}
}
//实现Callable接口
class MyThread2 implements Callable {
@Override
public Integer call() throws Exception {
System.out.println(Thread.currentThread().getName()+" come in callable");
return 200;
}
}
```
### 7.3)Future 接口
当 call() 方法完成时,结果必须存储在主线程已知的对象中,以便主线程可以知道该线程返回的结果。为此可以使用Future对象。
将Future视为保存结果的对象,它可能暂时不保存结果,但将来会保存(一旦 Callable返回)。
Future基本上是主线程可以跟踪进度以及其他线程的结果的一种方式;要实现此接口必须重写相关方法,这里列出了比较重要的方法,如下:
- public boolean cancel(boolean mayInterrupt):用于停止任务,如果尚未启动它将停止任务;如果已启动,则仅在mayInterrupt为true 时才会中断任务;
- public Object get() :抛出InterruptedException,ExecutionException, 用于获取任务的结果; 如果任务完成,它将立即返回结果,否则将等待任务完成,然后返回结果;
- public boolean isDone():如果任务完成,则返回true,否则返回false ;
可以看到Callable和Future做两件事:Callable与Runnable类似,因为它封装了要在另一个线程上运行的任务;而Future用于存储从另一个线程获得的结果【future也可以与Runnable一起使用】。
总结:要创建线程,需要Runnable/Callable;为了获得结果,需要future。
### 7.4)FutureTask
Java库具有具体的FutureTask类型,该类型实现Runnable和Future,并方便地将两种功能组合在一起。
可以通过为其构造函数提供Callable来创建 FutureTask;然后将FutureTask对象提供给Thread的构造函数以创建 Thread对象。从而间接地使用Callable创建线程。
#### 7.4.1)FutureTask核心原理
在主线程中需要执行比较耗时的操作时,但又不想阻塞主线程时,可以把这些任务交给Future对象在后台完成;
- 当主线程将来需要时,就可以通过Future对象获得后台作业的计算结果或者执行状态;
- 一般FutureTask多用于耗时的计算,主线程可以在完成自己的任务后,再去获取结果;
- 仅在计算完成时才能检索结果;如果计算尚未完成,则阻塞 get 方法 一旦计算完成,就不能再重新开始或取消计算 ;
- get方法而获取结果只有在计算完成时获取,否则会一直阻塞直到任务转入完成状态,然后会返回结果或者抛出异常;
- get只计算一次,因此get方法放到最后执行。
#### 7.4.2)FutureTask举例说明
例1:
> 老师上课口渴了,去买水不合适,讲课线程继续运行;
>
> 单开启线程让班长帮老师买水,把水买回来;
>
> 老师喝水时需要时候直接get水即可
例2:
> 3个同学计算, 1同学: 1+2...5 , 2同学:10+11+12....50, 3同学 100+200;
>
> 因为第2个同学计算量比较大,所以FutureTask单开启线程给2同学计算;
>
> 主线程先汇总1同学和3同学的计算结果 ,最后等2同学计算完成后,再统一汇总最终结果
核心思想:主线程中需要执行比较耗时的操作时,但又不想阻塞主线程,新增一个线程去处理耗时的操作,最终一次将所有操作汇总
### 7.5)Callable和Future代码示例
Callable 和 Future的代码示例如下:
```java
public class CallableDemo {
public static void main(String[] args) throws ExecutionException, InterruptedException {
//Runnable接口创建线程
new Thread(new MyThread1(), "AA").start();
//Callable接口,报错
// new Thread(new MyThread2(),"BB").start();
//FutureTask实现Callable接口
FutureTask futureTask1 = new FutureTask<>(new MyThread2());
//lam表达式:FutureTask实现Callable接口
FutureTask futureTask2 = new FutureTask<>(() -> {
System.out.println(Thread.currentThread().getName() + " come in callable");
return 1024;
});
//创建一个线程
new Thread(futureTask2, "lucy").start();
new Thread(futureTask1, "mary").start();
while(!futureTask2.isDone()) {
System.out.println("futureTask2.isNotDone(),wait.....");
}
//调用FutureTask的get方法
System.out.println(futureTask2.get());
System.out.println(futureTask1.get());
System.out.println(Thread.currentThread().getName() + " come over");
}
}
```
输出:
> futureTask2.isNotDone(),wait.....
> futureTask2.isNotDone(),wait.....
> ...
> futureTask2.isNotDone(),wait.....
> mary come in callable
> lucy come in callable
> futureTask2.isNotDone(),wait.....
> 1024
> 200
> main come over
### 7.6)小结
在主线程中需要执行比较耗时的操作,但又不想阻塞主线程时,可以把这些作业交给Future对象在后台完成:
- 当主线程将来需要时就可以通过Future 对象获得后台作业的计算结果或者执行状态 ;
- 一般FutureTask多用于耗时的计算,主线程可以在完成自己的任务后,再去获取结果;
- 仅在计算完成时才能检索结果,如果计算尚未完成则阻塞 get 方法,一旦计算完成,就不能再重新开始或取消计算;
- get方法而获取结果只有在计算完成时获取,否则会一直阻塞直到任务转入完成状态,然后会返回结果或者抛出异常;
- 只最终汇总计算一次。
## 八、JUC 三大辅助类
JUC中提供了三种常用的辅助类,通过这些辅助类可以很好的解决线程数量过多时Lock锁的频繁操作。这三种辅助类为:
1. CountDownLatch:减少计数
2. CyclicBarrier:循环栅栏
3. Semaphore:信号灯
### 8.1)减少计数 CountDownLatch
CountDownLatch类可以设置一个计数器,然后通过countDown方法来进行减1的操作,使用await方法等待计数器不大于0,然后继续执行await方法之后的语句。
- CountDownLatch主要有两个方法,当一个或多个线程调用await方法时,这些线程会阻塞
- 其它线程调用countDown方法会将计数器减1(调用countDown方法的线程 不会阻塞)
- 当计数器的值变为0时,因await方法阻塞的线程会被唤醒,继续执行
#### 8.1.1)使用场景
6个同学陆续离开教室后值班同学才可以关门
未添加CountDownLatch代码如下:
```java
//演示 CountDownLatch
public class CountDownLatchDemo {
//6个同学陆续离开教室之后,班长锁门
public static void main(String[] args) throws InterruptedException {
//6个同学陆续离开教室之后
for (int i = 1; i <= 6; i++) {
new Thread(() -> {
System.out.println(Thread.currentThread().getName() + " 号同学离开了教室");
}, String.valueOf(i)).start();
}
System.out.println(Thread.currentThread().getName() + " 班长锁门走人了");
}
}
```
输出:
> 1 号同学离开了教室
> 6 号同学离开了教室
> 5 号同学离开了教室
> 4 号同学离开了教室
> main 班长锁门走人了
> 3 号同学离开了教室
> 2 号同学离开了教室
发现上述结果不正确,还有人未离开但已经锁门了;
添加CountDownLatch代码如下:
```java
//演示 CountDownLatch
public class CountDownLatchDemo {
//6个同学陆续离开教室之后,班长锁门
public static void main(String[] args) throws InterruptedException {
//步骤一:创建CountDownLatch对象,设置初始值
CountDownLatch countDownLatch = new CountDownLatch(6);
//6个同学陆续离开教室之后
for (int i = 1; i <= 6; i++) {
new Thread(() -> {
System.out.println(Thread.currentThread().getName() + " 号同学离开了教室");
//步骤二:计数,计数器数值执行-1操作
countDownLatch.countDown();
}, String.valueOf(i)).start();
}
//步骤三:等待/阻塞操作,直至计数器变为0才可以执行后续操作
countDownLatch.await();
System.out.println(Thread.currentThread().getName() + " 班长锁门走人了");
}
}
```
输出:
> 2 号同学离开了教室
> 6 号同学离开了教室
> 5 号同学离开了教室
> 4 号同学离开了教室
> 1 号同学离开了教室
> 3 号同学离开了教室
> main 班长锁门走人了
发现上述结果正确,所有人都离开才会锁门;
### 8.2)循环栅栏 CyclicBarrier
CyclicBarrier看英文单词可以看出大概就是循环阻塞的意思,在使用中 CyclicBarrier的构造方法第一个参数是目标障碍数,每次执行CyclicBarrier一 次障碍数会加一,如果达到了目标障碍数,才会执行cyclicBarrier.await()之后 的语句;可以将CyclicBarrier理解为加1操作 。
#### 8.2.1)使用场景
集齐7颗龙珠就可以召唤神龙 ,代码如下:
```java
//集齐7颗龙珠就可以召唤神龙
public class CyclicBarrierDemo {
//创建固定值
private static final int NUMBER = 7;
public static void main(String[] args) {
//创建CyclicBarrier
CyclicBarrier cyclicBarrier =
new CyclicBarrier(NUMBER, () -> {
System.out.println("*****集齐7颗龙珠就可以召唤神龙");
});
//集齐七颗龙珠过程
for (int i = 1; i <= 7; i++) {
new Thread(() -> {
try {
System.out.println(Thread.currentThread().getName() + " 星龙被收集到了");
//阻塞状态,直至线程名Thread.currentThread().getName()与固定值 NUMBER相同才会执行 cyclicBarrier
cyclicBarrier.await();
} catch (Exception e) {
e.printStackTrace();
}
}, String.valueOf(i)).start();
}
}
}
```
输出:
> 1 星龙被收集到了
> 6 星龙被收集到了
> 2 星龙被收集到了
> 5 星龙被收集到了
> 4 星龙被收集到了
> 3 星龙被收集到了
> 7 星龙被收集到了
> *****集齐7颗龙珠就可以召唤神龙
说明:如果线程名Thread.currentThread().getName()与固定值 NUMBER一直都不相同,比如将上述代码中的循环次数修改为6
```
for (int i = 1; i <= 6; i++)
```
这样的话,线程名Thread.currentThread().getName()与固定值 NUMBER一直都不相同,线程将一直处于阻塞状态。
### 8.3)信号灯 Semaphore
Semaphore的构造方法中传入的第一个参数是最大信号量(可以看成最大线程池),每个信号量初始化为一个最多只能分发一个许可,使用acquire方 法获得许可,release方法释放许可 。
#### 8.3.1)使用场景
场景:抢车位, 6部汽车去争抢3个停车位 ,代码如下:
```java
//6辆汽车,停3个车位
public class SemaphoreDemo {
public static void main(String[] args) {
// 步骤一:创建Semaphore,设置许可数量
Semaphore semaphore = new Semaphore(3);
//模拟6辆汽车
for (int i = 1; i <= 6; i++) {
new Thread(() -> {
try {
//步骤二:获取许可,6个线程抢占3个车位,未获取许可的线程处于阻塞状态
semaphore.acquire();
System.out.println(Thread.currentThread().getName() + " 抢到了车位");
//设置随机停车时间
TimeUnit.SECONDS.sleep(new Random().nextInt(5));
System.out.println(Thread.currentThread().getName() + " ------离开了车位");
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
//步骤三:释放许可,之前未获取许可的线程被唤醒开始执行,去获取许可
semaphore.release();
}
}, String.valueOf(i)).start();
}
}
}
```
输出:
> 2 抢到了车位
> 4 抢到了车位
> 1 抢到了车位
> 1 ------离开了车位
> 3 抢到了车位
> 2 ------离开了车位
> 4 ------离开了车位
> 5 抢到了车位
> 6 抢到了车位
> 3 ------离开了车位
> 6 ------离开了车位
> 5 ------离开了车位
## 九、读写锁
### 9.1)读写锁介绍
现实中有这样一种场景:对共享资源有读和写的操作,且写操作没有读操作那么频繁;在没有写操作的时候,多个线程同时读一个资源没有任何问题,所以应该允许多个线程同时读取共享资源;但是如果一个线程想去写这些共享资源, 就不应该允许其他线程对该资源进行读和写的操作了。
针对这种场景,JAVA的并发包提供了读写锁ReentrantReadWriteLock, 它表示两个锁,一个是读操作相关的锁,称为共享锁;一个是写相关的锁,称为排他锁。
#### 9.1.1)读锁
读锁:是**共享锁**,会发生死锁的情况,如下图所示:

上图中:二者互相等待产生死锁
- 线程1修改数据的时候,需要等待线程2读取数据之后;
- 线程2修改数据的时候,同样需要等待线程1读取数据之后;
线程进入读锁的前提条件:
1. 没有其他线程的写锁
2. 没有写请求, 或者有写请求,但调用线程和持有锁的线程是同一个(可重入锁)
#### 9.1.2)写锁
写锁:是**独占锁**,会发生死锁的情况,如下图所示:

上图中:二者互相等待产生死锁
- 线程1对第一条数据进行写操作时候,同时对第二条数据进行修改操作【同一个线程中操作多条数据】需要等待线程2对第二条数据进行写操作之后才能够进行;
- 线程2对第二条数据进行写操作时候,同时对第一条数据进行修改操作【同一个线程中操作多条数据】需要等待线程1对第一条数据进行写操作之后才能够进行;;
线程进入写锁的前提条件:
1. 没有其他线程的读锁
2. 没有其他线程的写锁
#### 9.1.3)读写锁特性
读写锁有以下三个重要的特性:
(1)公平选择性:支持非公平(默认)和公平的锁获取方式,吞吐量还是非公平优于公平;
(2)重进入:读锁和写锁都支持线程重进入;
(3)锁降级:遵循获取写锁、获取读锁再释放写锁的次序,写锁能够降级成为读锁。
#### 9.1.4)读写锁的演变
读写锁:一个资源可以被多个读线程访问,或者可以被一个写线程访问,但是不能同时存在读写线程,读写是互斥的,读读是共享的;
读写锁的演变过程如下图:

### 9.2)ReentrantReadWriteLock
ReentrantReadWriteLock 类的整体结构如下:
```java
public class ReentrantReadWriteLock implements ReadWriteLock, java.io.Serializable {
private static final long serialVersionUID = -6992448646407690164L;
/** 读锁 */
private final ReentrantReadWriteLock.ReadLock readerLock;
/** 写锁 */
private final ReentrantReadWriteLock.WriteLock writerLock;
/** Performs all synchronization mechanics */
final Sync sync;
/** 使用默认(非公平)的排序属性创建一个新的 ReentrantReadWriteLock */
public ReentrantReadWriteLock() {
this(false);
}
/** 使用给定的公平策略创建一个新的 ReentrantReadWriteLock */
public ReentrantReadWriteLock(boolean fair) {
sync = fair ? new FairSync() : new NonfairSync();
readerLock = new ReadLock(this);
writerLock = new WriteLock(this);
}
/** 返回用于写入操作的锁 */
public ReentrantReadWriteLock.WriteLock writeLock() { return writerLock; }
/** 返回用于读取操作的锁 */
public ReentrantReadWriteLock.ReadLock readLock() { return readerLock; }
...
```
可以看到,ReentrantReadWriteLock实现了ReadWriteLock接口, ReadWriteLock接口定义了获取读锁和写锁的规范,具体需要实现类去实现; 同时其还实现了Serializable接口,表示可以进行序列化,在源代码中可以看到ReentrantReadWriteLock实现了自己的序列化逻辑。
### 9.3)读写锁案例
场景: 使用ReentrantReadWriteLock 对一个hashmap进行读和写操作
不添加读写锁ReentrantReadWriteLock ,代码如下:
```java
//资源类
class MyCache {
//创建map集合
private volatile Map map = new HashMap<>();
//放数据
public void put(String key, Object value) {
System.out.println(Thread.currentThread().getName() + " 正在写操作" + key);
try {
//暂停一会
TimeUnit.MICROSECONDS.sleep(300);
} catch (InterruptedException e) {
e.printStackTrace();
}
//放数据
map.put(key, value);
System.out.println(Thread.currentThread().getName() + " 写完了" + key);
}
//取数据
public Object get(String key) {
Object result = null;
System.out.println(Thread.currentThread().getName() + " 正在读取操作" + key);
//暂停一会
try {
TimeUnit.MICROSECONDS.sleep(300);
} catch (InterruptedException e) {
e.printStackTrace();
}
result = map.get(key);
System.out.println(Thread.currentThread().getName() + " 取完了" + key);
return result;
}
}
public class ReadWriteLockDemo {
public static void main(String[] args) throws InterruptedException {
MyCache myCache = new MyCache();
//创建3个线程放数据
for (int i = 1; i <= 3; i++) {
final int num = i;
new Thread(() -> {
myCache.put(num + "", num + "");
}, String.valueOf(i)).start();
}
//创建3个线程取数据
for (int i = 1; i <= 3; i++) {
final int num = i;
new Thread(() -> {
myCache.get(num + "");
}, String.valueOf(i)).start();
}
}
}
```
输出:存在问题——写操作还未结束就进行了读操作
> 2 正在写操作2
> 3 正在写操作3
> 1 正在写操作1
> 1 正在读取操作1
> 2 正在读取操作2
> 3 正在读取操作3
> 3 取完了3
> 2 取完了2
> 2 写完了2
> 1 取完了1
> 3 写完了3
> 1 写完了1
添加读写锁ReentrantReadWriteLock ,代码如下:
```java
//资源类
class MyCache {
//创建map集合
private volatile Map map = new HashMap<>();
//创建读写锁对象
private ReadWriteLock rwLock = new ReentrantReadWriteLock();
//放数据
public void put(String key, Object value) {
//添加写锁
rwLock.writeLock().lock();
try {
System.out.println(Thread.currentThread().getName() + " 正在写操作" + key);
//暂停一会
TimeUnit.MICROSECONDS.sleep(300);
//放数据
map.put(key, value);
System.out.println(Thread.currentThread().getName() + " 写完了" + key);
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
//释放写锁
rwLock.writeLock().unlock();
}
}
//取数据
public Object get(String key) {
//添加读锁
rwLock.readLock().lock();
Object result = null;
try {
System.out.println(Thread.currentThread().getName() + " 正在读取操作" + key);
//暂停一会
TimeUnit.MICROSECONDS.sleep(300);
result = map.get(key);
System.out.println(Thread.currentThread().getName() + " 取完了" + key);
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
//释放读锁
rwLock.readLock().unlock();
}
return result;
}
}
public class ReadWriteLockDemo {
public static void main(String[] args) throws InterruptedException {
MyCache myCache = new MyCache();
//创建3个线程放数据
for (int i = 1; i <= 3; i++) {
final int num = i;
new Thread(() -> {
myCache.put(num + "", num + "");
}, String.valueOf(i)).start();
}
TimeUnit.MICROSECONDS.sleep(300);
//创建3个线程取数据
for (int i = 1; i <= 3; i++) {
final int num = i;
new Thread(() -> {
myCache.get(num + "");
}, String.valueOf(i)).start();
}
}
}
```
输出:
> 2 正在写操作2
> 2 写完了2
> 1 正在写操作1
> 1 写完了1
> 3 正在写操作3
> 3 写完了3
> 1 正在读取操作1
> 2 正在读取操作2
> 3 正在读取操作3
> 2 取完了2
> 3 取完了3
> 1 取完了1
### 9.4)读写锁降级
读写锁降级:将写锁降级为读锁【读锁不能升级为写锁】
锁降级目的:为了提高数据的可见性,如果没有写操作就会读不到,读不到无法执行锁可重入的过程
jdk8中的锁降级流程如下图:

将写锁降级为读锁代码演示:
```java
//演示读写锁降级【写锁降级为读锁】
public class ReadWriteLockDemo2 {
public static void main(String[] args) {
//可重入读写锁对象
ReentrantReadWriteLock rwLock = new ReentrantReadWriteLock();
ReentrantReadWriteLock.ReadLock readLock = rwLock.readLock();//读锁
ReentrantReadWriteLock.WriteLock writeLock = rwLock.writeLock();//写锁
//锁降级
//1 获取写锁
writeLock.lock();
System.out.println("---write");
//2 获取读锁
readLock.lock();
System.out.println("---read");
//3 释放写锁
//writeLock.unlock();
//4 释放读锁
//readLock.unlock();
}
}
```
输出:
> ---write
> ---read
读锁不能升级为写锁代码演示:
```java
//演示读写锁降级【读锁不能升级为写锁】
public class ReadWriteLockDemo2 {
public static void main(String[] args) {
//可重入读写锁对象
ReentrantReadWriteLock rwLock = new ReentrantReadWriteLock();
ReentrantReadWriteLock.ReadLock readLock = rwLock.readLock();//读锁
ReentrantReadWriteLock.WriteLock writeLock = rwLock.writeLock();//写锁
//锁降级
//2 获取读锁
readLock.lock();
System.out.println("---read");
//1 获取写锁
writeLock.lock();
System.out.println("---write");
//3 释放写锁
//writeLock.unlock();
//4 释放读锁
//readLock.unlock();
}
}
```
输出:
> ---read
### 9.5)小结
- 在线程持有读锁的情况下,该线程不能取得写锁(因为获取写锁的时候,如果发现当前的读锁被占用,就马上获取失败,不管读锁是不是被当前线程持有);
- 在线程持有写锁的情况下,该线程可以继续获取读锁(获取读锁时如果发现写 锁被占用,只有写锁没有被当前线程占用的情况才会获取失败);
原因: 当线程获取读锁的时候,可能有其他线程同时也在持有读锁,因此不能把 获取读锁的线程“升级”为写锁;而对于获得写锁的线程,它一定独占了读写锁,因此可以继续让它获取读锁,当它同时获取了写锁和读锁后,还可以先释放写锁继续持有读锁,这样一个写锁就“降级”为了读锁。
## 十、阻塞队列
### 10.1)BlockingQueue 简介
Concurrent包中,BlockingQueue很好的解决了多线程中,如何高效安全“传输”数据的问题;通过这些高效并且线程安全的队列类,为我们快速搭建高质量的多线程程序带来极大的便利。
阻塞队列,顾名思义,首先它是一个队列, 通过一个共享的队列,可以使得数据由队列的一端输入,从另外一端输出,如下图所示:

阻塞情况如下:
- 当队列是空的,从队列中获取元素的操作将会被阻塞;
- 当队列是满的,从队列中添加元素的操作将会被阻塞;
- 试图从空的队列中获取元素的线程将会被阻塞,直到其他线程往空的队列插入新的元素;
- 试图向已满的队列中添加新元素的线程将会被阻塞,直到其他线程从队列中移除一个或多个元素或者完全清空,使队列变得空闲起来并后续新增
#### 10.1.1)常用的队列
常用的队列主要有以下两种:
- 先进先出(FIFO):先插入的队列的元素也最先出队列,类似于排队的功能,从某种程度上来说这种队列也体现了一种公平性 (队列)
- 后进先出(LIFO):后插入队列的元素最先出队列,这种队列优先处理最近发生的事件(栈)
#### 10.1.2)为什么需要BlockingQueue
在concurrent包发布以前,在多线程环境下,开发者都必须去自己控制阻塞线程的细节,尤其还要兼顾效率和线程安全,而这会给我们的程序带来不小的复杂度;
使用BlockingQueue 之后,开发者不需要关心什么时候需要阻塞线程,什么时候需要唤醒线程,因为这一切 都由BlockingQueue一手包办了。
#### 10.1.3)BlockingQueue的实际应用
多线程环境中,通过队列可以很容易实现数据共享,比如经典的“生产者”和 “消费者”模型中,通过队列可以很便利地实现两者之间的数据共享。
假设有若干生产者线程,另外又有若干个消费者线程。如果生产者线程需要把准备好的数据共享给消费者线程,利用队列的方式来传递数据,就可以很方便地解决它们之间的数据共享问题。
> 但如果生产者和消费者在某个时间段内,万一发生数据处理速度不匹配的情况如何解决呢?
>
> 1. 理想情况下,如果生产者产出数据的速度大于消费者消费的速度,并且当生产出来的数据累积到一定程度的时候,那么生产者必须暂停等待一下(阻塞生产者线程),以便等待消费者线程把累积的 数据处理完毕,反之亦然;
> 2. 当队列中没有数据的情况下,消费者端的所有线程都会被自动阻塞(挂起), 直到有数据放入队列;
> 3. 当队列中填满数据的情况下,生产者端的所有线程都会被自动阻塞(挂起), 直到队列中有空的位置,线程被自动唤醒
### 10.2)BlockingQueue 核心方法
BlockingQueue的核心方法,如下图所示:

#### 10.2.1)放入数据
1. offer(anObject):表示如果可能的话,将anObject加到BlockingQueue 里,即如果BlockingQueue可以容纳,则返回true,否则返回false;本方法不阻塞当前执行方法的线程)
2. offer(E o, long timeout, TimeUnit unit):可以设定等待的时间,如果在指定的时间内,还不能往队列中加入BlockingQueue,则返回失败
3. put(anObject):把anObject加到BlockingQueue 里,如果BlockQueue没有空间,则调用此方法的线程被阻断直到BlockingQueue里面有空间再继续
#### 10.2.2)获取数据
1. poll(time): 取走BlockingQueue里排在首位的对象,若不能立即取出,则可以等time参数规定的时间,还取不到时返回null;
2. poll(long timeout, TimeUnit unit):从BlockingQueue取出一个队首的对象, 如果在指定时间内,队列一旦有数据可取则立即返回队列中的数据;否则直到时间超时还没有数据可取则返回失败。
3. take():取走BlockingQueue 里排在首位的对象,若BlockingQueue为空,阻断进入等待状态直到BlockingQueue有新的数据被加入;
4. drainTo(): 一次性从BlockingQueue获取所有可用的数据对象(还可以指定 获取数据的个数),通过该方法可以提升获取数据效率,不需要多次分批加锁或释放锁。
### 10.3)常见的 BlockingQueue
#### 10.3.1) ArrayBlockingQueue(常用)
基于数组的阻塞队列实现,在ArrayBlockingQueue内部,维护了一个定长数组以便缓存队列中的数据对象,这是一个常用的阻塞队列,除了一个定长数组外,ArrayBlockingQueue 内部还保存着两个整形变量,分别标识着队列的头部和尾部在数组中的位置。
ArrayBlockingQueue在生产者放入数据和消费者获取数据,都是共用同一个锁对象,由此也意味着两者无法真正并行运行,这点尤其不同于 LinkedBlockingQueue;按照实现原理来分析,ArrayBlockingQueue 完全可以采用分离锁,从而实现生产者和消费者操作的完全并行运行。
Doug Lea之 所以没这样去做,也许是因为ArrayBlockingQueue的数据写入和获取操作已经足够轻巧,以至于引入独立的锁机制,除了给代码带来额外的复杂性外,其在性能上完全占不到任何便宜。
ArrayBlockingQueue和 LinkedBlockingQueue 间还有一个明显的不同之处在于,前者在插入或删除元素时不会产生或销毁任何额外的对象实例,而后者则会生成一个额外的Node对象。
这在长时间内需要高效并发地处理大批量数据的系统中,其对于 GC的影响还是存在一定的区别,而在创建ArrayBlockingQueue时,还可以控制对象的内部锁是否采用公平锁,默认采用非公平锁。
> 总结:由数组结构组成的有界阻塞队列
#### 10.3.2)LinkedBlockingQueue(常用)
基于链表的阻塞队列,同ArrayListBlockingQueue 类似,其内部也维持着一个数据缓冲队列(该队列由一个链表构成),当生产者往队列中放入一个数据时,队列会从生产者手中获取数据,并缓存在队列内部,而生产者立即返回;
只有当队列缓冲区达到最大值缓存容量时(LinkedBlockingQueue可以通过构造函数指定该值),才会阻塞生产者队列,直到消费者从队列中消费掉一份数据,生产者线程会被唤醒,反之对于消费者这端的处理也基于同样的原理。
而LinkedBlockingQueue 之所以能够高效的处理并发数据,还因为其对于生产者端和消费者端分别采用了独立的锁来控制数据同步,这也意味着在高并发的情况下生产者和消费者可以并行地操作队列中的数据,以此来提高整个队列的并发性能。
> 总结:由链表结构组成的有界(但大小默认值为 integer.MAX_VALUE)阻塞队列
#### 10.3.3)DelayQueue
DelayQueue中的元素只有当其指定的延迟时间到了,才能够从队列中获取到该元素。
DelayQueue是一个没有大小限制的队列,因此往队列中插入数据的操作(生产者)永远不会被阻塞,而只有获取数据的操作(消费者)才会被阻塞。
> 总结: 使用优先级队列实现的延迟无界阻塞队列
#### 10.3.4) PriorityBlockingQueue
基于优先级的阻塞队列(优先级的判断通过构造函数传入的Compator对象来决定),但需要注意的是PriorityBlockingQueue并不会阻塞数据生产者,而只会在没有可消费的数据时,阻塞数据的消费者。
因此使用的时候要特别注意,生产者生产数据的速度绝对不能快于消费者消费 数据的速度,否则时间一长,会最终耗尽所有的可用堆内存空间。
在实现PriorityBlockingQueue 时,内部控制线程同步的锁采用的是公平锁。
> 总结: 支持优先级排序的无界阻塞队列
### 10.4)ArrayBlockingQueue案例
#### 10.4.1) add && remove
ArrayBlockingQueue的 add && remove案例展示代码如下:
```java
//阻塞队列
public class BlockingQueueDemo1 {
public static void main(String[] args) throws InterruptedException {
//创建阻塞队列
BlockingQueue blockingQueue = new ArrayBlockingQueue<>(3);
//第一组: add && remove
System.out.println(blockingQueue.add("a"));
System.out.println(blockingQueue.add("b"));
System.out.println(blockingQueue.add("c"));
// System.out.println(blockingQueue.element());
// 队列已满,再添加元素会报错
System.out.println(blockingQueue.add("w"));
System.out.println(blockingQueue.remove());
System.out.println(blockingQueue.remove());
System.out.println(blockingQueue.remove());
// 队列已空,再删除元素会报错
System.out.println(blockingQueue.remove());
}
}
```
输出:队列已满,再添加元素会报错
> Exception in thread "main" java.lang.IllegalStateException: Queue full
> at java.util.AbstractQueue.add(AbstractQueue.java:98)
> at java.util.concurrent.ArrayBlockingQueue.add(ArrayBlockingQueue.java:312)
> at com.study.queue.BlockingQueueDemo1.main(BlockingQueueDemo1.java:19)
------
```java
//阻塞队列
public class BlockingQueueDemo1 {
public static void main(String[] args) throws InterruptedException {
//创建阻塞队列
BlockingQueue blockingQueue = new ArrayBlockingQueue<>(3);
//第一组: add && remove
System.out.println(blockingQueue.add("a"));
System.out.println(blockingQueue.add("b"));
System.out.println(blockingQueue.add("c"));
System.out.println(blockingQueue.element());
// 队列已满,再添加元素会报错
// System.out.println(blockingQueue.add("w"));
System.out.println(blockingQueue.remove());
System.out.println(blockingQueue.remove());
System.out.println(blockingQueue.remove());
// 队列已空,再删除元素会报错
System.out.println(blockingQueue.remove());
}
}
```
输出:队列已空,再删除元素会报错
> Exception in thread "main" java.util.NoSuchElementException
> at java.util.AbstractQueue.remove(AbstractQueue.java:117)
> at com.study.queue.BlockingQueueDemo1.main(BlockingQueueDemo1.java:22)
------
```java
//阻塞队列
public class BlockingQueueDemo1 {
public static void main(String[] args) throws InterruptedException {
//创建阻塞队列
BlockingQueue blockingQueue = new ArrayBlockingQueue<>(3);
//第一组: add && remove
System.out.println(blockingQueue.add("a"));
System.out.println(blockingQueue.add("b"));
System.out.println(blockingQueue.add("c"));
System.out.println(blockingQueue.element());
// 队列已满,再添加元素会报错
// System.out.println(blockingQueue.add("w"));
System.out.println(blockingQueue.remove());
System.out.println(blockingQueue.remove());
System.out.println(blockingQueue.remove());
// 队列已空,再删除元素会报错
// System.out.println(blockingQueue.remove());
}
}
```
输出:
> true
> true
> true
> a
> a
> b
> c
#### 10.4.2)offer && poll
ArrayBlockingQueue的 offer && poll 案例展示代码如下:
```java
//阻塞队列
public class BlockingQueueDemo2 {
public static void main(String[] args) throws InterruptedException {
//创建阻塞队列
BlockingQueue blockingQueue = new ArrayBlockingQueue<>(3);
//第二组:offer && poll
System.out.println(blockingQueue.offer("a"));
System.out.println(blockingQueue.offer("b"));
System.out.println(blockingQueue.offer("c"));
// 输出: false
System.out.println(blockingQueue.offer("www"));
System.out.println(blockingQueue.poll());
System.out.println(blockingQueue.poll());
System.out.println(blockingQueue.poll());
// 输出:null
System.out.println(blockingQueue.poll());
}
}
```
输出:
> true
> true
> true
> false
> a
> b
> c
> null
#### 10.4.3)put && take
ArrayBlockingQueue的 put && take 案例展示代码如下:
```java
//阻塞队列
public class BlockingQueueDemo3 {
public static void main(String[] args) throws InterruptedException {
//创建阻塞队列
BlockingQueue blockingQueue = new ArrayBlockingQueue<>(3);
//第三组:put && take
blockingQueue.put("a");
blockingQueue.put("b");
blockingQueue.put("c");
// 输出:程序一直未结束,因为空间不够,所以处于阻塞状态,直至将数据插入才结束
//blockingQueue.put("w");
System.out.println(blockingQueue.take());
System.out.println(blockingQueue.take());
System.out.println(blockingQueue.take());
// 输出:程序一直未结束,因为队列为空没有取到元素,所以处于阻塞状态,直至将数据取到才结束
// System.out.println(blockingQueue.take());
}
}
```
输出:
> a
> b
> c
#### 10.4.4)offer
ArrayBlockingQueue的offer方法案例展示代码如下:
```java
//阻塞队列
public class BlockingQueueDemo4 {
public static void main(String[] args) throws InterruptedException {
//创建阻塞队列
BlockingQueue blockingQueue = new ArrayBlockingQueue<>(3);
//第四组 :offer
System.out.println(blockingQueue.offer("a"));
System.out.println(blockingQueue.offer("b"));
System.out.println(blockingQueue.offer("c"));
// 输出:false
System.out.println(blockingQueue.offer("w",3L, TimeUnit.SECONDS));
}
}
```
输出:
> true
> true
> true
>
> 等待3秒后...
>
> false
## 十一、ThreadPool 线程池
### 11.1)线程池简介
线程池(英语:thread pool):一种线程使用模式。线程过多会带来调度开销, 进而影响缓存局部性和整体性能;而线程池维护着多个线程,等待着监督管理者分配可并发执行的任务;这避免了在处理短时间任务时创建与销毁线程的代价。
程池不仅能够保证内核的充分利用,还能防止过分调度。
#### 11.1.1)线程池简介
线程池做的工作只要是控制运行的线程数量,处理过程中将任务放入队列,然后在线程创建后启动这些任务,如果线程数量超过了最大数量, 超出数量的线程排队等候,等其他线程执行完毕,再从队列中取出任务来执行。
#### 11.1.2)线程池特点
线程池的主要特点为:
- 降低资源消耗:通过重复利用已创建的线程降低线程创建和销毁造成的销耗;
- 提高响应速度:当任务到达时,任务可以不需要等待线程创建就能立即执行;
- 提高线程的可管理性:线程是稀缺资源,如果无限制的创建,不仅会销耗系统资源,还会降低系统的稳定性,使用线程池可以进行统一的分配,调优和监控;
- Java 中的线程池是通过 Executor 框架实现的,该框架中用到了 Executor,Executors, ExecutorService,ThreadPoolExecutor这几个类 ,其关系如下图:

### 11.2)线程池参数说明
常用参数:
- corePoolSize:线程池的常驻线程数量(核心)
- maximumPoolSize:线程池能容纳的最大线程数
- keepAliveTime空闲:线程存活时间
- unit :线程存活的时间单位
- BlockingQueue workQueue :存放提交但未执行任务的队列 【阻塞队列】
- threadFactory: 创建线程的工厂类【线程工厂】
- handler :等待队列满后的拒绝策略
线程池中,有三个重要的参数,决定影响了拒绝策略:
corePoolSize - 核心线程数,也即最小的线程数;
workQueue - 阻塞队列;
maximumPoolSize - 最大线程数
当提交任务数大于 corePoolSize 的时候,会优先将任务放到 workQueue 阻塞队列中;当阻塞队列饱和后,会扩充线程池中线程数,直到达到 maximumPoolSize 最大线程数配置。此时再多余的任务,则会触发线程池的拒绝策略。
> 总结:当提交的任务数大于(workQueue.size() + maximumPoolSize ),就会触发线程池的拒绝策略
### 11.3)线程池底层工作原理
线程池底层工作原理如下图所示:

线程池底层工作流程:
1. 在创建了线程池后,线程池中的线程数为零;
2. 当调用execute()方法添加一个请求任务时,线程池会做出如下判断:
2.1) 如果正在运行的线程数量小于corePoolSize,那么马上创建线程运行这个任务;
2.2)如果正在运行的线程数量大于或等于corePoolSize,那么将这个任务放入队列;
2.3 )如果这个时候队列满了且正在运行的线程数量还小于 maximumPoolSize,那么还是要创建非核心线程立刻运行这个任务;
2.4 )如果队列满了且正在运行的线程数量大于或等于maximumPoolSize,那么线程池会启动饱和拒绝策略来执行。
3. 当一个线程完成任务时,它会从队列中取下一个任务来执行;
4. 当一个线程无事可做超过一定的时间(keepAliveTime)时,线程会判断:
4.1)如果当前运行的线程数大于corePoolSize,那么这个线程就被停掉;
4.2)所以线程池的所有任务完成后,它最终会收缩到corePoolSize的大小
### 11.4)线程池的拒绝策略
JDK内置的拒绝策略如下图所示:

说明:
1. AbortPolicy:线程池默认的拒绝策略;丢弃任务并抛出拒绝执行 RejectedExecutionException 异常信息;必须处理好抛出的异常,否则会打断当前的执 流程,影响后续的任务执行;
2. CallerRunsPolicy:当触发拒绝策略,只要线程池没有关闭的话,则使用调用线程直接运行任务,将任务退回给调用者;一般并发比较小,性能要求不高,不允许失败。但是由于调用者自己运行任务,如果任务提交速度过快,可能导致程序阻塞,性能损失;
3. DiscardOldestPolicy: 当触发拒绝策略,只要线程池没有关闭的话,丢弃阻塞队列 workQueue 中最老的一个任务,并将新任务加入提交;
4. DiscardPolicy: 直接丢弃,其他啥都不做
### 11.5)线程池使用方式
#### 11.5.1)newFixedThreadPool
##### 11.5.1.1)作用
【一池N线程】创建一个可重用固定线程数的线程池,以共享的无界队列方式来运行这些线程。
在任意点,在大多数线程会处于处理任务的活动状态;如果在所有线程处于活动状态时提交附加任务,则在有可用线程之前,附加任务将在队列中等待。
如果在关闭前的执行期间由于失败而导致任何线程终止,那么一个新线 程将代替它执行后续的任务(如果需要),在某个线程被显式地关闭之前,池中的线程将一直存在。
##### 11.5.1.2)特征
- 线程池中的线程处于一定的量,可以很好的控制线程的并发量;
- 线程可以重复被使用,在显示关闭之前,都将一直存在 ;
- 超出一定量的线程被提交时候需在队列中等待。
##### 11.5.1.3)适用场景
适用于可以预测线程数量的业务中,或者服务器负载较重,对线程数有严格限制的场景
##### 11.5.1.4)代码示例
```java
//演示newFixedThreadPool创建线程池【一池N线程】
public class ThreadPoolDemo1 {
public static void main(String[] args) {
//一池N线程
//5个窗口
ExecutorService threadPool1 = Executors.newFixedThreadPool(5);
//10个顾客请求
try {
for (int i = 1; i <= 10; i++) {
//执行
threadPool1.execute(() -> {
System.out.println("当前线程名:" + Thread.currentThread().getName() + " 办理业务");
});
}
} catch (Exception e) {
e.printStackTrace();
} finally {
//关闭
threadPool1.shutdown();
}
}
}
```
输出:
> 当前线程名:pool-1-thread-1 办理业务
> 当前线程名:pool-1-thread-5 办理业务
> 当前线程名:pool-1-thread-1 办理业务
> 当前线程名:pool-1-thread-4 办理业务
> 当前线程名:pool-1-thread-1 办理业务
> 当前线程名:pool-1-thread-2 办理业务
> 当前线程名:pool-1-thread-3 办理业务
> 当前线程名:pool-1-thread-1 办理业务
> 当前线程名:pool-1-thread-4 办理业务
> 当前线程名:pool-1-thread-5 办理业务
#### 11.5.2) newSingleThreadExecutor
##### 11.5.2.1)作用
【一池一线程】创建一个使用单个 worker 线程的 Executor,以无界队列方式来运行该线程,可保证顺序地执行各个任务,并且在任意给定的时间不会有多个线程是活动的。
与其他等效的 newFixedThreadPool不同,可保证无需重新配置此方法所返回的执行程序即 可使用其他的线程。
(注意,如果因为在关闭前的执行期间出现失败而终止了此单个线程, 那么如果需要,一个新线程将代替它执行后续的任务)。
> 一个任务一个任务执行,一池一线程
##### 11.5.2.2)特征
线程池中最多执行1个线程,之后提交的线程活动将会排在队列中以此执行
##### 11.5.2.3)特征
适用于需要保证顺序执行各个任务,并且在任意时间点,不会同时有多个 线程的场景
##### 11.5.2.4)代码示例
```java
//演示newSingleThreadExecutor创建线程池【一池一线程】
public class ThreadPoolDemo1 {
public static void main(String[] args) {
//一池一线程
//一个窗口
ExecutorService threadPool2 = Executors.newSingleThreadExecutor();
//10个顾客请求
try {
for (int i = 1; i <= 10; i++) {
//执行
threadPool2.execute(() -> {
System.out.println("当前线程名:" + Thread.currentThread().getName() + " 办理业务");
});
}
} catch (Exception e) {
e.printStackTrace();
} finally {
//关闭
threadPool2.shutdown();
}
}
}
```
输出:
> 当前线程名:pool-2-thread-1 办理业务
> 当前线程名:pool-2-thread-1 办理业务
> 当前线程名:pool-2-thread-1 办理业务
> 当前线程名:pool-2-thread-1 办理业务
> 当前线程名:pool-2-thread-1 办理业务
> 当前线程名:pool-2-thread-1 办理业务
> 当前线程名:pool-2-thread-1 办理业务
> 当前线程名:pool-2-thread-1 办理业务
> 当前线程名:pool-2-thread-1 办理业务
> 当前线程名:pool-2-thread-1 办理业务
#### 11.5.3) newCachedThreadPool
##### 11.5.3.1)作用
【一池可扩容线程】创建一个可缓存线程池,如果线程池长度超过处理需要,可灵活回收空闲线程,若无可回收,则新建线程。
> 线程池根据需求创建线程,可扩容,遇强则强
##### 11.5.3.2)特点
- 线程池中数量没有固定,可达到最大值(Interger. MAX_VALUE);
- 线程池中的线程可进行缓存重复利用和回收(回收默认时间为1分钟);
- 当线程池中,没有可用线程,会重新创建一个线程
##### 11.5.3.3)适用场景
适用于创建一个可无限扩大的线程池,服务器负载压力较轻,执行时间较 短,任务多的场景
##### 11.5.3.4)代码示例
```java
//演示newCachedThreadPool创建线程池【一池可扩容线程】
public class ThreadPoolDemo1 {
public static void main(String[] args) {
//一池可扩容线程
ExecutorService threadPool3 = Executors.newCachedThreadPool();
//10个顾客请求
try {
for (int i = 1; i <= 10; i++) {
//执行
threadPool3.execute(() -> {
System.out.println("当前线程名:" + Thread.currentThread().getName() + " 办理业务");
});
}
} catch (Exception e) {
e.printStackTrace();
} finally {
//关闭
threadPool3.shutdown();
}
}
}
```
输出:
> 当前线程名:pool-3-thread-2 办理业务
> 当前线程名:pool-3-thread-5 办理业务
> 当前线程名:pool-3-thread-3 办理业务
> 当前线程名:pool-3-thread-7 办理业务
> 当前线程名:pool-3-thread-4 办理业务
> 当前线程名:pool-3-thread-1 办理业务
> 当前线程名:pool-3-thread-8 办理业务
> 当前线程名:pool-3-thread-6 办理业务
> 当前线程名:pool-3-thread-10 办理业务
> 当前线程名:pool-3-thread-9 办理业务
### 11.6)自定义线程池
#### 11.6.1)自定义线程池的原因
项目中创建多线程时,使用常见的三种线程池创建方式,newFixedThreadPool【一池N线程】、newSingleThreadExecutor【一池一线程】、newCachedThreadPool【一池可扩容线程】都有一定问题,原因是FixedThreadPool和SingleThreadExecutor底层都是用 LinkedBlockingQueue实现的,这个队列最大长度为Integer.MAX_VALUE,容易导致OOM。
为什么不允许适用不允许Executors.的方式手动创建线程池,如下图 :

所以实际生产一般自己通过ThreadPoolExecutor的7个参数来自定义线程池
#### 11.6.2)创建自定义线程池
创建线程池推荐适用ThreadPoolExecutor及其7个参数手动创建:
- corePoolSize:线程池的核心线程数;
- maximumPoolSize:能容纳的最大线程数;
- keepAliveTime:空闲线程存活时间;
- unit:存活的时间单位;
- workQueue:存放提交但未执行任务的队列;
- threadFactory:创建线程的工厂类;
- handler:等待队列满后的拒绝策略
#### 11.6.3)自定义线程池代码示例
自定义线程池代码示例:
```java
//自定义线程池创建
public class ThreadPoolDemo2 {
public static void main(String[] args) {
// 自定义线程池
ExecutorService threadPool = new ThreadPoolExecutor(
2,
5,
2L,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(3),
Executors.defaultThreadFactory(),
new ThreadPoolExecutor.AbortPolicy()
);
//10个顾客请求
try {
for (int i = 1; i <= 10; i++) {
//执行
threadPool.execute(() -> {
System.out.println(Thread.currentThread().getName() + " 办理业务");
});
}
} catch (Exception e) {
e.printStackTrace();
} finally {
//关闭
threadPool.shutdown();
}
}
}
```
输出:
> pool-1-thread-2 办理业务
> pool-1-thread-5 办理业务
> pool-1-thread-5 办理业务
> pool-1-thread-4 办理业务
> pool-1-thread-1 办理业务
> pool-1-thread-3 办理业务
> pool-1-thread-5 办理业务
> pool-1-thread-2 办理业务
> java.util.concurrent.RejectedExecutionException: Task com.study.pool.ThreadPoolDemo2$$Lambda$1/1096979270@7ba4f24f rejected from java.util.concurrent.ThreadPoolExecutor@3b9a45b3[Running, pool size = 5, active threads = 3, queued tasks = 0, completed tasks = 5]
> at java.util.concurrent.ThreadPoolExecutor$AbortPolicy.rejectedExecution(ThreadPoolExecutor.java:2063)
> at java.util.concurrent.ThreadPoolExecutor.reject(ThreadPoolExecutor.java:830)
> at java.util.concurrent.ThreadPoolExecutor.execute(ThreadPoolExecutor.java:1379)
> at com.study.pool.ThreadPoolDemo2.main(ThreadPoolDemo2.java:23)
## 十二、Fork/Join 框架
### 12.1)Fork/Join框架简介
Fork/Join 可以将一个大的任务拆分成多个子任务进行并行处理,最后将子任务结果合并成最后的计算结果,并进行输出。
Fork/Join框架要完成两件事情:
- Fork:把一个复杂任务进行分拆,大事化小 ;
- Join:把分拆任务的结果进行合并
### 12.2)Fork/Join框架的实现过程
Fork/Join框架的实现过程如下:
1. 任务分割:首先Fork/Join框架需要把大的任务分割成足够小的子任务,如果 子任务比较大的话还要对子任务进行继续分割 ;
2. 执行任务并合并结果:分割的子任务分别放到双端队列里,然后几个启动线程分别从双端队列里获取任务执行;子任务执行完的结果都放在另外一个队列里, 启动一个线程从队列里取数据,然后合并这些数据。
如下图所示:

### 12.3)Fork/Join框架的实现原理
在Java的Fork/Join框架中,ForkJoinPool【分支合并池】由ForkJoinTask数组和ForkJoinWorkerThread数组组成;使用它们完成任务分割和合并结果操作 :
1. ForkJoinTask数组:负责将存放以及将程序提交给ForkJoinPool;需要通过ForkJoinPool来执行;要使用 Fork/Join 框架,首先需要创建一个 ForkJoin 任务,该类提供了在任务中执行fork和join的机制;通常情况下不需要直接集成ForkJoinTask类,只需要继承它的子类;
> Fork/Join框架提供了两个子类:
>
> - RecursiveAction:用于没有返回结果的任务;
>
> - RecursiveTask:用于有返回结果的任务,继承后可以实现递归(自己调自己)调用的任务
>
> 
2. ForkJoinWorkerThread数组:负责执行这些任务;
ForkJoinPool【分支合并池】类似于 线程池,如下图:

#### 12.3.1)Fork方法
Fork(任务分割):把一个大任务拆分成小任务,如果小任务还是太大还会继续拆分,直到足够小
实现原理: 当我们调用ForkJoinTask的fork方法时,程序会把任务放在ForkJoinWorkerThread的pushTask的workQueue中,异步地执行这个任务,然后立即返回结果 ,源码如下:
```java
public final ForkJoinTask fork() {
Thread t;
if ((t = Thread.currentThread()) instanceof ForkJoinWorkerThread)
((ForkJoinWorkerThread)t).workQueue.push(this);
else
ForkJoinPool.common.externalPush(this);
return this;
}
```
pushTask方法把当前任务存放在ForkJoinTask数组队列里。然后再调用 ForkJoinPool的signalWork()方法唤醒或创建一个工作线程来执行任务,源码如下:
```java
final void push(ForkJoinTask> task) {
ForkJoinTask>[] a; ForkJoinPool p;
int b = base, s = top, n;
if ((a = array) != null) { // ignore if queue removed
int m = a.length - 1; // fenced write for task visibility
U.putOrderedObject(a, ((m & s) << ASHIFT) + ABASE, task);
U.putOrderedInt(this, QTOP, s + 1);
if ((n = s - b) <= 1) {
if ((p = pool) != null)
p.signalWork(p.workQueues, this);
}
else if (n >= m)
growArray();
}
}
```
#### 12.3.2)join方法
Join(合并结果):阻塞当前线程并等待获取结果,将每个任务执行完成得到的结果进行汇总,返回最终结果;
ForkJoinTask 的 join方法的实现,源码如下:
```java
public final V join() {
int s;
if ((s = doJoin() & DONE_MASK) != NORMAL)
reportException(s);
return getRawResult();
}
```
通过调用doJoin方法,通过doJoin()方法得到当前任务的状态来判断返回 什么结果,任务状态有4种
```java
static final int NORMAL = 0xf0000000; // must be negative[已完成]
static final int CANCELLED = 0xc0000000; // must be < NORMAL[被取消]
static final int EXCEPTIONAL = 0x80000000; // must be < CANCELLED[出现异常]
static final int SIGNAL = 0x00010000; // must be >= 1 << 16[信号]
```
- 如果任务状态是已完成,则直接返回任务结果;
- 如果任务状态是被取消,则直接抛出CancellationException;
- 如果任务状态是抛出异常,则直接抛出对应的异常
doJoin方法的实现 ,源码如下:
```java
private int doJoin() {
int s; Thread t; ForkJoinWorkerThread wt; ForkJoinPool.WorkQueue w;
return (s = status) < 0 ? s :
((t = Thread.currentThread()) instanceof ForkJoinWorkerThread) ?
(w = (wt = (ForkJoinWorkerThread)t).workQueue).
tryUnpush(this) && (s = doExec()) < 0 ? s :
wt.pool.awaitJoin(w, this, 0L) :
externalAwaitDone();
}
```
doJoin()方法流程如下:
1. 首先通过查看任务的状态,看任务是否已经执行完成,如果执行完成,则直接返回任务状态;
2. 如果没有执行完,则从任务数组里取出任务并执行;
3. 如果任务顺利执行完成,则设置任务状态为NORMAL,如果出现异常,则记录异常,并将任务状态设置为EXCEPTIONAL
#### 12.3.3)Fork/Join框架的异常处理
ForkJoinTask在执行的时候可能会抛出异常,但是我们没办法在主线程里直接捕获异常,所以ForkJoinTask提供了isCompletedAbnormally()方法来检查任务是否已经抛出异常或已经被取消了,并且可以通过ForkJoinTask的 getException方法获取异常。
getException方法返回Throwable对象,如果任务被取消了则返回 CancellationException,如果任务没有完成或者没有抛出异常则返回null
### 12.4)Fork/Join案例
场景: 生成一个计算任务,计算1+2+3.........+100,10个数切分一个子任务
代码如下:
```java
class MyTask extends RecursiveTask {
//拆分差值不能超过10,计算10以内运算
private static final Integer VALUE = 10;
//拆分开始值
private int begin;
//拆分结束值
private int end;
//返回结果
private int result;
//创建有参数构造
public MyTask(int begin, int end) {
this.begin = begin;
this.end = end;
}
//拆分和合并过程
@Override
protected Integer compute() {
//判断相加两个数值是否大于10
if ((end - begin) <= VALUE) {
//相加操作
for (int i = begin; i <= end; i++) {
result = result + i;
}
} else {
//进一步拆分
//获取中间值
int middle = (begin + end) / 2;
//拆分左边
MyTask task01 = new MyTask(begin, middle);
//拆分右边
MyTask task02 = new MyTask(middle + 1, end);
//调用方法拆分
task01.fork();
task02.fork();
//合并结果
result = task01.join() + task02.join();
}
return result;
}
}
public class ForkJoinDemo {
public static void main(String[] args) throws ExecutionException, InterruptedException {
//创建MyTask对象
MyTask myTask = new MyTask(0, 100);
//创建分支合并池对象
ForkJoinPool forkJoinPool = new ForkJoinPool();
ForkJoinTask forkJoinTask = forkJoinPool.submit(myTask);
//获取最终合并之后结果
Integer result = forkJoinTask.get();
System.out.println("计算结果 = " + result);
//关闭池对象
forkJoinPool.shutdown();
}
}
```
输出:
> 计算结果 = 5050
### 12.5)小结
ForkJoin框架主要依赖ForkJoinTask来表示任务,而ForkJoinTask依赖ForkJoinPool来执行任务,ForkJoinPool会创建工作线程去执行任务,执行任务调用的是ForkJoinTask的子类RecursiveTask、RecursiveAction的compute方法。
**compute方法需要我们自己实现,在compute方法中可以定义如何拆分方法与组装结果**,在拆分成功后调用ForkJoinTask的fork方法可以把任务加载到当前工作线程的队列中。
工作线程执行完一个任务后会从队列中去获取任务继续执行,当队列中没有任务的时候则会从其他工作线程的队列中获取任务。
## 十三、CompletableFuture
### 13.1)CompletableFuture简介
CompletableFuture在 Java里面被用于异步编程,异步通常意味着非阻塞,可以使得我们的任务单独运行在与主线程分离的其他线程中,并且通过回调可以在主线程中得到异步任务的执行状态,是否完成,和是否异常等信息。
CompletableFuture实现了Future、CompletionStage接口,实现了Future 接口就可以兼容现在有线程池框架,而CompletionStage接口才是异步编程的接口抽象,里面定义多种异步方法,通过这两者集合从而打造出了强大的 CompletableFuture类。
### 13.2)Future与CompletableFuture
Futrue在Java里面,通常用来表示一个异步任务的引用,比如我们将任务提交到线程池里面,然后会得到一个Futrue,在Future里面有isDone() 方法来判断任务是否处理结束,还有get() 方法可以一直阻塞直到任务结束然后获取结果,但整体来说这种方式,还是同步的,因为需要客户端不断阻塞等待或者不断轮询才能知道任务是否完成。
#### 13.2.1)Future的主要缺点
Future的主要缺点如下:
1. 不支持手动完成:提交了一个任务,但是执行太慢了,通过其他路径已经获取到了任务结果, 现在没法把这个任务结果通知到正在执行的线程,所以必须主动取消或者一直等待它执行完成 ;
2. 不支持进一步的非阻塞调用:通过Future的get方法会一直阻塞到任务完成,但是想在获取任务之后执行额外的任务,因为Future不支持回调函数,所以无法实现这个功能;
3. 不支持链式调用:对于Future的执行结果,想继续传到下一个Future处理使用,从而形成 一个链式的pipline调用,这在Future中是没法实现的;
4. 不支持多个Future合并:有10个Future并行执行,我们想在所有的Future运行完毕之后执行某些函数,是没法通过Future实现的;
5. 不支持异常处理 :Future的API没有任何的异常处理的api,所以在异步运行时,如果出了问题是不好定位的;
### 13.3)CompletableFuture的代码示例
#### 13.3.1)异步调用【无返回值】
代码如下:
```java
public class CompletableFutureDemo {
public static void main(String[] args) throws Exception {
// 异步调用 无返回值
CompletableFuture completableFuture1 = CompletableFuture.runAsync(()->{
System.out.println(Thread.currentThread().getName()+" : CompletableFuture1");
});
completableFuture1.get();
}
}
```
输出:
> ForkJoinPool.commonPool-worker-1 : CompletableFuture1
#### 13.3.2)异步调用【有返回值】
代码如下:
```java
public class CompletableFutureDemo {
public static void main(String[] args) throws Exception {
//异步调用 有返回值
CompletableFuture completableFuture2 = CompletableFuture.supplyAsync(()->{
System.out.println(Thread.currentThread().getName()+" : CompletableFuture2");
//模拟异常
int i = 10/0;
return 1024;
});
completableFuture2.whenComplete((t,u)->{
// 返回值
System.out.println("------t="+t);
// 异常信息
System.out.println("------u="+u);
}).get();
}
}
```
输出:
> ForkJoinPool.commonPool-worker-1 : CompletableFuture2
> ------t=null
> ------u=java.util.concurrent.CompletionException: java.lang.ArithmeticException: / by zero
> Exception in thread "main" java.util.concurrent.ExecutionException: java.lang.ArithmeticException: / by zero
> at java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:357)
> at java.util.concurrent.CompletableFuture.get(CompletableFuture.java:1895)
> at com.atguigu.completable.CompletableFutureDemo.main(CompletableFutureDemo.java:28)
> Caused by: java.lang.ArithmeticException: / by zero
> at com.atguigu.completable.CompletableFutureDemo.lambda$main$1(CompletableFutureDemo.java:20)
> at java.util.concurrent.CompletableFuture$AsyncSupply.run(CompletableFuture.java:1590)
> at java.util.concurrent.CompletableFuture$AsyncSupply.exec(CompletableFuture.java:1582)
> at java.util.concurrent.ForkJoinTask.doExec(ForkJoinTask.java:289)
> at java.util.concurrent.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1056)
> at java.util.concurrent.ForkJoinPool.runWorker(ForkJoinPool.java:1692)
> at java.util.concurrent.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:157)