Alatke za komunikaciju u Java konkurentnom programiranju — Semaphore, Exchanger, CountDownLatch, CyclicBarrier, Phaser i druge
Lekcija 29: Alatke za komunikaciju
U JDK-u su dostupne neke uobičajene alatke za komunikaciju u konkurentnom programiranju koje možemo koristiti, kao što su CountDownLatch, Semaphore, Exchanger, CyclicBarrier, Phaser.
Sve se nalaze u JUC paketu. Najpre ćemo opširno sažeti koje to alatke postoje i kakva im je uloga, a zatim ćemo predstaviti njihove glavne načine upotrebe i principe.
| Klasa | Uloga |
|---|---|
| Semaphore | Ograničava broj niti |
| Exchanger | Razmena podataka između dve niti |
| CountDownLatch | Nit čeka dok brojač ne padne na 0, pa tek onda nastavlja rad |
| CyclicBarrier | Slično CountDownLatch-u, ali se može ponovo koristiti |
| Phaser | Unapređeni CyclicBarrier |
Semaphore
Semaphore u prevodu znači signal. Kao što ime kaže, funkcija koju ova klasa pruža jeste da više niti međusobno „šalje signale“. Taj „signal“ je podatak tipa int, koji se takođe može posmatrati kao vrsta „resursa“.
U konstruktoru se može proslediti ukupan broj početnih resursa, kao i da li se koristi „pravedan“ sinhronizator. Podrazumevano je nepravedan.
// Podrazumevano se koristi nepravedan
public Semaphore(int permits) {
sync = new NonfairSync(permits);
}
public Semaphore(int permits, boolean fair) {
sync = fair ? new FairSync(permits) : new NonfairSync(permits);
}Najvažnije metode su acquire i release. Metod acquire() traži jedan permit, dok release oslobađa jedan permit. Naravno, možete tražiti više acquire(int permits) ili oslobađati više release(int permits).
Pri svakom acquire, permits se smanjuje za jedan ili više. Ako padne na 0, svaka sledeća nit koja pozove acquire biće blokirana dok druga nit ne oslobodi permit.
Primer upotrebe Semaphore-a
Semaphore se često koristi u scenarijima sa ograničenim resursima, da bi se ograničio broj niti. Na primer, želim da ograničim da istovremeno samo 3 niti rade:
public class SemaphoreDemo {
static class MyThread implements Runnable {
private int value;
private Semaphore semaphore;
public MyThread(int value, Semaphore semaphore) {
this.value = value;
this.semaphore = semaphore;
}
@Override
public void run() {
try {
semaphore.acquire(); // dobavlja permit
System.out.println(String.format("Trenutna nit je %d, preostalo %d resursa, %d niti čeka",
value, semaphore.availablePermits(), semaphore.getQueueLength()));
// spavanje slučajnog vremena radi mešanja redosleda oslobađanja
Random random =new Random();
Thread.sleep(random.nextInt(1000));
System.out.println(String.format("Nit %d je oslobodila resurs", value));
} catch (InterruptedException e) {
e.printStackTrace();
} finally{
semaphore.release(); // oslobađa permit
}
}
}
public static void main(String[] args) {
Semaphore semaphore = new Semaphore(3);
for (int i = 0; i < 10; i++) {
new Thread(new MyThread(i, semaphore)).start();
}
}
}Izlaz:
Trenutna nit je 1, preostalo 2 resursa, 0 niti čeka Trenutna nit je 0, preostalo 1 resurs, 0 niti čeka Trenutna nit je 6, preostalo 0 resursa, 0 niti čeka Nit 6 je oslobodila resurs Trenutna nit je 2, preostalo 0 resursa, 6 niti čeka Nit 2 je oslobodila resurs Trenutna nit je 4, preostalo 0 resursa, 5 niti čeka Nit 0 je oslobodila resurs Trenutna nit je 7, preostalo 0 resursa, 4 niti čeka Nit 1 je oslobodila resurs Trenutna nit je 8, preostalo 0 resursa, 3 niti čeka Nit 7 je oslobodila resurs Trenutna nit je 5, preostalo 0 resursa, 2 niti čeka Nit 4 je oslobodila resurs Trenutna nit je 3, preostalo 0 resursa, 1 nit čeka Nit 8 je oslobodila resurs Trenutna nit je 9, preostalo 0 resursa, 0 niti čeka Nit 9 je oslobodila resurs Nit 5 je oslobodila resurs Nit 3 je oslobodila resurs
Vidimo da su u ovom pokretanju na početku resurse dobile niti 1, 0 i 6, dok su ostale ušle u red čekanja. Zatim, kada neka nit oslobodi resurs, neka nit iz reda čekanja dobije resurs.
Naravno, podrazumevani metod acquire niti uvodi u red čekanja i izbacuje izuzetak pri prekidu. Ali ima i metoda koji mogu zanemariti prekid ili ne ulaziti u blokirajući red:
// Zanemaruje prekid
public void acquireUninterruptibly()
public void acquireUninterruptibly(int permits)
// Ne ulazi u red čekanja, interno koristi CAS
public boolean tryAcquire
public boolean tryAcquire(int permits)
public boolean tryAcquire(int permits, long timeout, TimeUnit unit)
throws InterruptedException
public boolean tryAcquire(long timeout, TimeUnit unit)Princip Semaphore-a
Semaphore interno ima sinhronizator Sync koji nasleđuje AQS i prepisuje metod tryAcquireShared. U tom metodu pokušava se pribavljanje resursa.
Ako pribavljanje ne uspe (tražena količina resursa je veća od raspoložive), vraća se negativan broj (što označava neuspeh). Tada trenutna nit ulazi u red čekanja AQS-a.
Exchanger
Klasa Exchanger služi za razmenu podataka između dve niti. Podržava generičke tipove, što znači da možete prenositi bilo koje podatke između dve niti. Najpre jedan primer kako se koristi — na primer, razmena stringova između dve niti:
public class ExchangerDemo {
public static void main(String[] args) throws InterruptedException {
Exchanger<String> exchanger = new Exchanger<>();
new Thread(() -> {
try {
System.out.println("Ovo je nit A, dobila je podatke druge niti: "
+ exchanger.exchange("Ovo su podaci iz niti A"));
} catch (InterruptedException e) {
e.printStackTrace();
}
}).start();
System.out.println("U ovom trenutku nit A je blokirana, čeka podatke niti B");
Thread.sleep(1000);
new Thread(() -> {
try {
System.out.println("Ovo je nit B, dobila je podatke druge niti: "
+ exchanger.exchange("Ovo su podaci iz niti B"));
} catch (InterruptedException e) {
e.printStackTrace();
}
}).start();
}
}Izlaz:
U ovom trenutku nit A je blokirana, čeka podatke niti B Ovo je nit B, dobila je podatke druge niti: Ovo su podaci iz niti A Ovo je nit A, dobila je podatke druge niti: Ovo su podaci iz niti B
Vidimo da kada jedna nit pozove metod exchange, ulazi u blokirano stanje; tek kada i druga nit pozove exchange, prva nastavlja izvršavanje.
Pregledom izvornog koda vidi se da prebacivanje stanja čekanja ostvaruje preko park/unpark, ali pre upotrebe park/unpark radi se CAS provera, verovatno radi poboljšanja performansi.
Pošto Exchanger podržava generičke tipove, možemo prenositi bilo koje podatke, na primer IO tokove ili IO keš. Prema napomenama u JDK-u, mogu se izdvojiti sledeće osobine:
- spolja vidljive operacije ove klase su sinhronizovane;
- služi za razmenu podataka između uparenih niti;
- može se posmatrati kao obostrani sinhroni red;
- primenjiv u genetskim algoritmima, dizajnu protočnih linija i sl.
Klasa Exchanger ima i metod sa parametrom isteka vremena: ako u zadatom vremenu druga nit ne pozove exchange, izbacuje se izuzetak isteka vremena.
public V exchange(V x, long timeout, TimeUnit unit)Postavlja se pitanje: da li Exchanger može da razmenjuje podatke samo između dve niti? Šta se dešava ako tri niti pozovu exchange nad istom instancom? Odgovor je da će razmenu podataka obaviti samo prve dve niti, a treća će ući u blokirano stanje.
Treba napomenuti da se exchange može ponovo koristiti. Drugim rečima, dve niti mogu pomoću Exchanger-a neprestano razmenjivati podatke u memoriji.
CountDownLatch
Najpre da protumačimo značenje imena klase CountDownLatch. CountDown označava odbrojavanje nadole, a Latch znači „zastruga“. Neki je nazivaju i „barijera“. Uloga klase CountDownLatch se vrlo uklapa u to značenje: pretpostavimo da jedna nit pre izvršenja svog zadatka mora sačekati da druge niti završe neke preduslovne zadatke — tek kada su svi preduslovi završeni, može početi izvršenje zadatka te niti.
Metodi CountDownLatch-a su jednostavni:
// Konstruktor:
public CountDownLatch(int count)
public void await() // čeka
public boolean await(long timeout, TimeUnit unit) // čeka sa istekom vremena
public void countDown() // count - 1
public long getCount() // dobavlja preostali countPrimer CountDownLatch-a
Znamo da pri igranju igara, pre nego što igra zaista počne, obično se čeka da se završe neki preduslovi, poput „učitavanje podataka mape“, „učitavanje modela lika“, „učitavanje pozadinske muzike“ i sl. Tek kada je sve učitano, igrač zaista ulazi u igru. Hajde da simuliramo ovaj primer.
public class CountDownLatchDemo {
// Definiše nit preduslovnih zadataka
static class PreTaskThread implements Runnable {
private String task;
private CountDownLatch countDownLatch;
public PreTaskThread(String task, CountDownLatch countDownLatch) {
this.task = task;
this.countDownLatch = countDownLatch;
}
@Override
public void run() {
try {
Random random = new Random();
Thread.sleep(random.nextInt(1000));
System.out.println(task + " - zadatak završen");
countDownLatch.countDown();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
public static void main(String[] args) {
// Pretpostavimo da tri modula treba učitati
CountDownLatch countDownLatch = new CountDownLatch(3);
// Glavni zadatak
new Thread(() -> {
try {
System.out.println("Čekam na učitavanje podataka...");
System.out.println(String.format("Preostalo %d preduslovnih zadataka", countDownLatch.getCount()));
countDownLatch.await();
System.out.println("Podaci učitani, igra zvanično počinje!");
} catch (InterruptedException e) {
e.printStackTrace();
}
}).start();
// Preduslovni zadaci
new Thread(new PreTaskThread("Učitavanje podataka mape", countDownLatch)).start();
new Thread(new PreTaskThread("Učitavanje modela lika", countDownLatch)).start();
new Thread(new PreTaskThread("Učitavanje pozadinske muzike", countDownLatch)).start();
}
}Izlaz:
Čekam na učitavanje podataka... Preostalo 3 preduslovna zadatka Učitavanje modela lika - zadatak završen Učitavanje pozadinske muzike - zadatak završen Učitavanje podataka mape - zadatak završen Podaci učitani, igra zvanično počinje!
Princip CountDownLatch-a
Princip klase CountDownLatch je prilično jednostavan: interno se nalazi implementaciona klasa Sync koja nasleđuje AQS, a realizacija je vrlo jednostavna — verovatno najjednostavnija među potklasama AQS-a u JDK-u; zainteresovani čitaoci mogu pregledati izvorni kod te unutrašnje klase.
Treba napomenuti da vrednost brojača (count) u konstruktoru zapravo predstavlja broj niti koje zavrtnj treba da sačekaju. Ta vrednost se može postaviti samo jednom, a CountDownLatch ne pruža nikakav mehanizam za ponovno postavljanje te vrednosti.
CyclicBarrier
Ime CyclicBarrier razumemo kao „kružnu barijeru“. Malopre smo spomenuli da se jednom kada vrednost count CountDownLatch-a spusti na 0, ona ne može ponovo da se postavi, pa on služi samo jednokratno kao „barijera“. CyclicBarrier pak ima sve funkcije CountDownLatch-a, a dodatno omogućava resetovanje barijere metodom reset().
Ako učesnik (nit) za vreme čekanja bude prekinut jer je barijera oštećena, izbacuje se BrokenBarrierException. Metodom isBroken() može se proveriti da li je barijera oštećena.
- Ako neke niti već čekaju, poziv reset izaziva BrokenBarrierException kod niti koje čekaju. Pošto se pojavi BrokenBarrierException, čekanje više ne može da uspe.
- Ako je nit za vreme čekanja prekinuta, izbacuje se InterruptedException, koji se propagira svim ostalim nitima.
- Ako prilikom izvršenja akcije barijere nastupi izuzetak, taj izuzetak se propagira u trenutnu nit, a druge niti dobijaju BrokenBarrierException i barijera se oštećuje.
- Ako se premaši zadato vreme čekanja, trenutna nit izbacuje TimeoutException, a ostale niti dobijaju BrokenBarrierException.
Primer CyclicBarrier-a
Ponovo ćemo uzeti primer s igrama. Ako se igra sastoji od više „nivoa“, korišćenje CountDownLatch-a očigledno nije pogodno, jer bi za svaki nivo trebalo napraviti novu instancu. Umesto toga možemo koristiti CyclicBarrier da ostvarimo čekanje na učitavanje podataka za svaki nivo.
public class CyclicBarrierDemo {
static class PreTaskThread implements Runnable {
private String task;
private CyclicBarrier cyclicBarrier;
public PreTaskThread(String task, CyclicBarrier cyclicBarrier) {
this.task = task;
this.cyclicBarrier = cyclicBarrier;
}
@Override
public void run() {
// Pretpostavimo ukupno tri nivoa
for (int i = 1; i < 4; i++) {
try {
Random random = new Random();
Thread.sleep(random.nextInt(1000));
System.out.println(String.format("Zadatak nivoa %d — %s — završen", i, task));
cyclicBarrier.await();
} catch (InterruptedException | BrokenBarrierException e) {
e.printStackTrace();
}
}
}
}
public static void main(String[] args) {
CyclicBarrier cyclicBarrier = new CyclicBarrier(3, () -> {
System.out.println("Svi preduslovni zadaci nivoa završeni, igra kreće...");
});
new Thread(new PreTaskThread("Učitavanje podataka mape", cyclicBarrier)).start();
new Thread(new PreTaskThread("Učitavanje modela lika", cyclicBarrier)).start();
new Thread(new PreTaskThread("Učitavanje pozadinske muzike", cyclicBarrier)).start();
}
}Izlaz:
Zadatak nivoa 1 — Učitavanje podataka mape — završen Zadatak nivoa 1 — Učitavanje pozadinske muzike — završen Zadatak nivoa 1 — Učitavanje modela lika — završen Svi preduslovni zadaci nivoa završeni, igra kreće... Zadatak nivoa 2 — Učitavanje podataka mape — završen Zadatak nivoa 2 — Učitavanje pozadinske muzike — završen Zadatak nivoa 2 — Učitavanje modela lika — završen Svi preduslovni zadaci nivoa završeni, igra kreće... Zadatak nivoa 3 — Učitavanje modela lika — završen Zadatak nivoa 3 — Učitavanje podataka mape — završen Zadatak nivoa 3 — Učitavanje pozadinske muzike — završen Svi preduslovni zadaci nivoa završeni, igra kreće...
Obratite pažnju na razliku u odnosu na kod CountDownLatch-a. CyclicBarrier nema razdvojene await() i countDown(), već samo jedan jedini metod await().
Čim broj niti koje pozovu await postane jednak ukupnoj količini zadataka prosleđenoj u konstruktor (ovde 3), to znači da je barijera dostignuta. CyclicBarrier nam dozvoljava da pri dostizanju barijere izvršimo jedan zadatak — u konstruktor se može proslediti objekat tipa Runnable.
Gornji primer upravo pri dostizanju barijere ispisuje „Svi preduslovni zadaci nivoa završeni, igra kreće...“.
// Konstruktor
public CyclicBarrier(int parties) {
this(parties, null);
}
public CyclicBarrier(int parties, Runnable barrierAction) {
// konkretna realizacija
}Princip CyclicBarrier-a
Iako je funkcionalnost CyclicBarrier-a slična CountDownLatch-u, princip realizacije je potpuno drugačiji: CyclicBarrier interno čekanje/obaveštavanje ostvaruje pomoću Lock + Condition. Detalje možete pregledati u izvornom kodu metoda:
private int dowait(boolean timed, long nanos)Phaser
Phaser je alatka za sinhronizaciju uvedena u Javi 7, koja pruža sposobnost sinhronizacije nad dinamičkim brojem niti — za razliku od CyclicBarrier i CountDownLatch, koji zahtevaju prethodno poznavanje broja niti koje čekaju. Phaser je višefazni, što znači da može sinhronizovati više operacija u različitim fazama.
Malopre smo predstavili CyclicBarrier i uočili da se nakon prosleđivanja „ukupne količine zadataka“ parties u konstruktor ta vrednost više ne može menjati, a pri svakom pozivu await() može se potrošiti samo jedan parties. Phaser pak može dinamički da prilagođava ukupnu količinu zadataka!
Phaser je fazni, pa ima interni brojač faza. Kad god stignemo do kraja jedne faze, Phaser automatski prelazi u sledeću.
Tumačenje pojmova:
Party: u kontekstu Phasera, jedan party može biti nit ili zadatak. Kada na Phaseru registrujemo party, Phaser uvećava broj učesnika.
arrive: odgovara stanju jednog party-ja; na početku je unarrived, a pozivom
arriveAndAwaitAdvance()iliarriveAndDeregister()prelazi u stanje arrive; broj onih koji još nisu stigli može se dobiti prekogetUnarrivedParties().register: registruje novi party na Phaser.
deRegister: umanjuje jedan party.
phase: faza; kada svi registrovani party-ji stignu, pozvaće se metod
onAdvance()Phasera da bi se utvrdilo da li treba preći u sledeću fazu.
Phaser se završava na dva načina: kada su sve niti koje Phaser održava završile, ili kada onAdvance() vrati true.
Primer Phasera
Ponovo primer s igrama. Pretpostavimo da igra ima tri nivoa, ali samo prvi nivo ima tutorijal za početnike, pa treba učitati modul tutorijala. Drugi i treći nivo ga ne zahtevaju. Phaserom možemo ostvariti ovaj zahtev.
Kod:
public class PhaserDemo {
static class PreTaskThread implements Runnable {
private String task;
private Phaser phaser;
public PreTaskThread(String task, Phaser phaser) {
this.task = task;
this.phaser = phaser;
}
@Override
public void run() {
for (int i = 1; i < 4; i++) {
try {
// Od drugog nivoa nadalje preskače se učitavanje tutorijala
if (i >= 2 && "Učitavanje tutorijala za početnike".equals(task)) {
continue;
}
Random random = new Random();
Thread.sleep(random.nextInt(1000));
System.out.println(String.format("Nivo %d, treba učitati %d modula, trenutni modul [%s]",
i, phaser.getRegisteredParties(), task));
// Od drugog nivoa nadalje, preskače se učitavanje tutorijala
if (i == 1 && "Učitavanje tutorijala za početnike".equals(task)) {
System.out.println("Sledeći nivo uklanja modul [tutorijal za početnike]");
phaser.arriveAndDeregister(); // uklanja jedan modul
} else {
phaser.arriveAndAwaitAdvance();
}
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}
public static void main(String[] args) {
Phaser phaser = new Phaser(4) {
@Override
protected boolean onAdvance(int phase, int registeredParties) {
System.out.println(String.format("Priprema za nivo %d završena", phase + 1));
return phase == 3 || registeredParties == 0;
}
};
new Thread(new PreTaskThread("Učitavanje podataka mape", phaser)).start();
new Thread(new PreTaskThread("Učitavanje modela lika", phaser)).start();
new Thread(new PreTaskThread("Učitavanje pozadinske muzike", phaser)).start();
new Thread(new PreTaskThread("Učitavanje tutorijala za početnike", phaser)).start();
}
}Izlaz:
Nivo 1, treba učitati 4 modula, trenutni modul [Učitavanje pozadinske muzike] Nivo 1, treba učitati 4 modula, trenutni modul [Učitavanje tutorijala za početnike] Sledeći nivo uklanja modul [tutorijal za početnike] Nivo 1, treba učitati 3 modula, trenutni modul [Učitavanje podataka mape] Nivo 1, treba učitati 3 modula, trenutni modul [Učitavanje modela lika] Priprema za nivo 1 završena Nivo 2, treba učitati 3 modula, trenutni modul [Učitavanje podataka mape] Nivo 2, treba učitati 3 modula, trenutni modul [Učitavanje pozadinske muzike] Nivo 2, treba učitati 3 modula, trenutni modul [Učitavanje modela lika] Priprema za nivo 2 završena Nivo 3, treba učitati 3 modula, trenutni modul [Učitavanje modela lika] Nivo 3, treba učitati 3 modula, trenutni modul [Učitavanje podataka mape] Nivo 3, treba učitati 3 modula, trenutni modul [Učitavanje pozadinske muzike] Priprema za nivo 3 završena
Ovde treba obratiti pažnju na izlaz nivoa 1: u niti „Učitavanje tutorijala za početnike“ pozvan je arriveAndDeregister() kojim se umanjuje jedan party, pa naredne niti pomoću getRegisteredParties() dobijaju već izmenjenu vrednost parties. Ali u trenutnoj fazi (phase) i dalje je potrebno da svih 4 party-ja stignu da bi se okinula barijera. Tek od sledeće faze potrebno je da 3 party-ja stignu.
Klasa Phaser je vrlo korisna za kontrolu broja niti u određenoj fazi, ali ne mari za to koje tačno niti su stigli — kad god se dostigne trenutna vrednost parties, barijera se okida. Dakle, iako ovaj primer izričito koristi određenu nit (Učitavanje tutorijala za početnike) radi jasnijeg prikaza funkcionalnosti Phasera, sam Phaser nema sposobnost raspoznavanja koja je tačno nit u pitanju — njega zanima samo broj, na šta treba obratiti pažnju.
Princip Phasera
Princip klase Phaser znatno je složeniji. Interno koristi dve atomske klase zasnovane na Fork-Join okviru:
private final AtomicReference<QNode> evenQ;
private final AtomicReference<QNode> oddQ;
static final class QNode implements ForkJoinPool.ManagedBlocker {
// kod realizacije
}Zainteresovani čitaoci mogu pregledati izvorni kod JDK-a; ovde nećemo detaljnije ulaziti.
Kratak pregled
Sveukupno, CountDownLatch, CyclicBarrier i Phaser su jedan jači od drugog, ali i jedan složeniji od drugog; potrebno je razumno odabrati prema poslovnim potrebama.
Urednik: Chenmo Wang Er, deo sadržaja potiče izopen-source repozitorijuma prijatelja Xiao Qi Yinghuochong: Duboko i pristupačno o Java višenitnosti
Na GitHub-u je konačno stigao drugi PDF „Mali priručnik o konkurentnom programiranju“open-source baze znanja sa preko 17000 zvezdica „Ergov put ka naprednom Javom“! Obuhvata osnovne pojmove i načine upotrebe niti, memorijski model Java-e, synchronized, volatile, CAS, AQS, ReentrantLock, bazene niti, konkurentne kontejnere, ThreadLocal, model proizvođač-potrošač i druge teme koje su obavezne za intervjue i razvoj — ukupno preko 150 000 reči i preko 200 ručno nacrtanih ilustracija, što se može opisati kao pristupačno i sa humorom... Više detalja: Odlično, Ergov put ka naprednom konkurentnom programiranju.pdf
