Java多線程并發(fā)生產(chǎn)者消費(fèi)者設(shè)計(jì)模式實(shí)例解析
一、兩個線程一個生產(chǎn)者一個消費(fèi)者
需求情景
兩個線程,一個負(fù)責(zé)生產(chǎn),一個負(fù)責(zé)消費(fèi),生產(chǎn)者生產(chǎn)一個,消費(fèi)者消費(fèi)一個。
涉及問題
- 同步問題:如何保證同一資源被多個線程并發(fā)訪問時(shí)的完整性。常用的同步方法是采用標(biāo)記或加鎖機(jī)制。
- wait() / nofity() 方法是基類Object的兩個方法,也就意味著所有Java類都會擁有這兩個方法,這樣,我們就可以為任何對象實(shí)現(xiàn)同步機(jī)制。
- wait()方法:當(dāng)緩沖區(qū)已滿/空時(shí),生產(chǎn)者/消費(fèi)者線程停止自己的執(zhí)行,放棄鎖,使自己處于等待狀態(tài),讓其他線程執(zhí)行。
- notify()方法:當(dāng)生產(chǎn)者/消費(fèi)者向緩沖區(qū)放入/取出一個產(chǎn)品時(shí),向其他等待的線程發(fā)出可執(zhí)行的通知,同時(shí)放棄鎖,使自己處于等待狀態(tài)。
代碼實(shí)現(xiàn)(共三個類和一個main方法的測試類)
Resource.java
package com.demo.ProducerConsumer;
/**
* 資源
* @author lixiaoxi
*
*/
public class Resource {
/*資源序號*/
private int number = 0;
/*資源標(biāo)記*/
private boolean flag = false;
/**
* 生產(chǎn)資源
*/
public synchronized void create() {
if (flag) {//先判斷標(biāo)記是否已經(jīng)生產(chǎn)了,如果已經(jīng)生產(chǎn),等待消費(fèi);
try {
wait();//讓生產(chǎn)線程等待
} catch (InterruptedException e) {
e.printStackTrace();
}
}
number++;//生產(chǎn)一個
System.out.println(Thread.currentThread().getName() + "生產(chǎn)者------------" + number);
flag = true;//將資源標(biāo)記為已經(jīng)生產(chǎn)
notify();//喚醒在等待操作資源的線程(隊(duì)列)
}
/**
* 消費(fèi)資源
*/
public synchronized void destroy() {
if (!flag) {
try {
wait();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
System.out.println(Thread.currentThread().getName() + "消費(fèi)者****" + number);
flag = false;
notify();
}
}
Producer.java
package com.demo.ProducerConsumer;
/**
* 生產(chǎn)者
* @author lixiaoxi
*
*/
public class Producer implements Runnable{
private Resource resource;
public Producer(Resource resource) {
this.resource = resource;
}
@Override
public void run() {
while (true) {
try {
Thread.sleep(10);
} catch (InterruptedException e) {
e.printStackTrace();
}
resource.create();
}
}
}
Consumer.java
package com.demo.ProducerConsumer;
/**
* 消費(fèi)者
* @author lixiaoxi
*
*/
public class Consumer implements Runnable{
private Resource resource;
public Consumer(Resource resource) {
this.resource = resource;
}
@Override
public void run() {
while (true) {
try {
Thread.sleep(10);
} catch (InterruptedException e) {
e.printStackTrace();
}
resource.destroy();
}
}
}
ProducerConsumerTest.java
package com.demo.ProducerConsumer;
public class ProducerConsumerTest {
public static void main(String args[]) {
Resource resource = new Resource();
new Thread(new Producer(resource)).start();//生產(chǎn)者線程
new Thread(new Consumer(resource)).start();//消費(fèi)者線程
}
}
打印結(jié)果:

以上打印結(jié)果可以看出沒有任何問題。
二、多個線程,多個生產(chǎn)者和多個消費(fèi)者的問題
需求情景
四個線程,兩個個負(fù)責(zé)生產(chǎn),兩個個負(fù)責(zé)消費(fèi),生產(chǎn)者生產(chǎn)一個,消費(fèi)者消費(fèi)一個。
涉及問題
notifyAll()方法:當(dāng)生產(chǎn)者/消費(fèi)者向緩沖區(qū)放入/取出一個產(chǎn)品時(shí),向其他等待的所有線程發(fā)出可執(zhí)行的通知,同時(shí)放棄鎖,使自己處于等待狀態(tài)。
再次測試代碼
ProducerConsumerTest.java
package com.demo.ProducerConsumer;
public class ProducerConsumerTest {
public static void main(String args[]) {
Resource resource = new Resource();
new Thread(new Producer(resource)).start();//生產(chǎn)者線程
new Thread(new Producer(resource)).start();//生產(chǎn)者線程
new Thread(new Consumer(resource)).start();//消費(fèi)者線程
new Thread(new Consumer(resource)).start();//消費(fèi)者線程
}
}
運(yùn)行結(jié)果:


通過以上打印結(jié)果發(fā)現(xiàn)問題
147生產(chǎn)了一次,消費(fèi)了兩次。169生產(chǎn)了,而沒有消費(fèi)。
原因分析
當(dāng)兩個線程同時(shí)操作生產(chǎn)者生產(chǎn)或者消費(fèi)者消費(fèi)時(shí),如果有生產(chǎn)者或消費(fèi)者的兩個線程都wait()時(shí),再次notify(),由于其中一個線程已經(jīng)改變了標(biāo)記而另外一個線程再次往下直接執(zhí)行的時(shí)候沒有判斷標(biāo)記而導(dǎo)致的。if判斷標(biāo)記,只有一次,會導(dǎo)致不該運(yùn)行的線程運(yùn)行了。出現(xiàn)了數(shù)據(jù)錯誤的情況。
解決方案
while判斷標(biāo)記,解決了線程獲取執(zhí)行權(quán)后,是否要運(yùn)行!也就是每次wait()后再notify()時(shí)先再次判斷標(biāo)記。
代碼改進(jìn)(Resource中的 if -> while)
Resource.java
package com.demo.ProducerConsumer;
/**
* 資源
* @author lixiaoxi
*
*/
public class Resource {
/*資源序號*/
private int number = 0;
/*資源標(biāo)記*/
private boolean flag = false;
/**
* 生產(chǎn)資源
*/
public synchronized void create() {
while (flag) {//先判斷標(biāo)記是否已經(jīng)生產(chǎn)了,如果已經(jīng)生產(chǎn),等待消費(fèi);
try {
wait();//讓生產(chǎn)線程等待
} catch (InterruptedException e) {
e.printStackTrace();
}
}
number++;//生產(chǎn)一個
System.out.println(Thread.currentThread().getName() + "生產(chǎn)者------------" + number);
flag = true;//將資源標(biāo)記為已經(jīng)生產(chǎn)
notify();//喚醒在等待操作資源的線程(隊(duì)列)
}
/**
* 消費(fèi)資源
*/
public synchronized void destroy() {
while (!flag) {
try {
wait();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
System.out.println(Thread.currentThread().getName() + "消費(fèi)者****" + number);
flag = false;
notify();
}
}
運(yùn)行結(jié)果:

再次發(fā)現(xiàn)問題
打印到某個值比如生產(chǎn)完187,程序運(yùn)行卡死了,好像鎖死了一樣。
原因分析
notify:只能喚醒一個線程,如果本方喚醒了本方,沒有意義。而且while判斷標(biāo)記+notify會導(dǎo)致”死鎖”。
解決方案
notifyAll解決了本方線程一定會喚醒對方線程的問題。
最后代碼改進(jìn)(Resource中的 notify() -> notifyAll())
Resource.java
package com.demo.ProducerConsumer;
/**
* 資源
* @author lixiaoxi
*
*/
public class Resource {
/*資源序號*/
private int number = 0;
/*資源標(biāo)記*/
private boolean flag = false;
/**
* 生產(chǎn)資源
*/
public synchronized void create() {
while (flag) {//先判斷標(biāo)記是否已經(jīng)生產(chǎn)了,如果已經(jīng)生產(chǎn),等待消費(fèi);
try {
wait();//讓生產(chǎn)線程等待
} catch (InterruptedException e) {
e.printStackTrace();
}
}
number++;//生產(chǎn)一個
System.out.println(Thread.currentThread().getName() + "生產(chǎn)者------------" + number);
flag = true;//將資源標(biāo)記為已經(jīng)生產(chǎn)
notifyAll();//喚醒在等待操作資源的線程(隊(duì)列)
}
/**
* 消費(fèi)資源
*/
public synchronized void destroy() {
while (!flag) {
try {
wait();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
System.out.println(Thread.currentThread().getName() + "消費(fèi)者****" + number);
flag = false;
notifyAll();
}
}
運(yùn)行結(jié)果:

以上就大功告成了,沒有任何問題。
再來梳理一下整個流程。按照示例,生產(chǎn)者消費(fèi)者交替運(yùn)行,每次生產(chǎn)后都有對應(yīng)的消費(fèi)者,測試類創(chuàng)建實(shí)例,如果是生產(chǎn)者先運(yùn)行,進(jìn)入run()方法,進(jìn)入create()方法,flag默認(rèn)為false,number+1,生產(chǎn)者生產(chǎn)一個產(chǎn)品,flag置為true,同時(shí)調(diào)用notifyAll()方法,喚醒所有正在等待的線程,接下來如果還是生產(chǎn)者運(yùn)行呢?這是flag為true,進(jìn)入while循環(huán),執(zhí)行wait()方法,接下來如果是消費(fèi)者運(yùn)行的話,調(diào)用destroy()方法,這時(shí)flag為true,消費(fèi)者購買了一次產(chǎn)品,隨即將flag置為false,并喚醒所有正在等待的線程。這就是一次完整的多生產(chǎn)者對應(yīng)多消費(fèi)者的問題。
三、使用Lock和Condition來解決生產(chǎn)者消費(fèi)者問題
上面的代碼有一個問題,就是我們?yōu)榱吮苊馑械木€程都處于等待的狀態(tài),使用了notifyAll方法來喚醒所有的線程,即notifyAll喚醒的是自己方和對方線程。如果我需要只是喚醒對方的線程,比如:生產(chǎn)者只能喚醒消費(fèi)者的線程,消費(fèi)者只能喚醒生產(chǎn)者的線程。
在jdk1.5當(dāng)中為我們提供了多線程的升級解決方案:
1. 將同步synchronized替換成了Lock操作。
2. 將Object中的wait,notify,notifyAll方法替換成了Condition對象。
3. 可以只喚醒對方的線程。
完整代碼:
Resource1.java
package com.demo.ProducerConsumer;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
/**
* 資源
* @author lixiaoxi
*
*/
public class Resource1 {
/*資源序號*/
private int number = 0;
/*資源標(biāo)記*/
private boolean flag = false;
private Lock lock = new ReentrantLock();
//使用lock建立生產(chǎn)者的condition對象
private Condition condition_pro = lock.newCondition();
//使用lock建立消費(fèi)者的condition對象
private Condition condition_con = lock.newCondition();
/**
* 生產(chǎn)資源
*/
public void create() throws InterruptedException {
try{
lock.lock();
//先判斷標(biāo)記是否已經(jīng)生產(chǎn)了,如果已經(jīng)生產(chǎn),等待消費(fèi)
while(flag){
//生產(chǎn)者等待
condition_pro.await();
}
//生產(chǎn)一個
number++;
System.out.println(Thread.currentThread().getName() + "生產(chǎn)者------------" + number);
//將資源標(biāo)記為已經(jīng)生產(chǎn)
flag = true;
//生產(chǎn)者生產(chǎn)完畢后,喚醒消費(fèi)者的線程(注意這里不是signalAll)
condition_con.signal();
}finally{
lock.unlock();
}
}
/**
* 消費(fèi)資源
*/
public void destroy() throws InterruptedException{
try{
lock.lock();
//先判斷標(biāo)記是否已經(jīng)消費(fèi)了,如果已經(jīng)消費(fèi),等待生產(chǎn)
while(!flag){
//消費(fèi)者等待
condition_con.await();
}
System.out.println(Thread.currentThread().getName() + "消費(fèi)者****" + number);
//將資源標(biāo)記為已經(jīng)消費(fèi)
flag = false;
//消費(fèi)者消費(fèi)完畢后,喚醒生產(chǎn)者的線程
condition_pro.signal();
}finally{
lock.unlock();
}
}
}
Producer1.java
package com.demo.ProducerConsumer;
/**
* 生產(chǎn)者
* @author lixiaoxi
*
*/
public class Producer1 implements Runnable{
private Resource1 resource;
public Producer1(Resource1 resource) {
this.resource = resource;
}
@Override
public void run() {
while (true) {
try {
Thread.sleep(10);
resource.create();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}
Consumer1.java
package com.demo.ProducerConsumer;
/**
* 消費(fèi)者
* @author lixiaoxi
*
*/
public class Consumer1 implements Runnable{
private Resource1 resource;
public Consumer1(Resource1 resource) {
this.resource = resource;
}
@Override
public void run() {
while (true) {
try {
Thread.sleep(10);
resource.destroy();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}
ProducerConsumerTest1.java
package com.demo.ProducerConsumer;
public class ProducerConsumerTest1 {
public static void main(String args[]) {
Resource1 resource = new Resource1();
new Thread(new Producer1(resource)).start();//生產(chǎn)者線程
new Thread(new Producer1(resource)).start();//生產(chǎn)者線程
new Thread(new Consumer1(resource)).start();//消費(fèi)者線程
new Thread(new Consumer1(resource)).start();//消費(fèi)者線程
}
}
運(yùn)行結(jié)果:

四、總結(jié)
1、如果生產(chǎn)者、消費(fèi)者都是1個,那么flag標(biāo)記可以用if判斷。這里有多個,必須用while判斷。
2、在while判斷的同時(shí),notify函數(shù)可能喚醒本類線程(如一個消費(fèi)者喚醒另一個消費(fèi)者),這會導(dǎo)致所有消費(fèi)者忙等待,程序無法繼續(xù)往下執(zhí)行。使用notifyAll函數(shù)代替notify可以解決這個問題,notifyAll可以保證非本類線程被喚醒(消費(fèi)者線程能喚醒生產(chǎn)者線程,反之也可以),解決了忙等待問題。
小心假死
生產(chǎn)者/消費(fèi)者模型最終達(dá)到的目的是平衡生產(chǎn)者和消費(fèi)者的處理能力,達(dá)到這個目的的過程中,并不要求只有一個生產(chǎn)者和一個消費(fèi)者。可以多個生產(chǎn)者對應(yīng)多個消費(fèi)者,可以一個生產(chǎn)者對應(yīng)一個消費(fèi)者,可以多個生產(chǎn)者對應(yīng)一個消費(fèi)者。
假死就發(fā)生在上面三種場景下。假死指的是全部線程都進(jìn)入了WAITING狀態(tài),那么程序就不再執(zhí)行任何業(yè)務(wù)功能了,整個項(xiàng)目呈現(xiàn)停滯狀態(tài)。
比方說有生產(chǎn)者A和生產(chǎn)者B,緩沖區(qū)由于空了,消費(fèi)者處于WAITING。生產(chǎn)者B處于WAITING,生產(chǎn)者A被消費(fèi)者通知生產(chǎn),生產(chǎn)者A生產(chǎn)出來的產(chǎn)品本應(yīng)該通知消費(fèi)者,結(jié)果通知了生產(chǎn)者B,生產(chǎn)者B被喚醒,發(fā)現(xiàn)緩沖區(qū)滿了,于是繼續(xù)WAITING。至此,兩個生產(chǎn)者線程處于WAITING,消費(fèi)者處于WAITING,系統(tǒng)假死。
上面的分析可以看出,假死出現(xiàn)的原因是因?yàn)閚otify的是同類,所以非單生產(chǎn)者/單消費(fèi)者的場景,可以采取兩種方法解決這個問題:
(1)synchronized用notifyAll()喚醒所有線程、ReentrantLock用signalAll()喚醒所有線程。
(2)用ReentrantLock定義兩個Condition,一個表示生產(chǎn)者的Condition,一個表示消費(fèi)者的Condition,喚醒的時(shí)候調(diào)用相應(yīng)的Condition的signal()方法就可以了。
以上就是本文的全部內(nèi)容,希望對大家的學(xué)習(xí)有所幫助,也希望大家多多支持腳本之家。
相關(guān)文章
SpringBoot+thymeleaf+Echarts+Mysql 實(shí)現(xiàn)數(shù)據(jù)可視化讀取的示例
本文主要介紹了SpringBoot+thymeleaf+Echarts+Mysql 實(shí)現(xiàn)數(shù)據(jù)可視化讀取的示例,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2022-04-04
Java實(shí)現(xiàn)經(jīng)典俄羅斯方塊游戲
俄羅斯方塊是一個最初由阿列克謝帕吉特諾夫在蘇聯(lián)設(shè)計(jì)和編程的益智類視頻游戲。本文將利用Java實(shí)現(xiàn)這一經(jīng)典的小游戲,需要的可以參考一下2022-01-01
淺析java中 Spring MVC 攔截器作用及其實(shí)現(xiàn)
本篇文章主要介紹了java中SpringMVC 攔截器的使用及其實(shí)例,需要的朋友可以參考2017-04-04
idea創(chuàng)建Spring項(xiàng)目的方法步驟(圖文)
這篇文章主要介紹了idea創(chuàng)建Spring項(xiàng)目的方法步驟(圖文),小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧2019-01-01
Java之不通過構(gòu)造函數(shù)創(chuàng)建一個對象問題
這篇文章主要介紹了Java之不通過構(gòu)造函數(shù)創(chuàng)建一個對象問題,具有很好的參考價(jià)值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教2024-03-03
Java多線程+鎖機(jī)制實(shí)現(xiàn)簡單模擬搶票的項(xiàng)目實(shí)踐
鎖是一種同步機(jī)制,用于控制對共享資源的訪問,在線程獲取到鎖對象后,可以執(zhí)行搶票操作,本文主要介紹了Java多線程+鎖機(jī)制實(shí)現(xiàn)簡單模擬搶票的項(xiàng)目實(shí)踐,具有一定的參考價(jià)值,感興趣的可以了解一下2024-02-02
java 實(shí)現(xiàn)雙向鏈表實(shí)例詳解
這篇文章主要介紹了java 實(shí)現(xiàn)雙向鏈表實(shí)例詳解的相關(guān)資料,需要的朋友可以參考下2017-03-03
使用mongoTemplate實(shí)現(xiàn)多條件加分組查詢方式
這篇文章主要介紹了使用mongoTemplate實(shí)現(xiàn)多條件加分組查詢方式,具有很好的參考價(jià)值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教2022-06-06

