预见猿份
主题
首页面试题在线工具关于我们老苗一对一私教学员评价
实战项目
项目前置基础创新WMS项目Java微服务框架与实战云岚到家项目闪聚支付项目学成在线项目青橙电商项目JVM原理与实战调优分布式事务专题Java高频面试题MySQL从入门到精通Java数据结构与算法Java 并发编程(JUC)实战课程老苗一对一私教学员评价blog
blog
  • Java 并发编程(JUC)实战课程

    • 课程介绍
    • 01 多线程入门
    • 02 线程的常用方法
    • 03 线程池
    • 04 synchronized教程
    • 05 ReentrantLock教程
    • 06 读写锁
    • 07 JMM教程
    • 08 CAS教程
    • 09 wait&notify
    • 10 常用的并发集合
    • 11 异步任务控制
    • 12 常用并发工具







----- 到底线了 -----

第9章 线程协作-wait&notify ​

1.1 wait&notify实现线程协作 ​

1.1.1 线程协作案例需求分析 ​

在多线程开发中为了充分利用系统资源,将一个任务拆分为多个小任务,这样的话一个任务需要多线程协作来完成任务,比如线程A要处理的数据由线程B来提供;线程A处理完成去通知 线程B继续处理,线程与线程之间是协同完成任务而不是独立完成任务。

预见猿份视频处理程序由多线程协作完成视频处理,一个原始视频需要经过转码和切割,第一步先将视频统一转成mp4格式,第二步再将mp4视频切割成m3u8格式,一部分线程负责完成第一步,一部分线程负责完成第二步。

为了便于学习暂时将上边的需求简化为两个线程之间的协作:线程A完成视频 转mp4,线程B完成视频转m3u8,线程B需要等线程A先完成转mp4的处理然后将mp4转成m3u8,最终完成一个视频的处理。

如下图:

配图1

需求如下:

1、线程A负责转MP4,线程B负责转m3u8。

2、线程A先工作,线程B后工作。

3、线程A完成任务需要通知线程B去工作。

分析如下:

1、线程A和线程B的工作有先后

线程A负责转MP4并且先去执行任务,线程A完成任务后由线程B开始执行转M3u8任务,如果线程A和线程B同时启动,则可以通过一个共享变量控制视频的工作状态,比如设置如下共享变量:

mp4完成状态:false 表示未完成mp4转换,true表示完成mp4转换。

2、线程A完成任务需要通知线程B

根据当前所学知识线程之间的通信方法如下:

1)中断请求

线程A向线程B发起中断请求,线程A需要通过Thread对象向线程B发起中断请求,此方法的问题在于线程A和线程B的交互存在耦合,不易扩展,另外对于使用线程池来说需要从池子里拿到线程B的Thread对象,过程略显复杂。

2)阻塞队列

线程A将视频完成视频放入阻塞队列,线程B从阻塞队列获取视频进行处理,阻塞队列将线程A和线程B解耦合,阻塞队列在多线程通信中使用广泛,但是针对本需求线程A通过队列通知线程B有些笨重。

3)共享变量

可以设置一个共享变量存储视频转换状态,线程A完成转MP4工作后设置工作状态,线程B读取工作状态可知线程A是否完成工作。此方法较简单,本次测试使用此方法。

1.1.2 使用共享变量实现线程协作 ​

根据需求分析下边通过共享变量完成线程之间的协作。

代码如下:

1、首先创建视频类

这里创建一个内部类,类中设置两个状态变量:mp4完成状态、视频转换是否结束。

mp4完成状态:false 表示未完成mp4转换,true表示完成mp4转换。

视频转换是否结束:false表示视频转换未完成,true表示视频转换完成。

java
package com.yjoffer.javase.thread.waitnotify;

import com.yjoffer.javase.config.Logger;

import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;

//视频类
public class PbVideo{

		// 视频名称
		String name;

		// mp4完成状态
		boolean mp4finished = false;

		//m3u8完成状态
		boolean m3u8finished = false;

		//视频转换结束
		boolean finished = false;


		public PbVideo(String name) {
			this.name = name;
		}

		public String getName() {
			return name;
		}

		public void setName(String name) {
			this.name = name;
		}

		public boolean isMp4finished() {
			return mp4finished;
		}

		public void setMp4finished(boolean mp4finished) {
			this.mp4finished = mp4finished;
		}

		public boolean isFinished() {
			return finished;
		}

		public void setFinished(boolean finished) {
			this.finished = finished;
		}

		public boolean isM3u8finished() {
			return m3u8finished;
		}

		public void setM3u8finished(boolean m3u8finished) {
			this.m3u8finished = m3u8finished;
		}


	}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66

创建两个任务类,一个是转MP4、一个是转m3u8。

java
/**
	 * 转mp4
	 */
	 class Mp4RunnableV1 implements Runnable{
		//视频
		private PbVideo pbVideo;

		public Mp4RunnableV1(PbVideo pbVideo){
			this.pbVideo = pbVideo;
		}
		public void run() {
			if(pbVideo.isMp4finished()){
				return ;
			}
			logger.debug("开始转mp4...");
			try {
				//模拟视频转换过程
				TimeUnit.SECONDS.sleep(3);
			} catch (InterruptedException e) {
				e.printStackTrace();
			}
			pbVideo.setMp4finished(true);
			logger.debug("转mp4结束! ");

		}
	}
	/**
	 * 转M3u8
	 */
	 class M3u8RunnableV1 implements Runnable{

		private PbVideo pbVideo;
		public M3u8RunnableV1(PbVideo pbVideo){
			this.pbVideo = pbVideo;
		}
		public void run() {
			if(pbVideo.isM3u8finished()){
				return ;
			}
			while (!pbVideo.isMp4finished()){
				logger.debug("M3u8Runnable暂时休息");
				try {
					TimeUnit.SECONDS.sleep(1);
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
			}
			logger.debug("开始转m3u8...");
			try {
				TimeUnit.SECONDS.sleep(3);
			} catch (InterruptedException e) {
				e.printStackTrace();
			}
			pbVideo.setM3u8finished(true);
			pbVideo.setFinished(true);
			logger.debug("转m3u8结束! ");
		}
	}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58

最后编写测试方法:

java
/**
 * 线程协作测试
 * @author 预见猿份(www.yjoffer.com)
 *
 */
public class ThreadWaitnotifyTest {

	private static Logger logger = Logger.getLogger(ThreadWaitnotifyTest.class);


	...

	public static void test_videoconvert1(){
		ExecutorService threadPool = Executors.newCachedThreadPool();
		//视频对象
		PbVideo pbVideo = new PbVideo("预见猿份Java自学教程");
		//启动一个线程开始转mp4
		threadPool.execute(new Mp4RunnableV1(pbVideo));
		//启动一个线程开始转m3u8
		threadPool.execute(new M3u8RunnableV1(pbVideo));
		shutdown(threadPool);
	}
	....
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23

编写 main方法运行测试方法test_videoconvert(),输出 如下:

java
10:52:25.234 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - 开始转mp4...
10:52:25.235 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - M3u8Runnable暂时休息
10:52:26.256 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - M3u8Runnable暂时休息
10:52:27.256 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - M3u8Runnable暂时休息
10:52:28.256 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - M3u8Runnable暂时休息
10:52:28.256 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - 转mp4结束! 
10:52:29.256 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - 开始转m3u8...
10:52:32.256 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - 转m3u8结束!
1
2
3
4
5
6
7
8

从日志可以看出符合预期,虽然结果正确,但是方案1存在问题:

1、M3u8Runnable线程不能及时得到通知。

M3u8Runnable线程检查MP4转换未完成则休眠数秒,如果正在休眠时完成了MP4转换这时是得不到通知的,需要等休眠时间到下一轮循环方可得到通知。

有人说可以将休眠时间设置小一点就可以相对及时得到通知,此时又会存在另一个问题就是:线程B在没有工作时会频繁占用CPU。

1.1.3 wait&notify工作流程 ​

基于方案1存在问题,我们思考如何及时收到通知还不用频繁占用CPU。

首先使用一个同步锁让线程A和线程B去竞争,保证两者只能有一个去工作。

如果是线程A先抢到锁正符合要求,线程A完成任务释放锁,然后线程B拿到锁去工作。如下图:

配图2

如果是线程B先抢到了锁,此时线程B需要让出锁进入等待队列,线程A拿到锁去工作,完成工作后唤醒线程B。如下图:

配图3

线程B如何释放锁进入等待队列?线程A如何唤醒等待中的线程,具体的工作流程如下:

1、线程A和线程B争抢锁,线程B抢到锁,线程A在监视器的阻塞队列中,如下图:

配图4

2、线程B调用wait方法释放锁并进入等待队列,由于锁被释放线程A获取锁开始执行,如下图:

配图5

3、线程A处理完成调用notify方法,通知监视器的中正在等待的一个线程(注意只会通知到一个线程),线程A结束,此时等待队列 中只有线程B所以线程B收到通知并进入阻塞队列,如下图:

配图6

4、线程B争抢锁,抢到锁开始执行,如下图:

配图7

1.1.4 wait&notify入门程序 ​

根据上边的分析的wait&notify的工作流程,编写 入门程序如下:

java
//wait/notify入门程序
	public static void test_waitnotify_first(){
		Object obj = new Object();
		new Thread(()->{
			synchronized (obj){
				logger.debug("拿到锁");
				try {
					logger.debug("进入等待队列");
					obj.wait();
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
				logger.debug("被唤醒");
		}
		},"A").start();
		try {
			TimeUnit.SECONDS.sleep(1);
		} catch (InterruptedException e) {
			e.printStackTrace();
		}
		new Thread(()->{
			synchronized (obj){
				logger.debug("拿到锁");
				try {
					TimeUnit.SECONDS.sleep(2);
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
				obj.notify();
				logger.debug("去唤醒等待的线程");
			}
		},"B").start();

	}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34

运行:

java
13:13:40.393 FINE com.yjoffer.javase.thread2.waitnotify.ThreadWaitnotifyTest[A] - 拿到锁
13:13:40.417 FINE com.yjoffer.javase.thread2.waitnotify.ThreadWaitnotifyTest[A] - 进入等待队列
13:13:41.393 FINE com.yjoffer.javase.thread2.waitnotify.ThreadWaitnotifyTest[B] - 拿到锁
13:13:43.393 FINE com.yjoffer.javase.thread2.waitnotify.ThreadWaitnotifyTest[B] - 去唤醒等待的线程
13:13:43.393 FINE com.yjoffer.javase.thread2.waitnotify.ThreadWaitnotifyTest[A] - 被唤醒
1
2
3
4
5

wait/notify工作流程需要注意以下几点:

1、wait/notify都是Object类的方法。

2、等待队列与阻塞队列一样都与监视器关联,监视器与锁对象关联,所以可以认为等待队列与阻塞队列与锁对象关联。

3、线程只有获得锁才可以调用锁对象的wait方法,调用此方法表示自己释放锁并进入等待队列。

4、线程只有获得锁才可以调用锁对象的notify方法,调用此方法表示通知该锁对象关联的等待队列中一个线程(注意是通知其中一个线程),线程收到通知进入阻塞队列等待竞争锁,当竞争到锁后从当初调用wait方法处继续执行。

1.1.5 wait&notify实现线程协作案例 ​

根据wait/notify的工作流程,下边使用wait/notify实现本节的案例,代码如下:

重写两个任务类,一个是转MP4、一个是转m3u8。

java
/**
	 * 转mp4版本2,用wait/notify实现
	 */
	 class Mp4RunnableV2 implements Runnable{
		//视频
		private PbVideo pbVideo;

		public Mp4RunnableV2(PbVideo pbVideo){
			this.pbVideo = pbVideo;
		}
		public void run() {
			synchronized (pbVideo){
				if(pbVideo.isMp4finished()){
					return ;
				}
				logger.debug("开始转mp4...");
				try {
					//模拟视频转换过程
					TimeUnit.SECONDS.sleep(3);
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
				pbVideo.setMp4finished(true);
				logger.debug("转mp4结束! ");
				pbVideo.notify();
			}


		}
	}
	/**
	 * 转M3u8版本2,用wait/notify实现
	 */
	 class M3u8RunnableV2 implements Runnable{

		private PbVideo pbVideo;
		public M3u8RunnableV2(PbVideo pbVideo){
			this.pbVideo = pbVideo;
		}
		public void run() {
			synchronized (pbVideo){
				if(pbVideo.isM3u8finished()){
					return ;
				}
				while (!pbVideo.isMp4finished()){
					logger.debug("M3u8Runnable暂时休息");
					try {
						pbVideo.wait();
					} catch (InterruptedException e) {
						e.printStackTrace();
					}
				}
				logger.debug("开始转m3u8...");
				try {
					TimeUnit.SECONDS.sleep(3);
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
				pbVideo.setFinished(true);
				logger.debug("转m3u8结束! ");
			}


		}
	}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65

测试:

java
	public static void test_videoconvert2(){
		ExecutorService threadPool = Executors.newCachedThreadPool();
		//视频对象
		PbVideo pbVideo = new PbVideo("预见猿份Java自学教程");
		//启动一个线程开始转mp4
		threadPool.execute(new Mp4RunnableV2(pbVideo));
		//启动一个线程开始转m3u8
		threadPool.execute(new M3u8RunnableV2(pbVideo));
		shutdown(threadPool);
	}
1
2
3
4
5
6
7
8
9
10

运行程序,输出:

java
19:26:42.558 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - 开始转mp4...
19:26:45.579 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - 转mp4结束! 
19:26:45.579 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - 开始转m3u8...
19:26:48.580 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - 转m3u8结束!
1
2
3
4

1.1.6 wait与sleep的区别 ​

下面总结 wait与sleep的区别:

1)wait和notify是Object类的方法,sleep是Thread类的方法。

2)sleep在任意地方都可以执行,wait和notify只能在synchronized块中执行,这表示只是获取锁后方可执行。

3)调用wait后当前线程释放锁,如果在synchronized块中调用sleep是不会释放锁的。

4)sleep方法不需要被唤醒,休眠时间到自动唤醒,wait()方法需要被唤醒,其它线程调用notify/notifyAll唤醒正在wait的线程。

注意:wait()方法还有一个重载方法wait(long timeout),它可以指定一个超时时间,到达超时时间自动从等待队列到阻塞队列准备竞争锁,此时它不需要其它线程调用notify/notifyAll。

1.1.7 升级线程协作需求 ​

现在此案例的需求需要升级,如下:

1、增加视频加密功能。

2、视频切割后即转m3u8后进行视频加密。

根据升级的需求进行分析:

1、增加一个视频加密的线程。

2、三个线程的关系是:转mp4线程先工作,完成后通知转m3u8线程,m3u8线程完成后通知视频加密线程。

代码如下:

定义视频类,在原有基础上增加视频加密状态。

java
package com.yjoffer.javase.thread.waitnotify;
1

import com.yjoffer.javase.config.Logger;

import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;

/**

  • 线程协作测试
  • @author 预见猿份(www.yjoffer.com)

*/
public class PbVideo {

	// 视频名称
	String name;

	// mp4完成状态
	boolean mp4finished = false;

	//m3u8完成状态
	boolean m3u8finished = false;

	//视频转换结束
	boolean finished = false;

	//加密状态
	boolean encrypted = false;

	public PbVideo(String name) {
		this.name = name;
	}

	public String getName() {
		return name;
	}

	public void setName(String name) {
		this.name = name;
	}

	public boolean isMp4finished() {
		return mp4finished;
	}

	public void setMp4finished(boolean mp4finished) {
		this.mp4finished = mp4finished;
	}

	public boolean isFinished() {
		return finished;
	}

	public void setFinished(boolean finished) {
		this.finished = finished;
	}

	public boolean isM3u8finished() {
		return m3u8finished;
	}

	public void setM3u8finished(boolean m3u8finished) {
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
		this.m3u8finished = m3u8finished;
	}

	public boolean isEncrypted() {
		return encrypted;
	}

	public void setEncrypted(boolean encrypted) {
		this.encrypted = encrypted;
	}
}</code></pre>
1
2
3
4
5
6
7
8
9
10
11

定义三个任务类:

java
/**
	 * 转mp4版本3,用wait/notify实现
	 */
	 class Mp4RunnableV3 implements Runnable{
		//视频
		private PbVideo pbVideo;

		public Mp4RunnableV3(PbVideo pbVideo){
			this.pbVideo = pbVideo;
		}
		public void run() {
			synchronized (pbVideo){
				if(pbVideo.isMp4finished()){
					return ;
				}
				logger.debug("开始转mp4...");
				try {
					//模拟视频转换过程
					TimeUnit.SECONDS.sleep(3);
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
				pbVideo.setMp4finished(true);
				logger.debug("转mp4结束! ");
				//通知后续工作线程
				pbVideo.notify();
			}


		}
	}
	/**
	 * 转M3u8版本3,用wait/notify实现
	 */
	 class M3u8RunnableV3 implements Runnable{

		private PbVideo pbVideo;
		public M3u8RunnableV3(PbVideo pbVideo){
			this.pbVideo = pbVideo;
		}
		public void run() {
			synchronized (pbVideo){
				if(pbVideo.isM3u8finished()){
					return ;
				}
				while (!pbVideo.isMp4finished()){
					logger.debug("M3u8Runnable暂时休息");
					try {
						pbVideo.wait();
					} catch (InterruptedException e) {
						e.printStackTrace();
					}
				}
				if(pbVideo.isMp4finished() && !pbVideo.isM3u8finished()){
					logger.debug("开始转m3u8...");
					try {
						TimeUnit.SECONDS.sleep(3);
					} catch (InterruptedException e) {
						e.printStackTrace();
					}
					pbVideo.setM3u8finished(true);
					logger.debug("转m3u8结束! ");
					//通知后续工作线程
					pbVideo.notify();
				}else{
                    logger.debug("未达到转m3u8条件");
                }
			}


		}
	}
	/**
	 * 视频加密版本3,用wait/notify实现
	 */
	 class EncRunnableV3 implements Runnable{

		private PbVideo pbVideo;
		public EncRunnableV3(PbVideo pbVideo){
			this.pbVideo = pbVideo;
		}
		public void run() {
			synchronized (pbVideo){
				if(pbVideo.isEncrypted()){
					return ;
				}
				if (!pbVideo.isM3u8finished()){
					logger.debug("EncRunnable暂时休息");
					try {
						pbVideo.wait();
					} catch (InterruptedException e) {
						e.printStackTrace();
					}
				}
				if(pbVideo.isM3u8finished() && !pbVideo.isEncrypted()){
					logger.debug("开始加密...");
					try {
						TimeUnit.SECONDS.sleep(3);
					} catch (InterruptedException e) {
						e.printStackTrace();
					}
					pbVideo.setEncrypted(true);
					pbVideo.setFinished(true);
					logger.debug("加密结束! ");
				}else{
                    logger.debug("未达到加密条件");
            	}
			}


		}
	}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112

测试代码:

注意:先启动加密线程进行测试。

java
public static void test_videoconvert3(){
		ExecutorService threadPool = Executors.newCachedThreadPool();
		//视频对象
		PbVideo pbVideo = new PbVideo("预见猿份Java自学教程");
		//视频加密
        threadPool.execute(new EncRunnableV3(pbVideo));
        try {
            TimeUnit.SECONDS.sleep(1);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        //转mp4
        threadPool.execute(new Mp4RunnableV3(pbVideo));
        try {
            TimeUnit.SECONDS.sleep(1);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        //转m3u8
        threadPool.execute(new M3u8RunnableV3(pbVideo));
		shutdown(threadPool);
	}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22

运行程序,输出 :

java
09:58:04.864 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - EncRunnable暂时休息
09:58:04.893 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - M3u8Runnable暂时休息
09:58:05.866 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-3] - 开始转mp4...
09:58:08.866 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-3] - 转mp4结束! 
09:58:08.866 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - 没到加密条件!
1
2
3
4
5

从日志可以看出执行流程如下:

1、首先启动加密线程和m3u8线程,由于工作条件未达到两个线程都进入等待队列。

2、1秒后启动mp4线程,正常执行,转mp4结束后调用notify方法通知后续线程,结果加密线程收到通知。

3、由于加密线程的工作条件未达到,最终结束运行。

4、程序运行到最后m3u8线程还在等待队列等待通知。

此程序的问题有两个:

1、加密线程收到notify通知发现工作条件未达到却结束了线程,正确作法应该继续等待。

2、mp4线程完成工作后通知后续线程应该要通知到m3u8线程,结果却通知了加密线程。

1.1.8 等待重试功能 ​

解决第一个问题很简单只需要将if改为while即可,代码如下:

java
//如果工作条件没有达到继续进入下一次循环继续等待,直到工作条件达到为止
while (!pbVideo.isM3u8finished()){
    logger.debug("EncRunnable暂时休息");
    try {
    	pbVideo.wait();
    } catch (InterruptedException e) {
    	e.printStackTrace();
    }
}
1
2
3
4
5
6
7
8
9

运行程序,输出 如下:

java
10:10:31.882 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - EncRunnable暂时休息
10:10:31.903 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - M3u8Runnable暂时休息
10:10:32.882 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-3] - 开始转mp4...
10:10:35.883 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-3] - 转mp4结束! 
10:10:35.883 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - EncRunnable暂时休息
1
2
3
4
5

虽然程序最终结果还不正确,但是从日志可以看出加密线程的工作条件未达到时又一次开始休息,这是符合我们的预期的。

所以在使用wait/notify时注意工作条件的判断使用while,模板代码如下:

java
synchronized(obj){
	while(工作条件未达到){
		obj.wait();
	}
	//开始工作
}
1
2
3
4
5
6
7

1.1.9 解决虚假唤醒问题 ​

为什么会出现问题2呢?

因为notify方法会通知等待队列中的任意一个线程,由于这个是随机的所以mp4线程之后应该是m3u8线程去工作结果却通知到其它线程,这个现象叫虚假唤醒。

下图还原了程序的执行过程:

1)最开始m3u8线程和加密线程在等待队列

配图8

2)虚假唤醒后见下图

此时m3u8线程还在等待队列,加密线程在阻塞队列 等待获取锁。

配图9

解决虚假唤醒要使用notifyAll方法,此方法会通知等待队列中的所有线程,见下图:

mp4线程调用notifyAll方法后等待队列 中的所有线程到阻塞队列竞争获取锁。

即使是加密线程得到锁,它判断工作条件未达到又会进入等待队列同时释放锁 ,此时m3u8线程就会得到锁,m3u8线程完成工作后再次调用notifyAll,此时加密线程进入阻塞队列,最后获取锁进行视频加密。

配图10

所以,在使用wait/notify时除了注意判断工作条件的方法以外还要注意虚假唤醒问题,总结模板代码如下:

java
synchronized(obj){
	while(工作条件未达到){
		obj.wait();
	}
	//开始工作
	obj.notifyAll();
}
//另一个线程
synchronized(obj){
	while(工作条件未达到){
		obj.wait();
	}
	//开始工作
	obj.notifyAll();
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15

根据上边的分析修改后的完整代码如下:

java
/**
	 * 转mp4版本3,用wait/notify实现
	 */
	 class Mp4RunnableV3 implements Runnable{
		//视频
		private PbVideo pbVideo;

		public Mp4RunnableV3(PbVideo pbVideo){
			this.pbVideo = pbVideo;
		}
		public void run() {
			synchronized (pbVideo){
				if(pbVideo.isMp4finished()){
					return ;
				}
				logger.debug("开始转mp4...");
				try {
					//模拟视频转换过程
					TimeUnit.SECONDS.sleep(3);
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
				pbVideo.setMp4finished(true);
				logger.debug("转mp4结束! ");
				//通知后续工作线程
				pbVideo.notifyAll();
			}


		}
	}
	/**
	 * 转M3u8版本3,用wait/notify实现
	 */
	 class M3u8RunnableV3 implements Runnable{

		private PbVideo pbVideo;
		public M3u8RunnableV3(PbVideo pbVideo){
			this.pbVideo = pbVideo;
		}
		public void run() {
			synchronized (pbVideo){
				if(pbVideo.isM3u8finished()){
					return ;
				}
				while (!pbVideo.isMp4finished()){
					logger.debug("M3u8Runnable暂时休息");
					try {
						pbVideo.wait();
					} catch (InterruptedException e) {
						e.printStackTrace();
					}
				}
				if(pbVideo.isMp4finished() && !pbVideo.isM3u8finished()){
					logger.debug("开始转m3u8...");
					try {
						TimeUnit.SECONDS.sleep(3);
					} catch (InterruptedException e) {
						e.printStackTrace();
					}
					pbVideo.setM3u8finished(true);
					logger.debug("转m3u8结束! ");
					//通知后续工作线程
					pbVideo.notifyAll();
				}
			}


		}
	}
	/**
	 * 视频加密版本3,用wait/notify实现
	 */
	 class EncRunnableV3 implements Runnable{

		private PbVideo pbVideo;
		public EncRunnableV3(PbVideo pbVideo){
			this.pbVideo = pbVideo;
		}
		public void run() {
			synchronized (pbVideo){
				if(pbVideo.isEncrypted()){
					return ;
				}
				//如果工作条件没有达到继续进入下一次循环继续等待,直到工作条件达到为止
				while (!pbVideo.isM3u8finished()){
					logger.debug("EncRunnable暂时休息");
					try {
						pbVideo.wait();
					} catch (InterruptedException e) {
						e.printStackTrace();
					}
				}
				if(pbVideo.isM3u8finished() && !pbVideo.isEncrypted()){
					logger.debug("开始加密...");
					try {
						TimeUnit.SECONDS.sleep(3);
					} catch (InterruptedException e) {
						e.printStackTrace();
					}
					pbVideo.setEncrypted(true);
					pbVideo.setFinished(true);
					logger.debug("加密结束! ");
				}
			}


		}
	}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110

重新运行测试程序,输出 如下:

java
10:27:35.772 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - EncRunnable暂时休息
10:27:35.791 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - M3u8Runnable暂时休息
10:27:36.773 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-3] - 开始转mp4...
10:27:39.773 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-3] - 转mp4结束! 
10:27:39.773 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - 开始转m3u8...
10:27:42.774 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - 转m3u8结束! 
10:27:42.774 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - 开始加密...
10:27:45.774 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - 加密结束!
1
2
3
4
5
6
7
8

程序结果符合预期。

1.2 生产者与消费者模式 ​

1.2.1 生产者消费者模式介绍 ​

生产者与消费者模式是一种设计模式,设计模式是生产中对一些通用问题经过实践实验总结出来的最佳的解决方案,生产者与消费者模式是多线程交互协作的解决方案,它分为两个角色:生产者和消费者,阻塞队列可作为一个存放共享数据的容器,生产者不断的向容器存入数据,消费者不断的从容器中获取数据,就好比一个超市,售货员不断的向货架摆放商品,顾客不断的从货架拿走商品。

下图是生产者与消费模式的示意图:

配图11

生产者和消费者都可以是多线程,生产者负责向阻塞队列放数据,消费者负责从阻塞队列取数据。

生产者和消费者模式的好处是将将生产者与消费者解耦,生产者与消费者之间不直接交互,两者解耦有利于软件的维护和扩展。

生产者与消费者模式适用于耗时任务处理、程序解耦合处理等场景。

生产中常用的消息队列就是生产者与消费者模式。

1.2.2 生产者消费者模式代码实现 ​

下面利用所学的wait/notify尝试实现生产者与消费者模式,首先进行分析:

1、阻塞队列

阻塞队列在生产者与消费者模式中是一个中介,当生产者向阻塞队列放数据时队列已满则生产者线程会阻塞,当消费者从阻塞队列取数据时队列为空则消费者队列会阻塞,如果队列由空变为非空则消费者线程解除阻塞继续从队列取数据,当队列由满变为不满则生产者解除阻塞继续放数据。

2、生产者与消费者

下图是利用集合和wait/notify实现的阻塞队列流程图:

配图12

1、生产者执行流程:

1)获取锁

2)判断队列是否已满,如果已满则进入等待队列,否则向队列写入数据

3)如已向队列写入数据需要通知等待队列的消费者去获取数据

2、消费者执行流程:

1)获取锁

2)判断队列是否为空,如果为空则进入等待队列,否则从队列取数据

3)如果已从队列取数据需要通知等待队列的生产者去生产数据

3、代码如下:

1)定义阻塞队列

java
package com.yjoffer.javase.thread.waitnotify;

import com.yjoffer.javase.config.Logger;

import java.util.LinkedList;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;


/**
 * 线程协作测试
 * @author 预见猿份(www.yjoffer.com)
 *
 */
public class PbBlockingQueue<E>{

		//阻塞队列容器
		private LinkedList<E> queue = new LinkedList<>();

		//队列的容量
		private int capacity;

		public PbBlockingQueue(int capacity){
			this.capacity = capacity;
		}

		//存数据
		public void put(E e) throws InterruptedException {
			synchronized (queue){
				while (queue.size() == capacity){
					logger.debug("队列满暂时休息");
					queue.wait();
				}
				queue.addLast(e);
				logger.debug("生产:"+e);
				queue.notifyAll();
			}
		}
		//取数据
		public E take() throws InterruptedException {
			synchronized (queue){
				while (queue.isEmpty()){
					logger.debug("队列空暂时休息");
					queue.wait();
				}
				E e = queue.remove();
				logger.debug("消费---》"+e);
				queue.notifyAll();
				return e;
			}
		}
	}
	...
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56

2)测试生产者与消费者

下边生产者与消费者各启动一个线程:

java
//测试自定义阻塞队列
	public static void test_blockingqueue(){
		ExecutorService threadPool = Executors.newCachedThreadPool();
		//定义阻塞队列
		PbBlockingQueue<PbVideo> blockingQueue = new PbBlockingQueue<PbVideo>(1);
		//生产者线程每隔1秒生产一条数据
		threadPool.execute(()->{
1
2
3
4
5
6
7
		int i = 1;
		while (true){
			try {
				blockingQueue.put(new PbVideo("预见猿份自学Java教程"+(i++)));
				TimeUnit.SECONDS.sleep(1);
			} catch (InterruptedException e) {
				e.printStackTrace();
			}
		}
	});
	//消费线程每隔2秒获取一次数据
	threadPool.execute(()->{
		while (true){
			try {
				PbVideo video = blockingQueue.take();
				TimeUnit.SECONDS.sleep(2);
			} catch (InterruptedException e) {
				e.printStackTrace();
			}

		}

	});


}
public static void main(String[] args) {
	test_blockingqueue();
}</code></pre>
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29

运行程序,输出日志如下 :

java
17:07:32.507 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - 队列空暂时休息
17:07:32.527 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - 生产:PbVideo{name='预见猿份自学Java教程1'}
17:07:32.527 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - 消费---》PbVideo{name='预见猿份自学Java教程1'}
17:07:33.527 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - 生产:PbVideo{name='预见猿份自学Java教程2'}
17:07:34.527 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - 消费---》PbVideo{name='预见猿份自学Java教程2'}
17:07:34.528 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - 生产:PbVideo{name='预见猿份自学Java教程3'}
17:07:35.528 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - 队列满暂时休息
17:07:36.527 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - 消费---》PbVideo{name='预见猿份自学Java教程3'}
17:07:36.527 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - 生产:PbVideo{name='预见猿份自学Java教程4'}
17:07:37.528 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - 队列满暂时休息
17:07:38.527 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - 消费---》PbVideo{name='预见猿份自学Java教程4'}
17:07:38.527 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - 生产:PbVideo{name='预见猿份自学Java教程5'}
...
1
2
3
4
5
6
7
8
9
10
11
12
13

从日志可以看出生产者与消费者之间通过阻塞队列传输数据。

生产者与消费者模式不限制生产者与消费者的数量,可以一对多、多对一、多对多。

下边测试多对多的情况即多个生产者与多个消费者:

生产者线程:2个线程,每个线程每隔1秒生产一条数据

消费者线程:3个线程,每个线程每隔2秒消费一条数据

java
//测试自定义阻塞队列,多对多
	public static void test_blockingqueue_manytomany(){
		ExecutorService threadPool = Executors.newCachedThreadPool();
		//定义阻塞队列
		PbBlockingQueue<PbVideo> blockingQueue = new PbBlockingQueue<PbVideo>(2);
		AtomicInteger number = new AtomicInteger(0);
		for (int i = 0; i < 2; i++) {
			//生产者线程每隔1秒生产一条数据
			threadPool.execute(()->{
				while (true){
					try {
						blockingQueue.put(new PbVideo("预见猿份自学Java教程"+(number.incrementAndGet())));
						TimeUnit.SECONDS.sleep(1);
					} catch (InterruptedException e) {
						e.printStackTrace();
					}
				}
			});
		}

		for (int i = 0; i < 3; i++) {
			//消费线程每隔2秒获取一次数据
			threadPool.execute(()->{
				while (true){
					try {
						PbVideo video = blockingQueue.take();
						TimeUnit.SECONDS.sleep(2);
					} catch (InterruptedException e) {
						e.printStackTrace();
					}

				}

			});
		}


	}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38

1.3 Condition实现多等待队列 ​

1.3.1 wait&notify存在的缺陷 ​

wait&notify是多线程利用synchronized锁实现的线程协作的方法,一个线程获取锁调用wait()方法释放锁进入等待队列,另一个线程获取锁调用notify方法唤醒等待队列中的一个线程,为了避免虚假唤醒调用notifyAll唤醒等待队列中的所有线程,这是wait&notify的实现方法。

使用notifyAll是为了解决虚假唤醒,但是也存在不足,即唤醒了无关的线程。

有没有一种可以有目标的去唤醒线程呢?

使用ReentrantLock的Condition可以实现多个等待队列,使用synchronized实现的是一个锁对象维护一个等待队列,而ReentrantLock可以实现一个锁可维护多个等待队列,并且在唤醒时可以控制唤醒某个等待队列中的线程。

1.3.2 Condition介绍 ​

使用ReentrantLock锁也可以实现类似 wait&notify的线程协作的方法,并且它的功能更强大,使用Condition可以实现多个等待队列。

Condition用于条件控制,创建方法如下:

java
//创建一个锁
ReentrantLock reentrantLock = new ReentrantLock();
//创建一个条件控制对象,可创建多个条件控制对象
Condition a = reentrantLock.newCondition();
Condition b = reentrantLock.newCondition();
1
2
3
4
5

使用ReentrantLock锁需要使用await()、signal()、signalAll()替换wait()、notify()、notifyAll()。

Condition的方法如下:

java
// 去等待队列等待,相当于Object的wait()
void await() throws InterruptedException;	
// 唤醒在等待队列等待的单个线程,相当于Object的notify()
void signal();		
// 唤醒在等待队列等待的所有线程,相当于Object的notifyAll()
void signalAll();
1
2
3
4
5
6

每个condition对应一个等待队列,假设有两个condition:a和b,当线程t1调用了a的await()方法那么线程t1进入a等待队列。

1.3.3 Condition入门程序 ​

下边使用Condition实现线程协作的入门程序,依次启动线程A、B、C,再按C、B、A的顺序执行。

java
//condition入门程序,依次启动线程A、B、C,再按C、B、A的顺序执行
	public static void test_condition_first(){
		ReentrantLock lock = new ReentrantLock();
		Condition a = lock.newCondition();
		Condition b = lock.newCondition();

		new Thread(()->{
			lock.lock();
			try {
				logger.debug("拿到锁");
				try {
					logger.debug("进入等待队列");
					a.await();
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
				logger.debug("执行...");
			}finally {
				lock.unlock();
			}
		},"A").start();
		try {
			TimeUnit.SECONDS.sleep(1);
		} catch (InterruptedException e) {
			e.printStackTrace();
		}
		new Thread(()->{
			lock.lock();
			try {
				logger.debug("拿到锁");
				try {
					logger.debug("进入等待队列");
					b.await();
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
				logger.debug("执行...");
				logger.debug("唤醒A线程");
				a.signalAll();
			}finally {
				lock.unlock();
			}
		},"B").start();
		try {
			TimeUnit.SECONDS.sleep(1);
		} catch (InterruptedException e) {
			e.printStackTrace();
		}
		new Thread(()->{
			lock.lock();
			try {
				logger.debug("拿到锁");
				logger.debug("执行...");
				b.signalAll();
				logger.debug("唤醒B线程");
			}finally {
				lock.unlock();
			}
		},"C").start();

	}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61

输出 :

java
11:44:32.207 FINE com.yjoffer.javase.thread2.waitnotify.ThreadWaitnotifyTest[A] - 拿到锁
11:44:32.231 FINE com.yjoffer.javase.thread2.waitnotify.ThreadWaitnotifyTest[A] - 进入等待队列
11:44:33.208 FINE com.yjoffer.javase.thread2.waitnotify.ThreadWaitnotifyTest[B] - 拿到锁
11:44:33.208 FINE com.yjoffer.javase.thread2.waitnotify.ThreadWaitnotifyTest[B] - 进入等待队列
11:44:34.210 FINE com.yjoffer.javase.thread2.waitnotify.ThreadWaitnotifyTest[C] - 拿到锁
11:44:34.210 FINE com.yjoffer.javase.thread2.waitnotify.ThreadWaitnotifyTest[C] - 执行...
11:44:34.211 FINE com.yjoffer.javase.thread2.waitnotify.ThreadWaitnotifyTest[C] - 唤醒B线程
11:44:34.212 FINE com.yjoffer.javase.thread2.waitnotify.ThreadWaitnotifyTest[B] - 执行...
11:44:34.212 FINE com.yjoffer.javase.thread2.waitnotify.ThreadWaitnotifyTest[B] - 唤醒A线程
11:44:34.213 FINE com.yjoffer.javase.thread2.waitnotify.ThreadWaitnotifyTest[A] - 执行...
1
2
3
4
5
6
7
8
9
10

1.3.4 Condition实现视频转换案例 ​

下边使用Condition实现预见猿份视频转换案例。

1、synchronized更换为ReentrantLock锁

原程序中使用PbVideo对象作为锁对象,现在改为ReentrantLock锁,在PbVideo类中添加成员变量:

java
//ReentrankLock
ReentrantLock lock = new ReentrantLock();

//m3u8线程等待队列
Condition m3u8Set = lock.newCondition();
//加密线程等待队列
Condition encryptSet = lock.newCondition();
1
2
3
4
5
6
7

共定义两个等待队列,分别存放转m3u8线程和视频加密线程。

代码如下:

java
package com.yjoffer.javase.thread.waitnotify;

import com.yjoffer.javase.config.Logger;

import java.util.LinkedList;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;


/**
 * 线程协作测试
 * @author 预见猿份(www.yjoffer.com)
 *
 */
public  class PbVideo{
		// 视频名称
		String name;

		// mp4完成状态
		boolean mp4finished = false;

		//m3u8完成状态
		boolean m3u8finished = false;

		//视频转换结束
		boolean finished = false;

		//加密状态
		boolean encrypted = false;

		//ReentrankLock
		ReentrantLock lock = new ReentrantLock();

		//m3u8线程等待队列
		Condition m3u8Set = lock.newCondition();
		//加密线程等待队列
		Condition encryptSet = lock.newCondition();

		public PbVideo(String name) {
			this.name = name;
		}

		public String getName() {
			return name;
		}

		public void setName(String name) {
			this.name = name;
		}

		public boolean isMp4finished() {
			return mp4finished;
		}

		public void setMp4finished(boolean mp4finished) {
			this.mp4finished = mp4finished;
		}

		public boolean isFinished() {
			return finished;
		}

		public void setFinished(boolean finished) {
			this.finished = finished;
		}

		public boolean isM3u8finished() {
			return m3u8finished;
		}

		public void setM3u8finished(boolean m3u8finished) {
			this.m3u8finished = m3u8finished;
		}

		public boolean isEncrypted() {
			return encrypted;
		}

		public void setEncrypted(boolean encrypted) {
			this.encrypted = encrypted;
		}

		public ReentrantLock getLock() {
			return lock;
		}

		public Condition getM3u8Set() {
			return m3u8Set;
		}

		public Condition getEncryptSet() {
			return encryptSet;
		}

		@Override
		public String toString() {
			return "PbVideo{" +
					"name='" + name + '\'' +
					'}';
		}
	}
...
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106

2、测试程序中将synchronized改为ReentrantLock,使用wait()/notifyAll()的地方改为await()/signalAll()

新定义Mp4RunnableV4、M3u8RunnableV4、EncRunnableV4分别为Mp4转换任务类、M3u8转换任务类、视频加密任务类。

注意条件控制对象的使用方法:

mp4线程完成任务后调用m3u8线程的等待队列的signalAll()

m3u8线程工作条件不满足时调用m3u8线程的等待队列的await()表示去该队列等待。

m3u8线程完成任务后调用视频加密线程的等待队列的signalAll()

视频加密线程工作条件不满足时调用视频加密线程的等待队列的await()表示去该队列等待。

改造后的程序如下:

java
/**
	 * 转mp4版本3,用await/signal实现
	 */
 class Mp4RunnableV4 implements Runnable{
		//视频
		private PbVideo pbVideo;
1
2
3
4
5
6
	private ReentrantLock lock;

	public Mp4RunnableV4(PbVideo pbVideo){
		this.pbVideo = pbVideo;
		this.lock = pbVideo.getLock();
	}
	public void run() {
		lock.lock();
		try {
			if(pbVideo.isMp4finished()){
				return ;
			}
			logger.debug("开始转mp4...");
			try {
				//模拟视频转换过程
				TimeUnit.SECONDS.sleep(3);
			} catch (InterruptedException e) {
				e.printStackTrace();
			}
			pbVideo.setMp4finished(true);
			logger.debug("转mp4结束! ");
			//通知m3u8
			pbVideo.getM3u8Set().signalAll();
		}finally {
			lock.unlock();
		}
	}
}
/**
 * 转M3u8版本3,用await/signal实现
 */
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31

class M3u8RunnableV4 implements Runnable{

	private PbVideo pbVideo;
	private ReentrantLock lock;
	public M3u8RunnableV4(PbVideo pbVideo){
		this.pbVideo = pbVideo;
		this.lock = pbVideo.getLock();
	}
	public void run() {
		lock.lock();
		try {
			if(pbVideo.isM3u8finished()){
				return ;
			}
			while (!pbVideo.isMp4finished()){
				logger.debug("M3u8Runnable暂时休息");
				try {
					pbVideo.getM3u8Set().await();
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
			}
			if(pbVideo.isMp4finished() && !pbVideo.isM3u8finished()){
				logger.debug("开始转m3u8...");
				try {
					TimeUnit.SECONDS.sleep(3);
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
				pbVideo.setM3u8finished(true);
				logger.debug("转m3u8结束! ");
				//通知视频加密
				pbVideo.getEncryptSet().signalAll();
			}else{
				logger.debug("没到转m3u8条件!");
			}
		}finally {
			lock.unlock();
		}
	}
}
/**
 * 视频加密版本3,用await/signal实现
 */
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42

class EncRunnableV4 implements Runnable{

	private PbVideo pbVideo;
	private ReentrantLock lock;
	public EncRunnableV4(PbVideo pbVideo){
		this.pbVideo = pbVideo;
		this.lock = pbVideo.getLock();
	}
	public void run() {
		lock.lock();
		try {
			if(pbVideo.isEncrypted()){
				return ;
			}
			//如果工作条件没有达到继续进入下一次循环继续等待,直到工作条件达到为止
			while (!pbVideo.isM3u8finished()){
				logger.debug("EncRunnable暂时休息");
				try {
					pbVideo.getEncryptSet().await();
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
			}
			if(pbVideo.isM3u8finished() && !pbVideo.isEncrypted()){
				logger.debug("开始加密...");
				try {
					TimeUnit.SECONDS.sleep(3);
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
				pbVideo.setEncrypted(true);
				pbVideo.setFinished(true);
				logger.debug("加密结束! ");
			}else{
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
				logger.debug("没到加密条件!");
			}
		}finally {
			lock.unlock();
		}
	}
}</code></pre>
1
2
3
4
5
6
7

2、测试程序如下:

java
	public static void test_videoconvert(){
		ExecutorService threadPool = Executors.newCachedThreadPool();
		//视频对象
		PbVideo pbVideo = new PbVideo("预见猿份Java自学教程");
		//启动一个线程开始转mp4
		threadPool.execute(new Mp4RunnableV4(pbVideo));
		//启动一个线程开始加密
		threadPool.execute(new EncRunnableV4(pbVideo));
		//启动一个线程开始转m3u8
		threadPool.execute(new M3u8RunnableV4(pbVideo));

		shutdown(threadPool);
	}
1
2
3
4
5
6
7
8
9
10
11
12
13

运行程序,输出 :

java
11:23:34.535 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - 开始转mp4...
11:23:37.556 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-1] - 转mp4结束! 
11:23:37.557 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - EncRunnable暂时休息
11:23:37.557 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-3] - 开始转m3u8...
11:23:40.558 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-3] - 转m3u8结束! 
11:23:40.558 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - 开始加密...
11:23:43.558 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[pool-1-thread-2] - 加密结束!
1
2
3
4
5
6
7

1.4 park&unpark入门 ​

1.4.1 park&unpark介绍 ​

park、unpark是LockSupport的两个方法,它们与wait/notify一样可以实现线程的等待与唤醒,后边学习并发工具类时会用到park/unpark,这里学习它打好基础。

查询Locksupport类的API:

java
static void park() 
禁止当前线程进行线程调度,除非许可证可用。  
static void unpark(Thread thread) 
为给定的线程提供许可证(如果尚未提供)。
1
2
3
4

1.4.2 park&unpark测试 ​

下边启动两个线程t1、t2,t1线程调用park方法,t2线程去唤醒t1线程。

java
package com.yjoffer.javase.thread.waitnotify;

import com.yjoffer.javase.config.Logger;

import java.util.LinkedList;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.LockSupport;
import java.util.concurrent.locks.ReentrantLock;


/**
 * 线程协作测试
 * @author 预见猿份(www.yjoffer.com)
 *
 */
public class ThreadWaitnotifyTest {

	private static Logger logger = Logger.getLogger(ThreadWaitnotifyTest.class);
//测试park、unpark
	public static void test_parkunpark(){
		Thread t1 = new Thread(() -> {
			logger.debug("park...");
			//开始等待
			LockSupport.park();
			logger.debug("continue...");
		},"t1");
		Thread t2 = new Thread(() -> {

			try {
				TimeUnit.SECONDS.sleep(2);
			} catch (InterruptedException e) {
				e.printStackTrace();
			}
			logger.debug("unpark...");
			//唤醒t1
			LockSupport.unpark(t1);
		},"t2");

		t1.start();
		t2.start();
	}

	public static void main(String[] args) {
		test_parkunpark();
	}

}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51

运行程序输出如下:

java
20:24:44.43 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[t1] - park...
20:24:46.43 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[t2] - unpark...
20:24:46.43 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[t1] - continue...
1
2
3

从输出 可以看出t1线程执行park方法后等待,2秒后t2线程唤醒t1线程继续执行。

1.4.3 park&unpark原理 ​

wait/notify是基于监视器来设计(详细请参考"wait/notify入门"章节),park&unpark不需要获取锁即可被调用,park&unpark基于“许可”的设计思想:

1、每一个线程有一个parker对象,对象中维护一个变量_counter变量,调用park方法设置为0,调用unpark设置为1。

2、调用park过程:

尝试获取许可,如果_counter>0则将其更改为0并拿到许可,线程继续运行,如果counter<1则等待。

3、调用unpark过程:

将_counter更改为1,如果等待队列有等待的线程则去唤醒它。

多次调用unpark也仅是将counter更改为1。

下边修改测试程序来测试,先调用unpark方法再调用park方法。

java
//测试park、unpark,先调用unpark再调用park
	public static void test_parkunpark2(){
		Thread t1 = new Thread(() -> {
			try {
				TimeUnit.SECONDS.sleep(2);
			} catch (InterruptedException e) {
				e.printStackTrace();
			}
			logger.debug("park...");
			//开始等待
			LockSupport.park();
			logger.debug("continue...");
		},"t1");
		Thread t2 = new Thread(() -> {
			logger.debug("unpark...");
			//唤醒t1
			LockSupport.unpark(t1);
		},"t2");

		t1.start();
		t2.start();
	}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22

运行:

java
21:41:15.498 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[t2] - unpark...
21:41:17.498 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[t1] - park...
21:41:17.498 FINE com.yjoffer.javase.thread.waitnotify.ThreadWaitnotifyTest[t1] - continue...
1
2
3

从输出 可以看出,先调用unpark方法,2秒后调用park方法并没有等待而是直接继续执行。

← 08 CAS教程10 常用的并发集合 →








如果发现文档内容有错误或排版错乱,请及时联系站长老苗修改,不胜感激。联系我们
关于我们 | 隐私政策 | 豫ICP备2026003386号-4 | 豫公网安备41010202004008号
目录

本页无章节