Razumevanje modela proizvođač-potrošač od samih osnova
Model proizvođač-potrošač je veoma klasičan model višenitnog konkurentnog saradnog rada. Razumevanje problema proizvođač-potrošač može produbiti naše razumevanje konkurentnog programiranja.
Proizvođač-potrošač zapravo obuhvata dve vrste niti: jedna je nit proizvođača koja služi za proizvodnju podataka, a druga je nit potrošača koja služi za potrošnju podataka. Da bi se raskinula veza između proizvođača i potrošača, obično se koristi deljena oblast podataka, poput skladišta: proizvođač, nakon proizvodnje podataka, direktno ih stavlja u deljenu oblast podataka i ne mora da brine o ponašanju potrošača; potrošač samo uzima podatke iz deljene oblasti podataka i ne mora da brine o ponašanju proizvođača.

Ova deljena oblast podataka treba da ima sledeću funkciju konkurentne saradnje između niti:
- Ako je deljena oblast podataka puna, blokirati proizvođača da nastavi da proizvodi podatke.
- Ako je deljena oblast podataka prazna, blokirati potrošača da nastavi da troši podatke.
Prilikom implementacije problema proizvođač-potrošač mogu se koristiti tri načina:
- Mehanizam obaveštavanja porukama wait/notify iz klase Object.
- Mehanizam obaveštavanja porukama await/signal iz Condition klase koja ide uz Lock.
- Implementacija pomoću BlockingQueue.
Mehanizam obaveštavanja porukama wait/notify
Komunikacija između niti može se ostvariti putem metoda wait i notify ili notifyAll objekata klase Object.

Poziv metode wait blokiraće trenutnu nit dok druga nit ne pozove metodu notify ili notifyAll da je obavesti; tek tada trenutna nit može da se vrati iz metode wait i nastavi sa sledećim operacijama.
Ove smo pojmove zapravo obradili pri priči o Condition-u, sigurno se još sećate.
- wait
Ova metoda služi da stavi trenutnu nit u stanje mirovanja, sve dok ne dobije obaveštenje ili ne bude prekinuta.
Pre poziva wait, nit mora da pribavi monitor bravu tog objekta, odnosno wait metoda se može pozvati samo u sinhronizovanoj metodi ili sinhronizovanom bloku. Nakon poziva wait metode, trenutna nit pušta bravu. Ako pri pozivu wait metode nit nije pribavila bravu, biće bačen IllegalMonitorStateException. Tek kada ponovo pribavi bravu, trenutna nit može uspešno da se vrati iz metode wait.
- notify
Ova metoda se takođe mora pozvati u sinhronizovanoj metodi ili sinhronizovanom bloku, odnosno pre poziva nit mora da pribavi bravu na nivou objekta tog objekta. Ako pri pozivu notify ne drži odgovarajuću bravu, takođe će biti bačen IllegalMonitorStateException.
Ova metoda bira jednu nit iz WAITTING stanja da obavesti, čime se nit koja je pozvala wait premešta iz reda čekanja u red za sinhronizaciju i čeka priliku da ponovo pribavi bravu, čime nit koja je pozvala wait može da izađe iz metode wait.
Nakon poziva notify, trenutna nit ne pušta odmah bravu objekta; tek kada program izađe iz sinhronizovanog bloka, trenutna nit pušta bravu.
- notifyAll
Ova metoda radi na isti način kao notify, ali je jedna važna razlika: notifyAll čini da sve niti koje su čekale na tom objektu izađu iz WAITTING stanja, tako da se sve premeštaju iz reda čekanja u red za sinhronizaciju i čekaju sledeću priliku da pribave monitor bravu objekta.
Međutim, mehanizam obaveštavanja porukama wait/notify ima izvesnih problema.
1. Rano obaveštavanje notify
Propuštanje notify obaveštenja, odnosno: nit A još nije počela sa wait, a nit B je već pozvala notify, tako da obaveštenje niti B ne izaziva nikakvu reakciju. Kada nit B izađe iz sinhronizovanog bloka, a nit A tek tada počne sa wait, ona će biti trajno blokirana u čekanju, sve dok je ne prekine neka druga nit.
Sledeći primer koda simulira problem koji donosi rano obaveštenje notify:
public class EarlyNotify {
private static String lockObject = "";
public static void main(String[] args) {
WaitThread waitThread = new WaitThread(lockObject);
NotifyThread notifyThread = new NotifyThread(lockObject);
notifyThread.start();
try {
Thread.sleep(3000);
} catch (InterruptedException e) {
e.printStackTrace();
}
waitThread.start();
}
static class WaitThread extends Thread {
private String lock;
public WaitThread(String lock) {
this.lock = lock;
}
@Override
public void run() {
synchronized (lock) {
try {
System.out.println(Thread.currentThread().getName() + " ulazi u blok koda");
System.out.println(Thread.currentThread().getName() + " počinje wait");
lock.wait();
System.out.println(Thread.currentThread().getName() + " završava wait");
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}
static class NotifyThread extends Thread {
private String lock;
public NotifyThread(String lock) {
this.lock = lock;
}
@Override
public void run() {
synchronized (lock) {
System.out.println(Thread.currentThread().getName() + " ulazi u blok koda");
System.out.println(Thread.currentThread().getName() + " počinje notify");
lock.notify();
System.out.println(Thread.currentThread().getName() + " završava notify");
}
}
}
}U primeru su pokrenute dve niti, jedna je WaitThread, a druga NotifyThread. NotifyThread se prva pokreće i poziva metodu notify. Zatim se WaitThread pokreće i poziva metodu wait, ali pošto je obaveštenje već prošlo, metoda wait više ne može da dobije odgovarajuće obaveštenje, pa će WaitThread trajno biti blokirana u metodi wait — to je pojava ranog obaveštenja.
Rešenje ovog problema je dodavanje statusne oznake koja dozvoljava da waitThread pre poziva metode wait proveri da li se status promio. Ako je obaveštenje već poslano, WaitThread više ne ide na wait. Optimizovan kod je ispod:
public class EarlyNotify {
private static String lockObject = "";
private static boolean isWait = true;
public static void main(String[] args) {
WaitThread waitThread = new WaitThread(lockObject);
NotifyThread notifyThread = new NotifyThread(lockObject);
notifyThread.start();
try {
Thread.sleep(3000);
} catch (InterruptedException e) {
e.printStackTrace();
}
waitThread.start();
}
static class WaitThread extends Thread {
private String lock;
public WaitThread(String lock) {
this.lock = lock;
}
@Override
public void run() {
synchronized (lock) {
try {
while (isWait) {
System.out.println(Thread.currentThread().getName() + " ulazi u blok koda");
System.out.println(Thread.currentThread().getName() + " počinje wait");
lock.wait();
System.out.println(Thread.currentThread().getName() + " završava wait");
}
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}
static class NotifyThread extends Thread {
private String lock;
public NotifyThread(String lock) {
this.lock = lock;
}
@Override
public void run() {
synchronized (lock) {
System.out.println(Thread.currentThread().getName() + " ulazi u blok koda");
System.out.println(Thread.currentThread().getName() + " počinje notify");
lock.notifyAll();
isWait = false;
System.out.println(Thread.currentThread().getName() + " završava notify");
}
}
}
}Ovaj kod je dodao samo status isWait. NotifyThread nakon poziva metode notify ažurira status, a WaitThread pre poziva metode wait proverava status.
U ovom primeru, nakon poziva notify status isWait se menja u false, pa u WaitThread-u while petlja nakon provere isWait neće izvršiti metodu wait, čime se izbegava propuštanje usled ranog obaveštenja Notify-ja.
Rezime: pri korišćenju mehanizma čekanje/obaveštenje niti, obično treba kombinovati ga sa boolean promenljivom čija se vrednost menja pre notify, tako da wait nakon povratka može da izađe iz while petlje, ili da nakon propuštenog obaveštenja ne bude blokiran u metodi wait.
2. Promena uslova čekanja na wait
Ako nit za vreme čekanja dobije obaveštenje, ali se zatim uslov čekanja promenio i ne proveri ponovo uslov čekanja, to takođe može dovesti do greške u programu.
Hajde da ovo objasnimo primerom.
public class ConditionChange {
private static List<String> lockObject = new ArrayList();
public static void main(String[] args) {
Consumer consumer1 = new Consumer(lockObject);
Consumer consumer2 = new Consumer(lockObject);
Productor productor = new Productor(lockObject);
consumer1.start();
consumer2.start();
productor.start();
}
static class Consumer extends Thread {
private List<String> lock;
public Consumer(List lock) {
this.lock = lock;
}
@Override
public void run() {
synchronized (lock) {
try {
//Ovde korišćenjem if-a stvara se problem greške programa zbog promene uslova wait-a
if (lock.isEmpty()) {
System.out.println(Thread.currentThread().getName() + " lista je prazna");
System.out.println(Thread.currentThread().getName() + " poziva metodu wait");
lock.wait();
System.out.println(Thread.currentThread().getName() + " metoda wait je završena");
}
String element = lock.remove(0);
System.out.println(Thread.currentThread().getName() + " uzima prvi element: " + element);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}
static class Productor extends Thread {
private List<String> lock;
public Productor(List lock) {
this.lock = lock;
}
@Override
public void run() {
synchronized (lock) {
System.out.println(Thread.currentThread().getName() + " počinje da dodaje element");
lock.add(Thread.currentThread().getName());
lock.notifyAll();
}
}
}
}Baciće izuzetak:
Exception in thread "Thread-1" Thread-0 lista je prazna
Thread-0 poziva metodu wait
Thread-1 lista je prazna
Thread-1 poziva metodu wait
Thread-2 počinje da dodaje element
Thread-1 metoda wait je završena
java.lang.IndexOutOfBoundsException: Index: 0, Size: 0U ovom primeru je pokrenuto ukupno 3 niti: Consumer1, Consumer2 i Productor.
Nakon što Consumer1 pozove metodu wait, nit prelazi u WAITTING stanje i pušta bravu objekta.
Tada Consumer2 pribavlja bravu objekta i ulazi u sinhronizovani blok; kada dođe do metode wait, takođe pušta bravu objekta.
Zatim productor pribavlja bravu objekta, ulazi u sinhronizovani blok, umeće podatke u listu i putem metode notifyAll obaveštava Consumer1 i Consumer2 niti koje su u WAITING stanju.
Nakon što consumer1 dobije bravu objekta, izlazi iz metode wait, briše jedan element što listu čini praznom, metoda se završava, izlazi iz sinhronizovanog bloka i pušta bravu objekta.
U tom trenutku Consumer2 pribavlja bravu objekta, izlazi iz metode wait i nastavlja dalje. Kada Consumer2 izvrši lock.remove(0); doći će do greške, jer je lista već prazna.
Rešenje: Iz gore navedene analize se vidi da je Consumer2 prijavio grešku jer nakon izlaska iz metode wait nije proverio uslov čekanja, ali se uslov čekanja u međuvremenu promenio. Rešenje je da se nakon izlaska iz wait-a ponovo proveri uslov.
public class ConditionChange {
private static List<String> lockObject = new ArrayList();
public static void main(String[] args) {
Consumer consumer1 = new Consumer(lockObject);
Consumer consumer2 = new Consumer(lockObject);
Productor productor = new Productor(lockObject);
consumer1.start();
consumer2.start();
productor.start();
}
static class Consumer extends Thread {
private List<String> lock;
public Consumer(List lock) {
this.lock = lock;
}
@Override
public void run() {
synchronized (lock) {
try {
//Ovde korišćenjem if-a stvara se problem greške programa zbog promene uslova wait-a
while (lock.isEmpty()) {
System.out.println(Thread.currentThread().getName() + " lista je prazna");
System.out.println(Thread.currentThread().getName() + " poziva metodu wait");
lock.wait();
System.out.println(Thread.currentThread().getName() + " metoda wait je završena");
}
String element = lock.remove(0);
System.out.println(Thread.currentThread().getName() + " uzima prvi element: " + element);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}
static class Productor extends Thread {
private List<String> lock;
public Productor(List lock) {
this.lock = lock;
}
@Override
public void run() {
synchronized (lock) {
System.out.println(Thread.currentThread().getName() + " počinje da dodaje element");
lock.add(Thread.currentThread().getName());
lock.notifyAll();
}
}
}
}U poređenju sa prethodnim kodom, ovaj kod je samo if naredbu oko wait promenio u while petlju, tako da kada je lista prazna, nit nastavlja da čeka i neće nastaviti sa izvršavanjem koda koji briše elemente iz liste.
Rezime: pri korišćenju mehanizma čekanje/obaveštenje niti, wait metodu obično treba pozivati unutar while petlje. Zato je potrebno kombinovati je sa boolean promenljivom: kada je ispunjen uslov while petlje, ulazi se u while petlju i izvršava metoda wait; kada uslov while petlje nije ispunjen, izlazi se iz petlje i izvršava sledeći kod.
3. Stanje „lažne smrti"
Pojava: u situaciji sa više potrošača i više proizvođača, korišćenje metode notify može dovesti do „lažne smrti", odnosno svih niti se nalaze u stanju čekanja i ne mogu da budu probuđene.
Analiza uzroka: pretpostavimo da trenutno više niti proizvođača čeka blokirano u metodi wait. Kada jedna nit proizvođača pribavi bravu objekta i putem notify obavesti nit u WAITTING stanju, ako se probudi i dalje nit proizvođača, doći će do toga da sve niti proizvođača budu u stanju čekanja.
Rešenje: zameniti metodu notify metodom notifyAll; ako se koristi lock, zameniti metodu signal metodom signalAll.
Rezime: mehanizam obaveštavanja porukama koji pruža Object treba da poštuje sledeće uslove:
- Uvek proveravajte uslov u while petlji, a ne u if naredbi za proveru uslova wait-a.
- Koristite notifyAll umesto notify.
Osnovna paradigma korišćenja je ispod:
// The standard idiom for calling the wait method in Java
synchronized (sharedObject) {
while (condition) {
sharedObject.wait();
// (Releases lock, and reacquires on wakeup)
}
// do action based upon condition e.g. take or put into queue
}Implementacija proizvođač-potrošač pomoću wait/notifyAll
Kod za implementaciju proizvođača i potrošača pomoću wait/notifyAll je ispod:
public class ProductorConsumer {
public static void main(String[] args) {
LinkedList linkedList = new LinkedList();
ExecutorService service = Executors.newFixedThreadPool(15);
for (int i = 0; i < 5; i++) {
service.submit(new Productor(linkedList, 8));
}
for (int i = 0; i < 10; i++) {
service.submit(new Consumer(linkedList));
}
}
static class Productor implements Runnable {
private List<Integer> list;
private int maxLength;
public Productor(List list, int maxLength) {
this.list = list;
this.maxLength = maxLength;
}
@Override
public void run() {
while (true) {
synchronized (list) {
try {
while (list.size() == maxLength) {
System.out.println("Proizvođač " + Thread.currentThread().getName() + " lista je dostigla maksimalni kapacitet, poziva wait");
list.wait();
System.out.println("Proizvođač " + Thread.currentThread().getName() + " izlazi iz wait-a");
}
Random random = new Random();
int i = random.nextInt();
System.out.println("Proizvođač " + Thread.currentThread().getName() + " proizvodi podatak " + i);
list.add(i);
list.notifyAll();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}
}
static class Consumer implements Runnable {
private List<Integer> list;
public Consumer(List list) {
this.list = list;
}
@Override
public void run() {
while (true) {
synchronized (list) {
try {
while (list.isEmpty()) {
System.out.println("Potrošač " + Thread.currentThread().getName() + " lista je prazna, poziva wait");
list.wait();
System.out.println("Potrošač " + Thread.currentThread().getName() + " izlazi iz wait-a");
}
Integer element = list.remove(0);
System.out.println("Potrošač " + Thread.currentThread().getName() + " troši podatak: " + element);
list.notifyAll();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}
}
}Izlazni rezultat:
Proizvođač pool-1-thread-1 proizvodi podatak -232820990
Proizvođač pool-1-thread-1 proizvodi podatak 1432164130
Proizvođač pool-1-thread-1 proizvodi podatak 1057090222
Proizvođač pool-1-thread-1 proizvodi podatak 1201395916
Proizvođač pool-1-thread-1 proizvodi podatak 482766516
Proizvođač pool-1-thread-1 lista je dostigla maksimalni kapacitet, poziva wait
Potrošač pool-1-thread-15 izlazi iz wait-a
Potrošač pool-1-thread-15 troši podatak: 1237535349
Potrošač pool-1-thread-15 troši podatak: -1617438932
Potrošač pool-1-thread-15 troši podatak: -535396055
Potrošač pool-1-thread-15 troši podatak: -232820990
Potrošač pool-1-thread-15 troši podatak: 1432164130
Potrošač pool-1-thread-15 troši podatak: 1057090222
Potrošač pool-1-thread-15 troši podatak: 1201395916
Potrošač pool-1-thread-15 troši podatak: 482766516
Potrošač pool-1-thread-15 lista je prazna, poziva wait
Proizvođač pool-1-thread-5 izlazi iz wait-a
Proizvođač pool-1-thread-5 proizvodi podatak 1442969724
Proizvođač pool-1-thread-5 proizvodi podatak 1177554422
Proizvođač pool-1-thread-5 proizvodi podatak -133137235
Proizvođač pool-1-thread-5 proizvodi podatak 324882560
Proizvođač pool-1-thread-5 proizvodi podatak 2065211573
Proizvođač pool-1-thread-5 proizvodi podatak 253569900
Proizvođač pool-1-thread-5 proizvodi podatak 571277922
Proizvođač pool-1-thread-5 proizvodi podatak 1622323863
Proizvođač pool-1-thread-5 lista je dostigla maksimalni kapacitet, poziva wait
Potrošač pool-1-thread-10 izlazi iz wait-aImplementacija proizvođač-potrošač pomoću await/signalAll
Nasuprot metodama wait i notify/notifyAll iz klase Object, Condition pruža iste metode — await metodu i signal/signalAll metode. Ovo znanje smo obradili pri priči o Condition-u, sigurno se još sećate.
Ako se usvoji princip obaveštavanja porukama Condition-a za implementaciju modela proizvođač-potrošač, princip je isti kao kod wait/notifyAll. Evo koda:
public class ProductorConsumer {
private static ReentrantLock lock = new ReentrantLock();
private static Condition full = lock.newCondition();
private static Condition empty = lock.newCondition();
public static void main(String[] args) {
LinkedList linkedList = new LinkedList();
ExecutorService service = Executors.newFixedThreadPool(15);
for (int i = 0; i < 5; i++) {
service.submit(new Productor(linkedList, 8, lock));
}
for (int i = 0; i < 10; i++) {
service.submit(new Consumer(linkedList, lock));
}
}
static class Productor implements Runnable {
private List<Integer> list;
private int maxLength;
private Lock lock;
public Productor(List list, int maxLength, Lock lock) {
this.list = list;
this.maxLength = maxLength;
this.lock = lock;
}
@Override
public void run() {
while (true) {
lock.lock();
try {
while (list.size() == maxLength) {
System.out.println("Proizvođač " + Thread.currentThread().getName() + " lista je dostigla maksimalni kapacitet, poziva wait");
full.await();
System.out.println("Proizvođač " + Thread.currentThread().getName() + " izlazi iz wait-a");
}
Random random = new Random();
int i = random.nextInt();
System.out.println("Proizvođač " + Thread.currentThread().getName() + " proizvodi podatak " + i);
list.add(i);
empty.signalAll();
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
lock.unlock();
}
}
}
}
static class Consumer implements Runnable {
private List<Integer> list;
private Lock lock;
public Consumer(List list, Lock lock) {
this.list = list;
this.lock = lock;
}
@Override
public void run() {
while (true) {
lock.lock();
try {
while (list.isEmpty()) {
System.out.println("Potrošač " + Thread.currentThread().getName() + " lista je prazna, poziva wait");
empty.await();
System.out.println("Potrošač " + Thread.currentThread().getName() + " izlazi iz wait-a");
}
Integer element = list.remove(0);
System.out.println("Potrošač " + Thread.currentThread().getName() + " troši podatak: " + element);
full.signalAll();
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
lock.unlock();
}
}
}
}
}Izlazni rezultat:
Potrošač pool-1-thread-9 troši podatak: 1146627506
Potrošač pool-1-thread-9 troši podatak: 1508001019
Potrošač pool-1-thread-9 troši podatak: -600080565
Potrošač pool-1-thread-9 troši podatak: -1000305429
Potrošač pool-1-thread-9 troši podatak: -1270658620
Potrošač pool-1-thread-9 troši podatak: 1961046169
Potrošač pool-1-thread-9 troši podatak: -307680655
Potrošač pool-1-thread-9 lista je prazna, poziva wait
Potrošač pool-1-thread-13 izlazi iz wait-a
Potrošač pool-1-thread-13 lista je prazna, poziva wait
Potrošač pool-1-thread-10 izlazi iz wait-a
Proizvođač pool-1-thread-5 izlazi iz wait-a
Proizvođač pool-1-thread-5 proizvodi podatak -892558288
Proizvođač pool-1-thread-5 proizvodi podatak -1917220008
Proizvođač pool-1-thread-5 proizvodi podatak 2146351766
Proizvođač pool-1-thread-5 proizvodi podatak 452445380
Proizvođač pool-1-thread-5 proizvodi podatak 1695168334
Proizvođač pool-1-thread-5 proizvodi podatak 1979746693
Proizvođač pool-1-thread-5 proizvodi podatak -1905436249
Proizvođač pool-1-thread-5 proizvodi podatak -101410137
Proizvođač pool-1-thread-5 lista je dostigla maksimalni kapacitet, poziva wait
Proizvođač pool-1-thread-1 izlazi iz wait-a
Proizvođač pool-1-thread-1 lista je dostigla maksimalni kapacitet, poziva wait
Proizvođač pool-1-thread-4 izlazi iz wait-a
Proizvođač pool-1-thread-4 lista je dostigla maksimalni kapacitet, poziva wait
Proizvođač pool-1-thread-2 izlazi iz wait-a
Proizvođač pool-1-thread-2 lista je dostigla maksimalni kapacitet, poziva wait
Proizvođač pool-1-thread-3 izlazi iz wait-a
Proizvođač pool-1-thread-3 lista je dostigla maksimalni kapacitet, poziva wait
Potrošač pool-1-thread-9 izlazi iz wait-a
Potrošač pool-1-thread-9 troši podatak: -892558288Implementacija proizvođač-potrošač pomoću BlockingQueue
Kada smo pričali o BlockingQueue-u, rekli smo da je BlockingQueue veoma pogodan za implementaciju modela proizvođač-potrošač.
Razlog je taj što BlockingQueue pruža metode blokirajućeg umetanja i uklanjanja. Kada je kontejner reda pun, nit proizvođača se blokira dok red ne postane nepun; kada je kontejner reda prazan, nit potrošača se blokira dok red ne postane nepun.

Sa ovim redom, proizvođač samo treba da se fokusira na proizvodnju i ne mora da brine o ponašanju potrošnje potrošača, a kamoli da čeka da se nit potrošača završi; potrošač samo troši i ne mora da brine o tome kako proizvođač proizvodi, a kamoli da čeka da proizvođač proizvede.
Evo koda:
public class ProductorConsumer {
private static LinkedBlockingQueue<Integer> queue = new LinkedBlockingQueue<>();
public static void main(String[] args) {
ExecutorService service = Executors.newFixedThreadPool(15);
for (int i = 0; i < 5; i++) {
service.submit(new Productor(queue));
}
for (int i = 0; i < 10; i++) {
service.submit(new Consumer(queue));
}
}
static class Productor implements Runnable {
private BlockingQueue queue;
public Productor(BlockingQueue queue) {
this.queue = queue;
}
@Override
public void run() {
try {
while (true) {
Random random = new Random();
int i = random.nextInt();
System.out.println("Proizvođač " + Thread.currentThread().getName() + " proizvodi podatak " + i);
queue.put(i);
}
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
static class Consumer implements Runnable {
private BlockingQueue queue;
public Consumer(BlockingQueue queue) {
this.queue = queue;
}
@Override
public void run() {
try {
while (true) {
Integer element = (Integer) queue.take();
System.out.println("Potrošač " + Thread.currentThread().getName() + " troši podatak " + element);
}
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}Izlazni rezultat:
Potrošač pool-1-thread-7 troši podatak 1520577501
Proizvođač pool-1-thread-4 proizvodi podatak -127809610
Potrošač pool-1-thread-8 troši podatak 504316513
Proizvođač pool-1-thread-2 proizvodi podatak 1994678907
Potrošač pool-1-thread-11 troši podatak 1967302829
Proizvođač pool-1-thread-1 proizvodi podatak 369331507
Potrošač pool-1-thread-9 troši podatak 1994678907
Proizvođač pool-1-thread-2 proizvodi podatak -919544017
Potrošač pool-1-thread-12 troši podatak -127809610
Proizvođač pool-1-thread-4 proizvodi podatak 1475197572
Potrošač pool-1-thread-14 troši podatak -893487914
Proizvođač pool-1-thread-3 proizvodi podatak 906921688
Potrošač pool-1-thread-6 troši podatak -1292015016
Proizvođač pool-1-thread-5 proizvodi podatak -652105379
Proizvođač pool-1-thread-5 proizvodi podatak -1622505717
Proizvođač pool-1-thread-3 proizvodi podatak -1350268764
Potrošač pool-1-thread-7 troši podatak 906921688
Proizvođač pool-1-thread-4 proizvodi podatak 2091628867
Potrošač pool-1-thread-13 troši podatak 1475197572
Potrošač pool-1-thread-15 troši podatak -919544017
Proizvođač pool-1-thread-2 proizvodi podatak 564860122
Proizvođač pool-1-thread-2 proizvodi podatak 822954707
Potrošač pool-1-thread-14 troši podatak 564860122
Potrošač pool-1-thread-10 troši podatak 369331507
Proizvođač pool-1-thread-1 proizvodi podatak -245820912
Potrošač pool-1-thread-6 troši podatak 822954707
Proizvođač pool-1-thread-2 proizvodi podatak 1724595968
Proizvođač pool-1-thread-2 proizvodi podatak -1151855115
Potrošač pool-1-thread-12 troši podatak 2091628867
Proizvođač pool-1-thread-4 proizvodi podatak -1774364499
Proizvođač pool-1-thread-4 proizvodi podatak 2006106757
Potrošač pool-1-thread-14 troši podatak -1774364499
Proizvođač pool-1-thread-3 proizvodi podatak -1070853639
Potrošač pool-1-thread-9 troši podatak -1350268764
Potrošač pool-1-thread-11 troši podatak -1622505717
Proizvođač pool-1-thread-5 proizvodi podatak 355412953Vidi se da je implementacija proizvođač-potrošač pomoću BlockingQueue veoma koncizna — to je i prednost BlockingQueue-a.
Scenariji primene modela proizvođač-potrošač
Model proizvođač-potrošač se obično koristi da razdvoji stranu koja proizvodi podatke od strane koja ih troši, čime se procesi proizvodnje i potrošnje podataka međusobno raskidaju.
01. Okvir za izvršenje zadataka Executor:
Raskidanjem veze između podnošenja zadataka i izvršenja zadataka, operacija podnošenja zadataka odgovara proizvođaču, a operacija izvršenja zadataka odgovara potrošaču.
Na primer, korišćenje Executor-a za izgradnju web servera za obradu zahteva niti: proizvođač podnosi zadatke bazenu niti, a bazen niti kreira niti za obradu zadataka. Ako je broj zadataka koje treba pokrenuti veći od osnovnog broja niti u bazenu niti, zadaci se bacaju u red blokiranja (način sa bazenom niti + redom blokiranja je mnogo efikasniji od korišćenja samo jednog reda blokiranja, jer potrošač može da obradi direktno kada može, pa ne mora svaki potrošač prvo da izvadi zadatak iz reda blokiranja pa tek onda da ga izvrši).
02. Middleware za poruke MQ:
Tokom [popusta na 11. novembar], generiše se ogroman broj porudžbina, pa nije moguće istovremeno obraditi toliko porudžbina. Potrebno je porudžbine staviti u red, a zatim ih obraditi namenskim nitima.
Ovde korisničko naručivanje predstavlja proizvođača, a niti koje obrađuju porudžbine predstavljaju potrošača. Na primer, funkcija preuzimanja karata na 12306: prvo kontejner čuva korisničke porudžbine, a zatim namenske niti koje obrađuju porudžbine to rade polako, čime se u kratkom vremenskom roku može podržati visoko-konkurentna usluga.
03. Kada je vreme obrade zadatka relativno dugo:
Na primer, otpremanje priloga i njihova obrada — tada se otpremanje korisnika i obrada priloga mogu podeliti u dva procesa. Red se koristi za privremeno čuvanje korisničkih priloga, zatim se korisniku odmah vraća da je otpremanje uspelo, a zatim namenske niti obrađuju priloge u redu.
Prednosti modela proizvođač-potrošač:
- Raskidanje veze: raskida se veza između klase proizvođača i klase potrošača, eliminišu se zavisnosti u kodu i pojednostavljuje upravljanje opterećenjem.
- Ponovna upotreba: nezavisnim odvajanjem klase proizvođača i klase potrošača, omogućava se nezavisna ponovna upotreba i proširenje obe klase.
- Podešavanje broja konkurentnih jedinica: pošto su brzine obrade proizvođača i potrošača različite, može se podesiti broj konkurentnih jedinica, dajući više konkurentnih jedinica sporijoj strani kako bi se poboljšala brzina obrade zadataka.
- Asinhronost: za proizvođača i potrošača svako radi svoj posao — proizvođač samo mora da brine o tome da li u baferu ima podataka i ne mora da čeka da potrošač završi obradu; za potrošača, takođe se samo fokusira na sadržaj bafera i ne mora da prati proizvođača. Asinhronim načinom podržava se visoka konkurentnost: vremenski zahtevan proces se deli na dve faze — proizvodnju i potrošnju, tako da pošto je vreme izvršenja put kratko, proizvođač može da podrži visoku konkurentnost.
- Podrška za distribuirane sisteme: proizvođač i potrošač komuniciraju preko reda, pa ne moraju da se izvršavaju na istoj mašini. U distribuiranom okruženju, lista u redisu može poslužiti kao red, a potrošač samo treba da ispituje ima li podataka u redu. Takođe se podržava skalabilnost klastera — kada jedna mašina padne, neće dovesti do pada celog klastera.
Rezime
Ovaj tekst je uglavnom obradio mehanizam čekanje/obaveštenje niti, uključujući upotrebu metoda wait/notify/notifyAll, kao i primer koda implementacije modela proizvođač-potrošač pomoću wait/notifyAll.
Tu su i metode await/signalAll iz klase Condition i primer koda implementacije modela proizvođač-potrošač pomoću Condition-ovih await/signalAll. Na kraju je obrađen i primer koda implementacije modela proizvođač-potrošač pomoću BlockingQueue.
Urednik: Chenmo Wang Er. Deo sadržaja potiče iz GitHub repozitorijuma CL0610 https://github.com/CL0610/Java-concurrency, a deo slika i sadržaja iz ovog posta na Zhihu-u.
