# JUC **Repository Path**: jxxb/JUC ## Basic Information - **Project Name**: JUC - **Description**: JUC 高并发编程学习笔记 - **Primary Language**: Java - **License**: MulanPSL-2.0 - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 4 - **Created**: 2024-10-10 - **Last Updated**: 2024-10-10 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README [TOC] ## 一、JUC概述 ### 1.1)JUC 简介 在 Java 中,线程部分是一个重点,本文的JUC也是关于线程的。 JUC 就是java.util .concurrent工具包的简称,这是一个处理线程的工具包,JDK 1.5开始出现的。 ![输入图片说明](img/01.jpg) ### 1.2)进程与线程 #### 1.2.1)进程 进程(Process) 是计算机中的程序关于某数据集合上的一次运行活动,是系统进行资源分配和调度的基本单位,是操作系统结构的基础。 在当代面向线程设计的计算机结构中: - 程序是指令、数据及其组织形式的描述; - 进程是程序的实体,是计算机中的程序关于某数据集合上的一次运行活动,是系统进行资源分配和调度的基本单位,是操作系统结构的基础; - 进程是线程的容器; #### 1.2.2)线程 线程(thread) 是操作系统能够进行运算调度的最小单位。它被包含在进程之中,是进程中的实际运作单位。 一条线程指的是进程中一个单一顺序的控制流,一个进程中可以并发多个线程,每条线程并行执行不同的任务。 #### 1.2.3)总结 进程:指在系统中正在运行的一个应用程序;程序一旦运行就是进程;进程是资源分配的最小单位; 线程:系统分配处理器时间资源的基本单元,或者说进程之内独立执行的一个单元执行流;线程是程序执行的最小单位。 ### 1.3)线程的状态 #### 1.3.1)线程状态枚举类 ```java public enum State { /** * Thread state for a thread which has not yet started. * 新建 */ NEW, /** * Thread state for a runnable thread. A thread in the runnable * state is executing in the Java virtual machine but it may * be waiting for other resources from the operating system * such as processor. * 准备就绪 */ RUNNABLE, /** * Thread state for a thread blocked waiting for a monitor lock. * A thread in the blocked state is waiting for a monitor lock * to enter a synchronized block/method or * reenter a synchronized block/method after calling * {@link Object#wait() Object.wait}. * 阻塞 */ BLOCKED, /** * Thread state for a waiting thread. * A thread is in the waiting state due to calling one of the * following methods: * * *

A thread in the waiting state is waiting for another thread to * perform a particular action. * * For example, a thread that has called Object.wait() * on an object is waiting for another thread to call * Object.notify() or Object.notifyAll() on * that object. A thread that has called Thread.join() * is waiting for a specified thread to terminate. * 不见不散 */ WAITING, /** * Thread state for a waiting thread with a specified waiting time. * A thread is in the timed waiting state due to calling one of * the following methods with a specified positive waiting time: *

* 过时不候 */ TIMED_WAITING, /** * Thread state for a terminated thread. * The thread has completed execution. * 终结 */ TERMINATED; } ``` #### 1.3.2)wait/sleep 的区别 1. sleep是Thread的静态方法,wait是Object的方法,任何对象实例都能调用; 2. sleep不会释放锁,它也不需要占用锁;wait会释放锁,但调用它的前提是当前线程占有锁(即代码要在synchronized中); 3. 它们都可以被interrupted方法中断。 ### 1.4)并发与并行 #### 1.4.1)串行模式 串行表示所有任务都按先后顺序进行; 串行意味着必须先装完一车柴才能运送这车柴,只有运送到了,才能卸下这车柴,并且只有完成了这整个三个步 骤,才能进行下一个步骤; 串行是一次只能取得一个任务,并执行这个任务。 #### 1.4.2)并行模式 并行意味着可以同时取得多个任务,并同时去执行所取得的这些任务; 并行模式相当于将长长的一条队列,划分成了多条短队列,所以并行缩短了任务队列的长度; 并行的效率从代码层次上强依赖于多进程/多线程代码,从硬件角度上则依赖于多核CPU。 #### 1.4.3)并发 并发(concurrent)指的是多个程序可以同时运行的现象,多进程可以同时运行或者多指令可以同时运行; 并发的重点在于它是一种现象, 描述的是多进程同时运行的现象。 但实际上,对于单核心CPU来说,同一时刻只能运行一个线程。所以,这里的"同时运行"表示的不是真的同一时刻有多个线程运行的现象,这是并行的概念,而是提供一种功能让用户看来多个程序同时运行起来了,但实际上这些程序中的进程不是一直霸占CPU的,而是执行一会停一会。 要解决大并发问题,通常是将大任务分解成多个小任务, 由于操作系统对进程的调度是随机的,所以切分成多个小任务后,可能会从任一小任务处执行。 这可能会出现一些现象: 1. 可能出现一个小任务执行了多次,还没开始下个任务的情况。这时一般会采用队列或类似的数据结构来存放各个小任务的成果; 2. 可能出现还没准备好第一步就执行第二步的可能;这时一般采用多路复用或异步的方式,比如只有准备好产生了事件通知才执行某个任务。 3. 可以多进程/多线程的方式并行执行这些小任务;也可以单进程/单线程执行这 些小任务,这时很可能要配合多路复用才能达到较高的效率 #### 1.4.4)总结 并发:同一时刻多个线程在访问同一个资源,多个线程对一个点 ; > 例子:春运抢票 电商秒杀... 并行:多项工作一起执行,之后再汇总 > 例子:泡方便面,电水壶烧水,一边撕调料倒入桶中 ### 1.5)管程 管程(monitor)是保证了同一时刻只有一个进程在管程内活动,即管程内定义的操作在同一时刻只被一个进程调用(由编译器实现),但是这样并不能保证进程以设计的顺序执行; JVM中同步是基于进入和退出管程(monitor)对象实现的,每个对象都会有一个管程 (monitor)对象,管程(monitor)会随着java对象一同创建和销毁; 执行线程首先要持有管程对象,然后才能执行方法,当方法完成之后会释放管程,方法在执行时候会持有管程,其他线程无法再获取同一个管程。 总结: 管程 ——》monitor监视器【操作系统中】 ——》锁【Java中】; 它是一种同步机制,保证统一时间段内,只有一个线程能够访问被保护的资源,其余线程不能进行访问; JVM的同步机制是基于进入和退出的过程进行操作的,进入和退出就是使用管程对象实现的,其对JVM临界区进行操作,在进入时加锁,退出时解锁【每一个对象都有monitor管程对象,其随着 Java对象一起创建和销毁】 ### 1.6)用户线程和守护线程 用户线程:平时用到的普通线程、自定义线程 ; 守护线程:运行在后台,是一种特殊的线程,比如垃圾回收 代码演示: ```java //演示用户线程 public class Main { public static void main(String[] args) { Thread aa = new Thread(() -> { // isDaemon() :表示当前线程是用户线程还是守护线程, ture --> 守护线程 false --> 用户线程 System.out.println("当前线程名:"+ Thread.currentThread().getName() + "::" + Thread.currentThread().isDaemon()); while (true) { } }, "aa"); // 运行线程 aa.start(); System.out.println("主线程名:"+ Thread.currentThread().getName() + " over"); } } ``` 输出:【当主线程结束后,用户线程还在运行,则JVM存活 】 > 主线程名:main over > 当前线程名:aa::false ```java //演示守护线程 public class Main { public static void main(String[] args) { Thread aa = new Thread(() -> { // isDaemon() :表示当前线程是用户线程还是守护线程, ture --> 守护线程 false --> 用户线程 System.out.println("当前线程名:"+ Thread.currentThread().getName() + "::" + Thread.currentThread().isDaemon()); while (true) { } }, "aa"); //设置守护线程 aa.setDaemon(true); // 运行线程 aa.start(); System.out.println("主线程名:"+ Thread.currentThread().getName() + " over"); } } ``` 输出:【如果没有用户线程,都是守护线程,则JVM结束 】 > 主线程名:main over > 当前线程名:aa::true ## 二、Lock 接口 ### 2.1)Synchronized #### 2.1.1)Synchronized 关键字简介 synchronized 是 Java 中的关键字,是一种同步锁,它修饰的对象有以下几种: 1. 修饰一个代码块,被修饰的代码块称为同步语句块,其作用的范围是大括号{} 括起来的代码,作用的对象是调用这个代码块的对象; 2. 修饰一个方法,被修饰的方法称为同步方法,其作用的范围是整个方法,作用的对象是调用这个方法的对象; > 虽然可以使用synchronized来定义方法,但synchronized并不属于方法定义的一部分,因此,synchronized关键字不能被继承。 > > 如果在父类中的某个方法使用了synchronized关键字,而在子类中复写了这个方法,在子类中的这个方法默认情况下并不是同步的,而必须显式地在子类的这个方法中加上 synchronized关键字才可以。还可以在子类方法中调用父类中相应的方法,这样虽然子类中的方法不是同步的,但子类调用了父类的同步方法,因此, 子类的方法也就相当于同步了。 3. 修改一个静态的方法,其作用的范围是整个静态方法,作用的对象是这个类的所有对象; 4. 修改一个类,其作用的范围是synchronized后面括号括起来的部分,作用主的对象是这个类的所有对象。 #### 2.1.2)使用Synchronized 关键字示例 使用Synchronized 关键字实现卖票的示例: ```java package com.study.sync; //第一步 创建资源类,定义属性和和操作方法 class Ticket { //票数 private int number = 30; //操作方法:卖票 public synchronized void sale() { //判断:是否有票 if (number > 0) { System.out.println(Thread.currentThread().getName() + " : 卖出:" + (number--) + " 剩下:" + number); } } } public class SaleTicket { //第二步 创建多个线程,调用资源类的操作方法 public static void main(String[] args) { //创建Ticket对象 Ticket ticket = new Ticket(); //创建三个线程 // 创建线程 AA new Thread(new Runnable() { @Override public void run() { //调用卖票方法 for (int i = 0; i < 40; i++) { ticket.sale(); } } }, "AA").start(); // 创建线程 BB new Thread(new Runnable() { @Override public void run() { //调用卖票方法 for (int i = 0; i < 40; i++) { ticket.sale(); } } }, "BB").start(); // 创建线程 CC new Thread(new Runnable() { @Override public void run() { //调用卖票方法 for (int i = 0; i < 40; i++) { ticket.sale(); } } }, "CC").start(); } } ``` 输出: > com.study.sync.SaleTicket > 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 > BB : 卖出:15 剩下:14 > BB : 卖出:14 剩下:13 > BB : 卖出:13 剩下:12 > BB : 卖出:12 剩下:11 > BB : 卖出:11 剩下:10 > BB : 卖出:10 剩下:9 > BB : 卖出:9 剩下:8 > BB : 卖出:8 剩下:7 > BB : 卖出:7 剩下:6 > BB : 卖出:6 剩下:5 > BB : 卖出:5 剩下:4 > BB : 卖出:4 剩下:3 > BB : 卖出:3 剩下:2 > BB : 卖出:2 剩下:1 > BB : 卖出:1 剩下:0 #### 2.1.3)Synchronized 关键字说明 如果一个代码块被synchronized修饰了,当一个线程获取了对应的锁,并执行该代码块时,其他线程便只能一直等待,等待获取锁的线程释放锁,而这里获取锁的线程释放锁只会有两种情况: 1. 获取锁的线程执行完了该代码块,然后线程释放对锁的占有; 2. 线程执行发生异常,此时JVM会让线程自动释放锁。 那么如果这个获取锁的线程由于要等待IO或者其他原因(比如调用sleep 方法)被阻塞了,但是又没有释放锁,其他线程便只能干巴巴地等待,试想一 下,这多么影响程序执行效率。 因此就需要有一种机制可以不让等待的线程一直无期限地等待下去(比如只等 待一定的时间或者能够响应中断),通过Lock就可以办到。 ### 2.2)Lock Lock锁实现提供了比使用同步方法和语句可以获得的更广泛的锁操作;Lock允许更灵活的结构,能具有非常不同的属性,并且能支持多个关联的条件对象;Lock提供了比synchronized更多的功能。 #### 2.2.1)使用Lock 接口示例 使用Lock 接口实现卖票的示例: ```java package com.study.lock; import java.util.concurrent.locks.ReentrantLock; //第一步 创建资源类,定义属性和和操作方法 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(); } } ``` 输出: > com.study.sync.SaleTicket > AA :卖出30 剩余:29 > AA :卖出29 剩余:28 > AA :卖出28 剩余:27 > AA :卖出27 剩余:26 > AA :卖出26 剩余:25 > BB :卖出25 剩余:24 > AA :卖出24 剩余:23 > CC :卖出23 剩余:22 > BB :卖出22 剩余:21 > AA :卖出21 剩余:20 > CC :卖出20 剩余:19 > BB :卖出19 剩余:18 > AA :卖出18 剩余:17 > BB :卖出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 #### 2.2.2)Lock和synchronized的不同 Lock和synchronized有以下几点不同: 1. Lock是一个接口,而synchronized是Java中的关键字,synchronized是内置的语言实现; 2. synchronized在发生异常时,会自动释放线程占有的锁,因此不会导致死锁现象发生;而Lock在发生异常时,如果没有主动通过unLock()去释放锁,则很可能造成死锁现象,因此使用Lock时需要在finally块中释放锁; 3. Lock可以让等待锁的线程响应中断,而synchronized却不行,使用 synchronized时,等待的线程会一直等待下去,不能够响应中断; 4. 通过Lock可以知道有没有成功获取锁,而synchronized却无法办到; 5. Lock可以提高多个线程进行读操作的效率,在性能上来说,如果竞争资源不激烈,两者的性能是差不多的,而当竞争资源非常激烈时(即有大量线程同时竞争),此时Lock的性能要远远优于synchronized。 ## 三、线程间通信 线程间通信的模型有两种:共享内存和消息传递,以下方式都是基本这两种模型来实现的。 ### 3.1)线程间通信案例 两个线程,一个线程对当前数值加1,另一个线程对当前数值减1,要求用线程间通信 #### 3.1.1)synchronized 方案 ```java package com.study.sync; //第一步 创建资源类,定义属性和操作方法 class Share { //初始值 private int number = 0; //+1的方法 public synchronized void incr() throws InterruptedException { //第二步 判断 干活 通知 //判断:number值是否是0,如果不是0,等待 // if (number != 0) { --> 会产生虚假唤醒问题,即只对第一次进入的条件进行了判断,使用while判断可以避免虚假唤醒 while (number != 0) { //在哪里睡,就在哪里醒 this.wait(); } //干活:如果number值是0,就+1操作 number++; System.out.println("当前线程名:" + Thread.currentThread().getName() + " :: " + number); //通知:通知其他线程 this.notifyAll(); } //-1的方法 public synchronized void decr() throws InterruptedException { //判断:number值是否是1,如果不是1,等待 // if (number != 1) { --> 会产生虚假唤醒问题,即只对第一次进入的条件进行了判断,使用while判断可以避免虚假唤醒 while (number != 1) { this.wait(); } //干活:如果number值是1,就-1操作 number--; System.out.println("当前线程名:" + Thread.currentThread().getName() + " :: " + number); //通知:通知其他线程 this.notifyAll(); } } public class ThreadDemo1 { //第三步 创建多个线程,调用资源类的操作方法 public static void main(String[] args) { Share share = new Share(); //创建线程 // 创建线程 AA new Thread(() -> { for (int i = 1; i <= 10; i++) { try { share.incr(); //+1 } catch (InterruptedException e) { e.printStackTrace(); } } }, "AA").start(); // 创建线程 BB new Thread(() -> { for (int i = 1; i <= 10; i++) { try { share.decr(); //-1 } catch (InterruptedException e) { e.printStackTrace(); } } }, "BB").start(); // 创建线程 CC new Thread(() -> { for (int i = 1; i <= 10; i++) { try { share.incr(); //+1 } catch (InterruptedException e) { e.printStackTrace(); } } }, "CC").start(); // 创建线程 DD new Thread(() -> { for (int i = 1; i <= 10; i++) { try { share.decr(); //-1 } catch (InterruptedException e) { e.printStackTrace(); } } }, "DD").start(); } } ``` 输出: > com.study.sync.ThreadDemo1 > 当前线程名:AA :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:BB :: 0 > 当前线程名:CC :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:BB :: 0 > 当前线程名:CC :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:DD :: 0 > 当前线程名:AA :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:DD :: 0 #### 3.1.2)Lock 方案 ```java package com.study.lock; import java.util.concurrent.locks.Condition; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; //第一步 创建资源类,定义属性和操作方法 class Share { private int number = 0; //创建Lock private Lock lock = new ReentrantLock(); private Condition condition = lock.newCondition(); //+1 方法 public void incr() throws InterruptedException { //上锁 lock.lock(); try { //判断:number值是否是0,如果不是0,等待 // if (number != 0) { --> 会产生虚假唤醒问题,即只对第一次进入的条件进行了判断,使用while判断可以避免虚假唤醒 while (number != 0) { condition.await(); } //干活:如果number值是0,就+1操作 number++; System.out.println("当前线程名:" + Thread.currentThread().getName() + " :: " + number); //通知:通知其他线程 condition.signalAll(); } finally { //解锁 lock.unlock(); } } //-1 方法 public void decr() throws InterruptedException { //上锁 lock.lock(); try { //判断:number值是否是1,如果不是1,等待 // if (number != 1) { --> 会产生虚假唤醒问题,即只对第一次进入的条件进行了判断,使用while判断可以避免虚假唤醒 while (number != 1) { condition.await(); } //干活:如果number值是1,就-1操作 number--; System.out.println("当前线程名:" + Thread.currentThread().getName() + " :: " + number); //通知:通知其他线程 condition.signalAll(); } finally { //解锁 lock.unlock(); } } } public class ThreadDemo2 { public static void main(String[] args) { Share share = new Share(); // 创建线程 AA new Thread(() -> { for (int i = 1; i <= 10; i++) { try { // 调用 +1 方法 share.incr(); } catch (InterruptedException e) { e.printStackTrace(); } } }, "AA").start(); // 创建线程 BB new Thread(() -> { for (int i = 1; i <= 10; i++) { try { // 调用 -1 方法 share.decr(); } catch (InterruptedException e) { e.printStackTrace(); } } }, "BB").start(); // 创建线程 CC new Thread(() -> { for (int i = 1; i <= 10; i++) { try { // 调用 +1 方法 share.incr(); } catch (InterruptedException e) { e.printStackTrace(); } } }, "CC").start(); // 创建线程 DD new Thread(() -> { for (int i = 1; i <= 10; i++) { try { // 调用 -1 方法 share.decr(); } catch (InterruptedException e) { e.printStackTrace(); } } }, "DD").start(); } } ``` 输出: > com.study.lock.ThreadDemo2 > 当前线程名:AA :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:BB :: 0 > 当前线程名:AA :: 1 > 当前线程名:BB :: 0 > 当前线程名:CC :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:DD :: 0 > 当前线程名:CC :: 1 > 当前线程名:DD :: 0 ### 3.2)多线程编程步骤 1. 创建资源类,在资源类中创建属性和操作方法 2. 在资源类中操作方法 > 1. 判断 > 2. 干活 > 3. 通知 3. 创建多个线程,调用资源类的操作方法 4. 防止虚假唤醒问题 ## 四、线程间定制化通信 ### 4.1)线程间定制化通信案例 A线程打印2次AA,B线程打印4次BB,C线程打印6次CC,按照此顺序循环3轮 ### 4.2)线程间定制化通信案例代码实现 #### 4.2.1)Lock接口 ```java public interface Lock { void lock(); void lockInterruptibly() throws InterruptedException; boolean tryLock(); boolean tryLock(long time, TimeUnit unit) throws InterruptedException; void unlock(); Condition newCondition(); } ``` ##### 4.2.1.1)lock() lock()方法是平常使用得最多的一个方法,就是用来获取锁。 如果锁已被其他线程获取,则进行等待;采用Lock,必须主动去释放锁,并且在发生异常时,不会自动释放锁。 因此一般来说,使用Lock必须在try{}catch{}块中进行,并且将释放锁的操作放在finally块中进行,以保证锁一定被被释放,防止死锁的发生。 通常使用Lock 来进行同步的话,是以下面这种形式去使用的: ```java Lock lock = new ReentrantLock(); //上锁 lock.lock(); try{ //处理任务 }catch(Exception ex){ ... }finally{ //释放锁 lock.unlock(); } ``` ##### 4.2.1.2)newCondition 关键字synchronized与wait()/notify()这两个方法一起使用可以实现等待/通知模式,Lock锁的newContition()方法返回Condition对象,Condition类也可以实现等待/通知模式。 用notify()通知时,JVM会随机唤醒某个等待的线程,使用Condition类可以进行选择性通知,Condition比较常用的两个方法: 1. await()会使当前线程等待,同时会释放锁,当其他线程调用signal()时,线程会重 新获得锁并继续执行; 2. signal()用于唤醒一个等待的线程。 > 注意:在调用Condition的await()/signal()方法前,也需要线程持有相关的Lock锁,调用await()后线程会释放这个锁,在singal()调用后会从当前 Condition对象的等待队列中,唤醒一个线程,唤醒的线程尝试获得锁,一旦获得锁成功就继续执行。 #### 4.2.2)代码实现 流程分析如下图: ![输入图片说明](img/02.jpg) 代码如下: ```java package com.study.lock; import java.util.concurrent.locks.Condition; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; //第一步 创建资源类 class ShareResource { //定义标志位 // 1 AA 2 BB 3 CC private int flag = 1; //创建Lock锁 private Lock lock = new ReentrantLock(); //创建三个condition private Condition c1 = lock.newCondition(); private Condition c2 = lock.newCondition(); private Condition c3 = lock.newCondition(); //打印2次,参数第几轮 public void print2(int loop) throws InterruptedException { //上锁 lock.lock(); try { //判断:flag值是否是1,如果不是1,等待 while (flag != 1) { //等待 c1.await(); } //干活:如果flag值是1,就输出信息 for (int i = 1; i <= 2; i++) { System.out.println("当前线程名:" + Thread.currentThread().getName() + " :: " + i + " :轮数:" + loop); } //通知:通知BB线程 //修改标志位为 2 -> BB flag = 2; //通知BB线程 c2.signal(); } finally { //释放锁 lock.unlock(); } } //打印4次,参数第几轮 public void print4(int loop) throws InterruptedException { lock.lock(); try { //判断:flag值是否是2,如果不是2,等待 while (flag != 2) { c2.await(); } //干活:如果flag值是2,就输出信息 for (int i = 1; i <= 4; i++) { System.out.println("当前线程名:" + Thread.currentThread().getName() + " :: " + i + " :轮数:" + loop); } //通知:通知CC线程 //修改标志位为 3 -> CC flag = 3; //通知CC线程 c3.signal(); } finally { lock.unlock(); } } //打印6次,参数第几轮 public void print6(int loop) throws InterruptedException { lock.lock(); try { //判断:flag值是否是3,如果不是3,等待 while (flag != 3) { c3.await(); } //干活:如果flag值是3,就输出信息 for (int i = 1; i <= 6; i++) { System.out.println("当前线程名:" + Thread.currentThread().getName() + " :: " + i + " :轮数:" + loop); } //通知:通知AA线程 //修改标志位为 1 -> AA flag = 1; //通知AA线程 c1.signal(); } finally { lock.unlock(); } } } public class ThreadDemo3 { public static void main(String[] args) { ShareResource shareResource = new ShareResource(); // 创建线程 AA new Thread(() -> { // 循环调用 3 次 for (int i = 1; i <= 3; i++) { try { shareResource.print2(i); } catch (InterruptedException e) { e.printStackTrace(); } } }, "AA").start(); // 创建线程 BB new Thread(() -> { // 循环调用 3 次 for (int i = 1; i <= 3; i++) { try { shareResource.print4(i); } catch (InterruptedException e) { e.printStackTrace(); } } }, "BB").start(); // 创建线程 CC new Thread(() -> { // 循环调用 3 次 for (int i = 1; i <= 3; i++) { try { shareResource.print6(i); } catch (InterruptedException e) { e.printStackTrace(); } } }, "CC").start(); } } ``` 输出: > com.study.lock.ThreadDemo3 > 当前线程名:AA :: 1 :轮数:1 > 当前线程名:AA :: 2 :轮数:1 > 当前线程名:BB :: 1 :轮数:1 > 当前线程名:BB :: 2 :轮数:1 > 当前线程名:BB :: 3 :轮数:1 > 当前线程名:BB :: 4 :轮数:1 > 当前线程名:CC :: 1 :轮数:1 > 当前线程名:CC :: 2 :轮数:1 > 当前线程名:CC :: 3 :轮数:1 > 当前线程名:CC :: 4 :轮数:1 > 当前线程名:CC :: 5 :轮数:1 > 当前线程名:CC :: 6 :轮数:1 > 当前线程名:AA :: 1 :轮数:2 > 当前线程名:AA :: 2 :轮数:2 > 当前线程名:BB :: 1 :轮数:2 > 当前线程名:BB :: 2 :轮数:2 > 当前线程名:BB :: 3 :轮数:2 > 当前线程名:BB :: 4 :轮数:2 > 当前线程名:CC :: 1 :轮数:2 > 当前线程名:CC :: 2 :轮数:2 > 当前线程名:CC :: 3 :轮数:2 > 当前线程名:CC :: 4 :轮数:2 > 当前线程名:CC :: 5 :轮数:2 > 当前线程名:CC :: 6 :轮数:2 > 当前线程名:AA :: 1 :轮数:3 > 当前线程名:AA :: 2 :轮数:3 > 当前线程名:BB :: 1 :轮数:3 > 当前线程名:BB :: 2 :轮数:3 > 当前线程名:BB :: 3 :轮数:3 > 当前线程名:BB :: 4 :轮数:3 > 当前线程名:CC :: 1 :轮数:3 > 当前线程名:CC :: 2 :轮数:3 > 当前线程名:CC :: 3 :轮数:3 > 当前线程名:CC :: 4 :轮数:3 > 当前线程名:CC :: 5 :轮数:3 > 当前线程名:CC :: 6 :轮数:3 ## 五、集合的线程安全 ### 5.1)ArrayList集合线程 #### 5.1.1)ArrayList集合线程不安全案例展示 代码如下: ```java public class ArrayListErrorDemo { public static void main(String[] args) { //创建ArrayList集合 List list = new ArrayList<>(); // 同时创建30个线程同时向集合中添加和获取元素 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(); } } } ``` 输出: > Exception in thread "27" Exception in thread "17" java.util.ConcurrentModificationException > at java.util.ArrayList$Itr.checkForComodification(ArrayList.java:909) > at java.util.ArrayList$Itr.next(ArrayList.java:859) > 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.ArrayListErrorDemo.lambda$main$0(ArrayListErrorDemo.java:24) > at java.lang.Thread.run(Thread.java:748) > java.util.ConcurrentModificationException > at java.util.ArrayList$Itr.checkForComodification(ArrayList.java:909) > at java.util.ArrayList$Itr.next(ArrayList.java:859) > 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.ArrayListErrorDemo.lambda$main$0(ArrayListErrorDemo.java:24) > at java.lang.Thread.run(Thread.java:748) 发现报错,java.util.ConcurrentModificationException ——》并发修改异常 > 问题: 为什么会出现并发修改异常? 查看ArrayList的add方法源码 ,如下: ```java /** * Appends the specified element to the end of this list. * * @param e element to be appended to this list * @return true (as specified by {@link Collection#add}) */ public boolean add(E e) { ensureCapacityInternal(size + 1); // Increments modCount!! elementData[size++] = e; return true; } ``` add() 方法没有添加 synchronized 关键字,即没有对多线程情况进行处理 #### 5.1.2)ArrayList集合线程不安全解决方案 ##### 5.1.2.1)解决方案1-Vector Vector 是矢量队列,它是JDK1.0版本添加的类;继承于AbstractList,实现了List, RandomAccess, Cloneable 这些接口。 它是一个队列,支持相关的添加、删除、修改、遍历等功能。 Vector 实现了RandmoAccess接口,提供了随机访问功能。 可以通过元素的序号快速获取元素对象;这就是快速随机访问。 代码如下: ```java // Vector解决ArrayList集合线程不安全问题 public class VectorDemo { public static void main(String[] args) { //创建Vector集合 List list = new Vector<>(); 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(); } } } ``` 输出: > [2e6f3116, 6fe43ea3, 02fca385, b678c6a3, dd308329, ab1ed47c] > [2e6f3116, 6fe43ea3, 02fca385, b678c6a3, dd308329, ab1ed47c, 6940b25f, de6b1ff7, 762f2d14, c88c6146, 8eb58040, 0d628c6c, b0c3849e, 4bed8696, 11c2ac32, ea177a6a, 6c0e1f8f, 50040ab8, 7e3a1e0d, 9162aee6] > ... > [2e6f3116, 6fe43ea3, 02fca385, b678c6a3, dd308329, ab1ed47c, 6940b25f, de6b1ff7, 762f2d14, c88c6146, 8eb58040, 0d628c6c, b0c3849e, 4bed8696, 11c2ac32, ea177a6a, 6c0e1f8f, 50040ab8, 7e3a1e0d, 9162aee6, 952f0041] 和ArrayList不同,Vector中的操作是线程安全的,为什么不会出现并发修改异常? 查看Vector的add()方法源码 ,如下: ```java /** * Appends the specified element to the end of this Vector. * * @param e element to be appended to this Vector * @return {@code true} (as specified by {@link Collection#add}) * @since 1.2 */ public synchronized boolean add(E e) { modCount++; ensureCapacityHelper(elementCount + 1); elementData[elementCount++] = e; return true; } ``` Vector的add()方法添加了 synchronized关键字,所以不会出现并发修改异常 ##### 5.1.2.2)解决方案2-Collections Collections提供了方法synchronizedList保证list是同步线程安全的,代码如下: ```java // Collections解决ArrayList集合线程不安全问题 public class CollectionsDemo { public static void main(String[] args) { //Collections解决 List 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的核心技术就是 **写时复制技术**,流程如下图: ![输入图片说明](img/03.jpg) 读集合数据时是并发读,写集合数据时是独立写,这样就避免了在写数据时导致的并发修改异常问题 代码如下: ```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)什么是死锁 死锁:两个或者两个以上进程在执行过程中,因为争夺资源而造成一种互相等待的现象,如果没有外力干涉,它们就无法在执行下去,会卡在这里,如下图所示: ![输入图片说明](img/04.jpg) 线程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)读锁 读锁:是**共享锁**,会发生死锁的情况,如下图所示: ![输入图片说明](img/09.jpg) 上图中:二者互相等待产生死锁 - 线程1修改数据的时候,需要等待线程2读取数据之后; - 线程2修改数据的时候,同样需要等待线程1读取数据之后; 线程进入读锁的前提条件: 1. 没有其他线程的写锁 2. 没有写请求, 或者有写请求,但调用线程和持有锁的线程是同一个(可重入锁) #### 9.1.2)写锁 写锁:是**独占锁**,会发生死锁的情况,如下图所示: ![输入图片说明](img/10.jpg) 上图中:二者互相等待产生死锁 - 线程1对第一条数据进行写操作时候,同时对第二条数据进行修改操作【同一个线程中操作多条数据】需要等待线程2对第二条数据进行写操作之后才能够进行; - 线程2对第二条数据进行写操作时候,同时对第一条数据进行修改操作【同一个线程中操作多条数据】需要等待线程1对第一条数据进行写操作之后才能够进行;; 线程进入写锁的前提条件: 1. 没有其他线程的读锁 2. 没有其他线程的写锁 #### 9.1.3)读写锁特性 读写锁有以下三个重要的特性: (1)公平选择性:支持非公平(默认)和公平的锁获取方式,吞吐量还是非公平优于公平; (2)重进入:读锁和写锁都支持线程重进入; (3)锁降级:遵循获取写锁、获取读锁再释放写锁的次序,写锁能够降级成为读锁。 #### 9.1.4)读写锁的演变 读写锁:一个资源可以被多个读线程访问,或者可以被一个写线程访问,但是不能同时存在读写线程,读写是互斥的,读读是共享的; 读写锁的演变过程如下图: ![输入图片说明](img/11.jpg) ### 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中的锁降级流程如下图: ![输入图片说明](img/13.jpg) 将写锁降级为读锁代码演示: ```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很好的解决了多线程中,如何高效安全“传输”数据的问题;通过这些高效并且线程安全的队列类,为我们快速搭建高质量的多线程程序带来极大的便利。 阻塞队列,顾名思义,首先它是一个队列, 通过一个共享的队列,可以使得数据由队列的一端输入,从另外一端输出,如下图所示: ![输入图片说明](img/14.jpg) 阻塞情况如下: - 当队列是空的,从队列中获取元素的操作将会被阻塞; - 当队列是满的,从队列中添加元素的操作将会被阻塞; - 试图从空的队列中获取元素的线程将会被阻塞,直到其他线程往空的队列插入新的元素; - 试图向已满的队列中添加新元素的线程将会被阻塞,直到其他线程从队列中移除一个或多个元素或者完全清空,使队列变得空闲起来并后续新增 #### 10.1.1)常用的队列 常用的队列主要有以下两种: - 先进先出(FIFO):先插入的队列的元素也最先出队列,类似于排队的功能,从某种程度上来说这种队列也体现了一种公平性 (队列) - 后进先出(LIFO):后插入队列的元素最先出队列,这种队列优先处理最近发生的事件(栈) #### 10.1.2)为什么需要BlockingQueue 在concurrent包发布以前,在多线程环境下,开发者都必须去自己控制阻塞线程的细节,尤其还要兼顾效率和线程安全,而这会给我们的程序带来不小的复杂度; 使用BlockingQueue 之后,开发者不需要关心什么时候需要阻塞线程,什么时候需要唤醒线程,因为这一切 都由BlockingQueue一手包办了。 #### 10.1.3)BlockingQueue的实际应用 多线程环境中,通过队列可以很容易实现数据共享,比如经典的“生产者”和 “消费者”模型中,通过队列可以很便利地实现两者之间的数据共享。 假设有若干生产者线程,另外又有若干个消费者线程。如果生产者线程需要把准备好的数据共享给消费者线程,利用队列的方式来传递数据,就可以很方便地解决它们之间的数据共享问题。 > 但如果生产者和消费者在某个时间段内,万一发生数据处理速度不匹配的情况如何解决呢? > > 1. 理想情况下,如果生产者产出数据的速度大于消费者消费的速度,并且当生产出来的数据累积到一定程度的时候,那么生产者必须暂停等待一下(阻塞生产者线程),以便等待消费者线程把累积的 数据处理完毕,反之亦然; > 2. 当队列中没有数据的情况下,消费者端的所有线程都会被自动阻塞(挂起), 直到有数据放入队列; > 3. 当队列中填满数据的情况下,生产者端的所有线程都会被自动阻塞(挂起), 直到队列中有空的位置,线程被自动唤醒 ### 10.2)BlockingQueue 核心方法 BlockingQueue的核心方法,如下图所示: ![输入图片说明](img/17.jpg) #### 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这几个类 ,其关系如下图: ![输入图片说明](img/18.jpg) ### 11.2)线程池参数说明 常用参数: - corePoolSize:线程池的常驻线程数量(核心) - maximumPoolSize:线程池能容纳的最大线程数 - keepAliveTime空闲:线程存活时间 - unit :线程存活的时间单位 - BlockingQueue workQueue :存放提交但未执行任务的队列 【阻塞队列】 - threadFactory: 创建线程的工厂类【线程工厂】 - handler :等待队列满后的拒绝策略 线程池中,有三个重要的参数,决定影响了拒绝策略: corePoolSize - 核心线程数,也即最小的线程数; workQueue - 阻塞队列; maximumPoolSize - 最大线程数 当提交任务数大于 corePoolSize 的时候,会优先将任务放到 workQueue 阻塞队列中;当阻塞队列饱和后,会扩充线程池中线程数,直到达到 maximumPoolSize 最大线程数配置。此时再多余的任务,则会触发线程池的拒绝策略。 > 总结:当提交的任务数大于(workQueue.size() + maximumPoolSize ),就会触发线程池的拒绝策略 ### 11.3)线程池底层工作原理 线程池底层工作原理如下图所示: ![输入图片说明](img/19.jpg) 线程池底层工作流程: 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内置的拒绝策略如下图所示: ![输入图片说明](img/20.jpg) 说明: 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.的方式手动创建线程池,如下图 : ![输入图片说明](img/21.jpg) 所以实际生产一般自己通过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. 执行任务并合并结果:分割的子任务分别放到双端队列里,然后几个启动线程分别从双端队列里获取任务执行;子任务执行完的结果都放在另外一个队列里, 启动一个线程从队列里取数据,然后合并这些数据。 如下图所示: ![输入图片说明](img/26.jpg) ### 12.3)Fork/Join框架的实现原理 在Java的Fork/Join框架中,ForkJoinPool【分支合并池】由ForkJoinTask数组和ForkJoinWorkerThread数组组成;使用它们完成任务分割和合并结果操作 : 1. ForkJoinTask数组:负责将存放以及将程序提交给ForkJoinPool;需要通过ForkJoinPool来执行;要使用 Fork/Join 框架,首先需要创建一个 ForkJoin 任务,该类提供了在任务中执行fork和join的机制;通常情况下不需要直接集成ForkJoinTask类,只需要继承它的子类; > Fork/Join框架提供了两个子类: > > - RecursiveAction:用于没有返回结果的任务; > > - RecursiveTask:用于有返回结果的任务,继承后可以实现递归(自己调自己)调用的任务 > > ![输入图片说明](img/23.jpg) 2. ForkJoinWorkerThread数组:负责执行这些任务; ForkJoinPool【分支合并池】类似于 线程池,如下图: ![输入图片说明](img/22.jpg) #### 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)