Pametni trikovi — istovremeno programiranje sa CompletableFuture je zaista elegantno

Podsetimo se Future-a kroz primer
U nekim poslovnim scenarijima moramo koristiti više niti da asinhrono izvršimo zadatke i ubrzamo njihovo izvršavanje.
U JDK5 dodat je interfejs Future koji opisuje rezultat asinhronog izračunavanja.
Iako Future i njegove metode pružaju mogućnost asinhronog izvršavanja zadataka, dobijanje rezultata je vrlo nezgodno — moramo pozivom Future.get() blokirati pozivajuću nit, ili pak pozivom Future.isDone ispitivati u petlji da li je zadatak završen, pa tek onda preuzeti rezultat.
Nijedan od ovih načina obrade nije naročito elegantan; odgovarajući kod je ispod:
@Test
public void testFuture() throws ExecutionException, InterruptedException {
ExecutorService executorService = Executors.newFixedThreadPool(5);
Future<String> future = executorService.submit(() -> {
Thread.sleep(2000);
return "hello";
});
System.out.println(future.get());
System.out.println("end");
}Uz to, Future ne može da reši scenarije gde više asinhronih zadataka međusobno zavisi; jednostavnije rečeno, glavna nit mora da sačeka da podzadaci u drugim nitima završe pre nego što nastavi. U tom slučaju pomislite na „CountDownLatch" — zaista može da reši problem; kod je ispod.
Definišimo dva Future-a: prvi preko korisničkog ID-a dohvata podatke o korisniku, a drugi preko ID-a proizvoda podatke o proizvodu.
@Test
public void testCountDownLatch() throws InterruptedException, ExecutionException {
ExecutorService executorService = Executors.newFixedThreadPool(5);
CountDownLatch downLatch = new CountDownLatch(2);
long startTime = System.currentTimeMillis();
Future<String> userFuture = executorService.submit(() -> {
//simulacija trajanja upita 500 ms
Thread.sleep(500);
downLatch.countDown();
return "Korisnik A";
});
Future<String> goodsFuture = executorService.submit(() -> {
//simulacija trajanja upita 500 ms
Thread.sleep(400);
downLatch.countDown();
return "Proizvod A";
});
downLatch.await();
//simulacija trajanja glavnog programa
Thread.sleep(600);
System.out.println("Dohvaćeni podaci o korisniku:" + userFuture.get());
System.out.println("Dohvaćeni podaci o proizvodu:" + goodsFuture.get());
System.out.println("Ukupno vreme" + (System.currentTimeMillis() - startTime) + "ms");
}„Rezultat rada"
Dohvaćeni podaci o korisniku:Korisnik A
Dohvaćeni podaci o proizvodu:Proizvod A
Ukupno vreme1110msIz rezultata se vidi da su oba rezultata dohvaćena; da nismo koristili asinhronu obradu, vreme bi bilo 500+400+600 = 1500, a sa asinhronom obradom stvarno utrošeno je samo 1110.
Međutim, od Jave 8 nadalje ovo ne smatram elegantnim rešenjem; upoznajmo se sa upotrebom CompletableFuture-a.
Isti primer realizovan preko CompletableFuture-a
@Test
public void testCompletableInfo() throws InterruptedException, ExecutionException {
long startTime = System.currentTimeMillis();
//poziv korisničkog servisa za dohvatanje osnovnih podataka o korisniku
CompletableFuture<String> userFuture = CompletableFuture.supplyAsync(() ->
//simulacija trajanja upita 500 ms
{
try {
Thread.sleep(500);
} catch (InterruptedException e) {
e.printStackTrace();
}
return "Korisnik A";
});
//poziv servisa proizvoda za dohvatanje osnovnih podataka o proizvodu
CompletableFuture<String> goodsFuture = CompletableFuture.supplyAsync(() ->
//simulacija trajanja upita 500 ms
{
try {
Thread.sleep(400);
} catch (InterruptedException e) {
e.printStackTrace();
}
return "Proizvod A";
});
System.out.println("Dohvaćeni podaci o korisniku:" + userFuture.get());
System.out.println("Dohvaćeni podaci o proizvodu:" + goodsFuture.get());
//simulacija trajanja glavnog programa
Thread.sleep(600);
System.out.println("Ukupno vreme" + (System.currentTimeMillis() - startTime) + "ms");
}Rezultat rada
Dohvaćeni podaci o korisniku:Korisnik A
Dohvaćeni podaci o proizvodu:Proizvod A
Ukupno vreme1112msPomoću CompletableFuture-a lako se ostvaruje funkcionalnost CountDownLatch-a; ali mislite da je to sve? Daleko od toga, CompletableFuture je mnogo moćniji.
Na primer, može: zadatak 1 da se završi pa tek onda zadatak 2; čak i rezultat zadatka 1 može biti argument zadatka 2 i sl.; upoznajmo se sa API-jem CompletableFuture-a.
Načini kreiranja CompletableFuture-a
1. Često korišćena 4 načina kreiranja
U izvornom kodu CompletableFuture-a postoje četiri statičke metode za izvršavanje asinhronih zadataka:
public static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier){..}
public static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier,Executor executor){..}
public static CompletableFuture<Void> runAsync(Runnable runnable){..}
public static CompletableFuture<Void> runAsync(Runnable runnable,Executor executor){..}Obično kreiramo CompletableFuture pomoću gore navedenih statičkih metoda; objasnimo i njihove razlike:
- „supplyAsync" izvršava zadatak i podržava povratnu vrednost.
- „runAsync" izvršava zadatak i nema povratnu vrednost.
①, „metoda supplyAsync"
//Koristi podrazumevani ugrađeni bazen niti ForkJoinPool.commonPool() i na osnovu supplier-a gradi zadatak za izvršavanje
public static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier)
//Sa prilagođenim bazenom niti, na osnovu supplier-a gradi zadatak za izvršavanje
public static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier, Executor executor)②, „metoda runAsync"
//Koristi podrazumevani ugrađeni bazen niti ForkJoinPool.commonPool() i na osnovu runnable-a gradi zadatak za izvršavanje
public static CompletableFuture<Void> runAsync(Runnable runnable)
//Sa prilagođenim bazenom niti, na osnovu runnable-a gradi zadatak za izvršavanje
public static CompletableFuture<Void> runAsync(Runnable runnable, Executor executor)2. Četiri načina dobijanja rezultata
Za dobijanje rezultata klasa CompletableFuture nudi četiri načina:
//Način jedan
public T get()
//Način dva
public T get(long timeout, TimeUnit unit)
//Način tri
public T getNow(T valueIfAbsent)
//Način četiri
public T join()Objašnjenje:
- „get() i get(long timeout, TimeUnit unit)" => već su postojali u Future-u; drugi nudi rad sa istekom vremena; ako se rezultat ne dobije u zadatom roku, baciće izuzetak isteka.
- „getNow" => odmah dobija rezultat bez blokiranja; ako je izračunavanje rezultata završeno, vraće rezultat ili izuzetak nastao tokom izračunavanja; ako nije završeno, vraće zadatu vrednost valueIfAbsent.
- „join" => u metodi se ne baca izuzetak.
Primer:
@Test
public void testCompletableGet() throws InterruptedException, ExecutionException {
CompletableFuture<String> cp1 = CompletableFuture.supplyAsync(() -> {
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
e.printStackTrace();
}
return "Proizvod A";
});
// Test metode getNow
System.out.println(cp1.getNow("Proizvod B"));
//Test metode join
CompletableFuture<Integer> cp2 = CompletableFuture.supplyAsync((() -> 1 / 0));
System.out.println(cp2.join());
System.out.println("-----------------------------------------------------");
//Test metode get
CompletableFuture<Integer> cp3 = CompletableFuture.supplyAsync((() -> 1 / 0));
System.out.println(cp3.get());
}„Rezultat rada":
- Prvi rezultat je „Proizvod B", jer rezultat ne može odmah da se dobije — nit prvo spava 1 sekundu.
- Metod join ne baca izuzetak unutar metode, ali će rezultat izvršavanja baciti izuzetak; bačeni izuzetak je CompletionException.
- Metod get baca izuzetak unutar metode, a rezultat izvršavanja baca ExecutionException.
Metode asinhronog povratnog poziva

1,thenRun/thenRunAsync
Jednostavnije rečeno, „nakon prvog zadatka izvrši i drugi zadatak, a drugi zadatak takođe nema povratnu vrednost".
Primer
@Test
public void testCompletableThenRunAsync() throws InterruptedException, ExecutionException {
long startTime = System.currentTimeMillis();
CompletableFuture<Void> cp1 = CompletableFuture.runAsync(() -> {
try {
//izvrši zadatak A
Thread.sleep(600);
} catch (InterruptedException e) {
e.printStackTrace();
}
});
CompletableFuture<Void> cp2 = cp1.thenRun(() -> {
try {
//izvrši zadatak B
Thread.sleep(400);
} catch (InterruptedException e) {
e.printStackTrace();
}
});
// Test metode get
System.out.println(cp2.get());
//simulacija trajanja glavnog programa
Thread.sleep(600);
System.out.println("Ukupno vreme" + (System.currentTimeMillis() - startTime) + "ms");
}
//Rezultat rada
/**
* null
* Ukupno vreme1610ms
*/„Koja je razlika između thenRun i thenRunAsync?"
Ako pri izvršavanju prvog zadatka prosledite prilagođeni bazen niti:
- Poziv thenRun pri izvršavanju drugog zadatka znači da drugi zadatak i prvi koriste isti bazen niti.
- Poziv thenRunAsync pri izvršavanju drugog zadatka znači da prvi zadatak koristi bazen niti koji ste prosledili, a drugi zadatak koristi ForkJoin bazen niti.
Napomena: Isto važi za thenAccept i thenAcceptAsync, thenApply i thenApplyAsync i sl., koje ćemo opisati u nastavku.
2,thenAccept/thenAcceptAsync
Nakon prvog zadatka, izvršava se drugi zadatak povratnog poziva; rezultat prvog zadatka prosleđuje se kao argument metodi povratnog poziva, ali metoda povratnog poziva nema povratnu vrednost.
Primer
@Test
public void testCompletableThenAccept() throws ExecutionException, InterruptedException {
long startTime = System.currentTimeMillis();
CompletableFuture<String> cp1 = CompletableFuture.supplyAsync(() -> {
return "dev";
});
CompletableFuture<Void> cp2 = cp1.thenAccept((a) -> {
System.out.println("Rezultat prethodnog zadatka: " + a);
});
cp2.get();
}3, thenApply/thenApplyAsync
Nakon prvog zadatka, izvršava se drugi zadatak povratnog poziva; rezultat prvog zadatka prosleđuje se kao argument metodi povratnog poziva, a metoda povratnog poziva ima povratnu vrednost.
@Test
public void testCompletableThenApply() throws ExecutionException, InterruptedException {
CompletableFuture<String> cp1 = CompletableFuture.supplyAsync(() -> {
return "dev";
}).thenApply((a) -> {
if(Objects.equals(a,"dev")){
return "dev";
}
return "prod";
});
System.out.println("Trenutno okruženje:" + cp1.get());
//Izlaz: Trenutno okruženje:dev
}Povratni poziv za izuzetke
CompletableFuture, bilo da se zadatak normalno završi ili izazove izuzetak, uvek poziva funkciju povratnog poziva „whenComplete".
- „Normalno završen": whenComplete vraća rezultat isti kao i nadređeni zadatak, izuzetak je null;
- „Sa izuzetkom": whenComplete vraća rezultat null, a izuzetak je onaj iz nadređenog zadatka;
Dakle, pri pozivu get(), ako je normalno završen dobija se rezultat; ako nastupi izuzetak, baca se izuzetak koji treba obraditi.
1,Samo whenComplete
@Test
public void testCompletableWhenComplete() throws ExecutionException, InterruptedException {
CompletableFuture<Double> future = CompletableFuture.supplyAsync(() -> {
if (Math.random() < 0.5) {
throw new RuntimeException("Greška");
}
System.out.println("Normalno završeno");
return 0.11;
}).whenComplete((aDouble, throwable) -> {
if (aDouble == null) {
System.out.println("whenComplete aDouble is null");
} else {
System.out.println("whenComplete aDouble is " + aDouble);
}
if (throwable == null) {
System.out.println("whenComplete throwable is null");
} else {
System.out.println("whenComplete throwable is " + throwable.getMessage());
}
});
System.out.println("Konačan rezultat = " + future.get());
}Pri normalnom završetku, bez izuzetka:
Normalno završeno
whenComplete aDouble is 0.11
whenComplete throwable is null
Konačan rezultat = 0.11Pri izuzetku: get() baca izuzetak
whenComplete aDouble is null
whenComplete throwable is java.lang.RuntimeException: Greška
java.util.concurrent.ExecutionException: java.lang.RuntimeException: Greška
at java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:357)
at java.util.concurrent.CompletableFuture.get(CompletableFuture.java:1895)2,whenComplete + exceptionally primer
@Test
public void testWhenCompleteExceptionally() throws ExecutionException, InterruptedException {
CompletableFuture<Double> future = CompletableFuture.supplyAsync(() -> {
if (Math.random() < 0.5) {
throw new RuntimeException("Greška");
}
System.out.println("Normalno završeno");
return 0.11;
}).whenComplete((aDouble, throwable) -> {
if (aDouble == null) {
System.out.println("whenComplete aDouble is null");
} else {
System.out.println("whenComplete aDouble is " + aDouble);
}
if (throwable == null) {
System.out.println("whenComplete throwable is null");
} else {
System.out.println("whenComplete throwable is " + throwable.getMessage());
}
}).exceptionally((throwable) -> {
System.out.println("Izuzetak u exceptionally: " + throwable.getMessage());
return 0.0;
});
System.out.println("Konačan rezultat = " + future.get());
}Kada nastupi izuzetak, exceptionally hvata taj izuzetak i vraća podrazumevanu vrednost 0.0.
whenComplete aDouble is null
whenComplete throwable is java.lang.RuntimeException: Greška
Izuzetak u exceptionally: java.lang.RuntimeException: Greška
Konačan rezultat = 0.0Povratni poziv za kombinaciju više zadataka

1,AND odnos kombinacije
thenCombine / thenAcceptBoth / runAfterBoth znače: „kad i zadatak jedan i zadatak dva završe, izvrši zadatak tri".
Razlika je u sledećem:
- „runAfterBoth" ne prosleđuje rezultat izvršavanja kao argument metode i nema povratnu vrednost
- „thenAcceptBoth": rezultate oba zadatka prosleđuje kao argumente metode u navedeni metod i nema povratnu vrednost
- „thenCombine": rezultate oba zadatka prosleđuje kao argumente metode u navedeni metod i ima povratnu vrednost
@Test
public void testCompletableThenCombine() throws ExecutionException, InterruptedException {
//Kreira bazen niti
ExecutorService executorService = Executors.newFixedThreadPool(10);
//Pokreni asinhroni zadatak 1
CompletableFuture<Integer> task = CompletableFuture.supplyAsync(() -> {
System.out.println("Asinhroni zadatak 1, trenutna nit je: " + Thread.currentThread().getId());
int result = 1 + 1;
System.out.println("Asinhroni zadatak 1 završen");
return result;
}, executorService);
//Pokreni asinhroni zadatak 2
CompletableFuture<Integer> task2 = CompletableFuture.supplyAsync(() -> {
System.out.println("Asinhroni zadatak 2, trenutna nit je: " + Thread.currentThread().getId());
int result = 1 + 1;
System.out.println("Asinhroni zadatak 2 završen");
return result;
}, executorService);
//Kombinovanje zadataka
CompletableFuture<Integer> task3 = task.thenCombineAsync(task2, (f1, f2) -> {
System.out.println("Izvršavanje zadatka 3, trenutna nit je: " + Thread.currentThread().getId());
System.out.println("Rezultat zadatka 1: " + f1);
System.out.println("Rezultat zadatka 2: " + f2);
return f1 + f2;
}, executorService);
Integer res = task3.get();
System.out.println("Konačan rezultat: " + res);
}„Rezultat rada"
Asinhroni zadatak 1, trenutna nit je: 17
Asinhroni zadatak 1 završen
Asinhroni zadatak 2, trenutna nit je: 18
Asinhroni zadatak 2 završen
Izvršavanje zadatka 3, trenutna nit je: 19
Rezultat zadatka 1: 2
Rezultat zadatka 2: 2
Konačan rezultat: 42,OR odnos kombinacije
applyToEither / acceptEither / runAfterEither znače: „od dva zadatka, čim jedan završi, izvrši zadatak tri".
Razlika je u sledećem:
- „runAfterEither": ne prosleđuje rezultat izvršavanja kao argument metode i nema povratnu vrednost
- „acceptEither": rezultat već završenog zadatka prosleđuje kao argument metode u navedeni metod i nema povratnu vrednost
- „applyToEither": rezultat već završenog zadatka prosleđuje kao argument metode u navedeni metod i ima povratnu vrednost
Primer
@Test
public void testCompletableEitherAsync() {
//Kreira bazen niti
ExecutorService executorService = Executors.newFixedThreadPool(10);
//Pokreni asinhroni zadatak 1
CompletableFuture<Integer> task = CompletableFuture.supplyAsync(() -> {
System.out.println("Asinhroni zadatak 1, trenutna nit je: " + Thread.currentThread().getId());
int result = 1 + 1;
System.out.println("Asinhroni zadatak 1 završen");
return result;
}, executorService);
//Pokreni asinhroni zadatak 2
CompletableFuture<Integer> task2 = CompletableFuture.supplyAsync(() -> {
System.out.println("Asinhroni zadatak 2, trenutna nit je: " + Thread.currentThread().getId());
int result = 1 + 2;
try {
Thread.sleep(3000);
} catch (InterruptedException e) {
e.printStackTrace();
}
System.out.println("Asinhroni zadatak 2 završen");
return result;
}, executorService);
//Kombinovanje zadataka
task.acceptEitherAsync(task2, (res) -> {
System.out.println("Izvršavanje zadatka 3, trenutna nit je: " + Thread.currentThread().getId());
System.out.println("Rezultat prethodnog zadatka: " + res);
}, executorService);
}Rezultat rada
//Iz rezultata se vidi da asinhroni zadatak 2 nije ni završio; zadatak 3 dohvata rezultat zadatka 1
Asinhroni zadatak 1, trenutna nit je: 17
Asinhroni zadatak 1 završen
Asinhroni zadatak 2, trenutna nit je: 18
Izvršavanje zadatka 3, trenutna nit je: 19
Rezultat prethodnog zadatka: 2Napomena
Ako gorešnji broj jezgara niti promenite u 1, odnosno
ExecutorService executorService = Executors.newFixedThreadPool(1);Rezultat rada će biti sledeći — videćete da zadatak 3 uopšte nije izvršen; očigledno je zadatak 3 direktno odbačen.
Asinhroni zadatak 1, trenutna nit je: 17
Asinhroni zadatak 1 završen
Asinhroni zadatak 2, trenutna nit je: 173,Kombinacija više zadataka
- „allOf": čeka da se svi zadaci završe
- „anyOf": čim se jedan zadatak završi
Primer
allOf: čeka da se svi zadaci završe
@Test
public void testCompletableAallOf() throws ExecutionException, InterruptedException {
//Kreira bazen niti
ExecutorService executorService = Executors.newFixedThreadPool(10);
//Pokreni asinhroni zadatak 1
CompletableFuture<Integer> task = CompletableFuture.supplyAsync(() -> {
System.out.println("Asinhroni zadatak 1, trenutna nit je: " + Thread.currentThread().getId());
int result = 1 + 1;
System.out.println("Asinhroni zadatak 1 završen");
return result;
}, executorService);
//Pokreni asinhroni zadatak 2
CompletableFuture<Integer> task2 = CompletableFuture.supplyAsync(() -> {
System.out.println("Asinhroni zadatak 2, trenutna nit je: " + Thread.currentThread().getId());
int result = 1 + 2;
try {
Thread.sleep(3000);
} catch (InterruptedException e) {
e.printStackTrace();
}
System.out.println("Asinhroni zadatak 2 završen");
return result;
}, executorService);
//Pokreni asinhroni zadatak 3
CompletableFuture<Integer> task3 = CompletableFuture.supplyAsync(() -> {
System.out.println("Asinhroni zadatak 3, trenutna nit je: " + Thread.currentThread().getId());
int result = 1 + 3;
try {
Thread.sleep(4000);
} catch (InterruptedException e) {
e.printStackTrace();
}
System.out.println("Asinhroni zadatak 3 završen");
return result;
}, executorService);
//Kombinovanje zadataka
CompletableFuture<Void> allOf = CompletableFuture.allOf(task, task2, task3);
//Čeka da se svi zadaci završe
allOf.get();
//Dohvatanje rezultata zadataka
System.out.println("Rezultat task: " + task.get());
System.out.println("Rezultat task2: " + task2.get());
System.out.println("Rezultat task3: " + task3.get());
}anyOf: čim se jedan zadatak završi
@Test
public void testCompletableAnyOf() throws ExecutionException, InterruptedException {
//Kreira bazen niti
ExecutorService executorService = Executors.newFixedThreadPool(10);
//Pokreni asinhroni zadatak 1
CompletableFuture<Integer> task = CompletableFuture.supplyAsync(() -> {
int result = 1 + 1;
return result;
}, executorService);
//Pokreni asinhroni zadatak 2
CompletableFuture<Integer> task2 = CompletableFuture.supplyAsync(() -> {
int result = 1 + 2;
return result;
}, executorService);
//Pokreni asinhroni zadatak 3
CompletableFuture<Integer> task3 = CompletableFuture.supplyAsync(() -> {
int result = 1 + 3;
return result;
}, executorService);
//Kombinovanje zadataka
CompletableFuture<Object> anyOf = CompletableFuture.anyOf(task, task2, task3);
//Čim se jedan zadatak završi
Object o = anyOf.get();
System.out.println("Rezultat završenog zadatka: " + o);
}Na šta treba paziti pri upotrebi CompletableFuture-a

Dok nam CompletableFuture olakšava asinhrono programiranje i čini kod elegantnijim, treba obratiti pažnju i na neke stvari pri njegovoj upotrebi.
1,Future mora dohvatiti povratnu vrednost da bi se dobio podatak o izuzetku
@Test
public void testWhenCompleteExceptionally() {
CompletableFuture<Double> future = CompletableFuture.supplyAsync(() -> {
if (1 == 1) {
throw new RuntimeException("Greška");
}
return 0.11;
});
//Ako se ne doda poziv get(), ne vidi se informacija o izuzetku
//future.get();
}Future mora dohvatiti povratnu vrednost da bi se dobio podatak o izuzetku. Bez poziva get()/join() ne vidi se informacija o izuzetku.
Pazite na to pri upotrebi; razmislite da li treba dodati try...catch... ili koristiti metod exceptionally.
2,Metod get() CompletableFuture-a je blokirajući
Metod get() CompletableFuture-a je blokirajući; ako se koristi za dobijanje povratne vrednosti asinhronog poziva, potrebno je dodati vreme isteka.
//Protivprimer
CompletableFuture.get();
//Pravilan primer
CompletableFuture.get(5, TimeUnit.SECONDS);3,Ne preporučuje se upotreba podrazumevanog bazena niti
U kodu CompletableFuture se koristi podrazumevani „ForkJoin bazen niti", čiji je broj niti „broj jezgara CPU-a umanjen za 1". Pri velikom broju zahteva, ako je obrada složena, odziv može biti spor. Generalno se preporučuje korišćenje prilagođenog bazena niti i optimizacija njegovih parametara.
4,Pri prilagođenom bazenu niti, obratite pažnju na strategiju zasićenja
Metod get() CompletableFuture-a je blokirajući i uglavnom se preporučuje future.get(5, TimeUnit.SECONDS). Takođe se preporučuje i prilagođeni bazen niti.
Međutim, ako je strategija odbijanja bazena niti DiscardPolicy ili DiscardOldestPolicy, kada se bazen niti zasiti, zadatak se direktno odbacuje bez bacanja izuzetka. Stoga se preporučuje da strategija bazena niti CompletableFuture-a bude AbortPolicy, a za dugotrajne asinhrone niti dobro obezbediti izolaciju bazena niti.
Zahvalnica
- https://www.cnblogs.com/wtzbk/p/15953570.html
- https://www.cnblogs.com/cyan-orange/p/15339145.html
- https://www.cnblogs.com/ludongguoa/p/15316488.html
Reference: obrada: Marko Marković
