CodeGym /Kursy /JAVA 25 SELF /thenCompose + niestandardowy Executor + timeouty

thenCompose + niestandardowy Executor + timeouty

JAVA 25 SELF
Poziom 55 , Lekcja 4
Dostępny

1. thenCompose vs. thenApply: różnice i kiedy czego używać

W programowaniu asynchronicznym w Javie (przez CompletableFuture) często trzeba wykonywać łańcuchy działań. Służą do tego dwa podobne metody: thenApply i thenCompose. Jednak działają one inaczej!

thenApply

Metoda thenApply jest używana, gdy kolejny krok — to prosta transformacja wartości, bez uruchamiania nowych operacji asynchronicznych. Odbiera wynik poprzedniego kroku, przetwarza go i zwraca nową wartość (nie CompletableFuture).

Jeśli znasz Stream API, to thenApply zachowuje się podobnie do map: bierze wynik, stosuje funkcję i zwraca przekształconą wersję.

Przykład:

CompletableFuture<String> cf = CompletableFuture.supplyAsync(() -> "42");
CompletableFuture<Integer> lengthFuture = cf.thenApply(s -> s.length());
// lengthFuture zawiera 2 (długość napisu "42")

Mówiąc prościej, thenApply — to sposób na powiedzenie: „kiedy wynik będzie gotowy, zrób z nim to”.

thenCompose

  • Używany, gdy następny krok — to kolejna asynchroniczna operacja (zwraca CompletableFuture).
  • Pozwala „spłaszczać” zagnieżdżone CompletableFuture (analog flatMap).
  • Jeśli użyjesz thenApply z funkcją asynchroniczną, otrzymasz CompletableFuture<CompletableFuture<T>> — niewygodne!

Przykład:

CompletableFuture<String> cf = CompletableFuture.supplyAsync(() -> "user42");

// Załóżmy, że na podstawie nazwy użytkownika chcemy pobrać jego zamówienia (asynchronicznie)
CompletableFuture<List<Order>> ordersFuture = cf.thenCompose(username -> fetchOrdersAsync(username));
// fetchOrdersAsync zwraca CompletableFuture<List<Order>>

Wizualnie:

  • thenApply: CF<String>thenApply(s -> s.length())CF<Integer>
  • thenCompose: CF<User>thenCompose(u -> fetchOrdersAsync(u.id))CF<List<Order>>

Kiedy czego używać?

  • Funkcja zwraca zwykłą wartość — użyj thenApply.
  • Funkcja zwraca CompletableFuture — użyj thenCompose.

Przykład błędu:

cf.thenApply(username -> fetchOrdersAsync(username)); // Otrzymasz CF<CF<List<Order>>>
cf.thenCompose(username -> fetchOrdersAsync(username)); // Otrzymasz CF<List<Order>>

2. Zarządzanie pulą wątków (Executor): po co i jak używać własnego Executora

Domyślnie: ForkJoinPool.commonPool()

Gdy zapisujesz CompletableFuture.supplyAsync(...) lub thenApplyAsync(...) bez podania Executor, Java używa wspólnej puli wątków — ForkJoinPool.commonPool(). To wygodne, ale nie zawsze odpowiednie:

  • Jeśli masz wiele długich lub blokujących operacji (żądania sieciowe, praca z plikami), wspólna pula może się „zapchać”, a wszystkie zadania będą czekać.
  • Czasem trzeba odizolować zadania o różnych priorytetach albo ograniczyć liczbę jednocześnie działających wątków.

Kiedy potrzebny jest własny Executor?

  • Długie, blokujące operacje (np. zapytania do bazy danych, żądania HTTP, odczyt plików).
  • Izolacja zadań: aby zadania użytkownika nie przeszkadzały systemowym.
  • Ograniczenie zasobów: np. nie uruchamiać więcej niż 10 równoległych pobrań.

Jak stworzyć własnego Executora

Zwykle używa się ThreadPoolExecutor lub fabryk z Executors:

ExecutorService myExecutor = Executors.newFixedThreadPool(10);

Jak używać własnego Executora z CompletableFuture

  • W metodach supplyAsync, runAsync, thenApplyAsync, thenComposeAsync i innych można przekazać drugi argument — waszego Executor.

Przykłady:

CompletableFuture<String> cf = CompletableFuture.supplyAsync(
    () -> loadDataFromNetwork(), myExecutor
);

cf.thenApplyAsync(data -> processData(data), myExecutor)
  .thenAcceptAsync(result -> System.out.println(result), myExecutor);

Ważne: jeśli nie podasz Executor, zostanie użyty ForkJoinPool.commonPool().

Kiedy wystarczy domyślny Executor?

  • Dla krótkich zadań CPU-bound (proste obliczenia).
  • Gdy nie ma znaczenia, w jakim wątku wykonywane jest zadanie.

3. Obsługa timeoutów: orTimeout i completeOnTimeout

Operacje asynchroniczne mogą się zawiesić albo trwać zbyt długo (np. gdy serwer nie odpowiada). Aby nie czekać w nieskończoność, w CompletableFuture są metody do pracy z timeoutami.

orTimeout

  • Kończy CompletableFuture wyjątkiem TimeoutException, jeśli operacja nie zakończy się w podanym czasie.
  • Nie anuluje faktycznie wykonywanego zadania, ale dalsza część łańcucha otrzyma błąd.

Składnia:

cf.orTimeout(3, TimeUnit.SECONDS)
  .exceptionally(ex -> {
      System.out.println("Przekroczenie czasu: " + ex);
      return null;
  });

Przykład:

CompletableFuture<String> cf = CompletableFuture.supplyAsync(() -> {
    Thread.sleep(5000); // symulujemy długą operację
    return "OK";
});

cf.orTimeout(2, TimeUnit.SECONDS)
  .exceptionally(ex -> {
      System.out.println("Błąd: " + ex);
      return "TIMEOUT";
  });

Wynik:

Po 2 sekundach zostanie rzucony TimeoutException, a exceptionally obsłuży błąd.

completeOnTimeout

  • Kończy CompletableFuture podaną wartością, jeśli operacja nie zakończy się w czasie limitu.
  • Nie zgłasza wyjątku, lecz zwraca „zapasową” wartość.

Składnia:

cf.completeOnTimeout("DEFAULT", 2, TimeUnit.SECONDS);

Przykład:

CompletableFuture<String> cf = CompletableFuture.supplyAsync(() -> {
    Thread.sleep(5000);
    return "OK";
});

cf.completeOnTimeout("TIMEOUT", 2, TimeUnit.SECONDS)
  .thenAccept(System.out::println); // Po 2 sekundach wypisze "TIMEOUT"

Porównanie orTimeout i completeOnTimeout

Metoda Co robi przy timeoutcie? Jak jest obsługiwane dalej?
orTimeout
Kończy z TimeoutException Można obsłużyć przez exceptionally/handle
completeOnTimeout
Kończy z podaną wartością thenAccept/thenApply otrzyma tę wartość

4. Praktyka: przykład z thenCompose, niestandardowym Executorem i timeoutem

Zadanie:

  • Pobrać użytkownika po id (asynchronicznie, z opóźnieniem).
  • Następnie asynchronicznie pobrać listę zamówień użytkownika (również z opóźnieniem).
  • Użyć niestandardowego Executor.
  • Dodać timeout na pobieranie zamówień.
import java.util.concurrent.*;
import java.util.*;

public class AsyncDemo {
    static ExecutorService ioExecutor = Executors.newFixedThreadPool(4);

    // Imitacja asynchronicznego pobierania użytkownika
    static CompletableFuture<String> fetchUserAsync(int userId) {
        return CompletableFuture.supplyAsync(() -> {
            sleep(1000);
            return "user" + userId;
        }, ioExecutor);
    }

    // Imitacja asynchronicznego pobierania zamówień użytkownika
    static CompletableFuture<List<String>> fetchOrdersAsync(String username) {
        return CompletableFuture.supplyAsync(() -> {
            sleep(3000); // Długa operacja!
            return List.of("order1", "order2");
        }, ioExecutor);
    }

    static void sleep(long ms) {
        try { Thread.sleep(ms); } catch (InterruptedException ignored) {}
    }

    public static void main(String[] args) {
        fetchUserAsync(42)
            .thenCompose(username ->
                fetchOrdersAsync(username)
                    .orTimeout(2, TimeUnit.SECONDS) // Timeout na pobieranie zamówień
                    .exceptionally(ex -> {
                        System.out.println("Nie udało się pobrać zamówień: " + ex);
                        return List.of();
                    })
            )
            .thenAccept(orders -> System.out.println("Zamówienia: " + orders))
            .join(); // Czekamy na zakończenie całego łańcucha

        ioExecutor.shutdown();
    }
}

Co się dzieje:

  • Pobieramy użytkownika (1 sekunda).
  • Pobieramy zamówienia (3 sekundy, ale timeout 2 sekundy).
  • Jeśli się nie wyrobimy — łapiemy TimeoutException, zwracamy pustą listę.
  • Wszystko działa na niestandardowym Executor.

Wynik:

Nie udało się pobrać zamówień: java.util.concurrent.TimeoutException
Zamówienia: []

Jeśli zmniejszysz opóźnienie w fetchOrdersAsync do 1_000 ms — zobaczysz rzeczywiste zamówienia.

5. Typowe błędy i niuanse

Błąd nr 1: Użycie thenApply zamiast thenCompose dla operacji asynchronicznych.
Jeśli funkcja zwraca CompletableFuture, a Ty zastosujesz thenApply, otrzymasz zagnieżdżony typ CompletableFuture<CompletableFuture<T>>. To skomplikuje łańcuch i spowoduje niepotrzebne „opakowania”. Rozwiązanie: użyj thenCompose, aby „spłaszczyć” wynik do CompletableFuture<T>.

Błąd nr 2: Uruchamianie długich lub zadań IO bez własnego Executor.
Domyślnie zadania są wykonywane w ForkJoinPool.commonPool(). Jeśli go przeciążysz, opóźnienia zaczną rosnąć, a inne zadania w aplikacji mogą zwolnić. Rozwiązanie: twórz własny ExecutorService i przekazuj go do supplyAsync/thenApplyAsync.

Błąd nr 3: Oczekiwanie, że orTimeout anuluje wykonywanie zadania.
orTimeout jedynie kończy CompletableFuture wyjątkiem z powodu timeoutu, ale samo zadanie nadal działa w tle. Rozwiązanie: jeśli trzeba zatrzymać wykonywanie, użyj cancel(true) lub własnych mechanizmów przerywania.

Błąd nr 4: Nieprawidłowe rozumienie zakresu działania timeoutu.
orTimeout i completeOnTimeout działają tylko dla jednego konkretnego kroku łańcucha, a nie dla całego łańcucha. Rozwiązanie: jeśli potrzebujesz ogólnego timeoutu dla całego łańcucha, owiń go w osobny CompletableFuture i zastosuj timeout do niego.

Błąd nr 5: Nie zamknięto ExecutorService.
Jeśli po wykonaniu zadań nie wywołasz shutdown()/shutdownNow() na ExecutorService, wątki będą dalej działać i program może „zawisnąć”. Rozwiązanie: zawsze zamykaj ExecutorService w finally albo używaj try-with-resources w Java 21+.

1
Ankieta/quiz
Asynchroniczne programowanie, poziom 55, lekcja 4
Niedostępny
Asynchroniczne programowanie
Asynchroniczne programowanie
Komentarze
TO VIEW ALL COMMENTS OR TO MAKE A COMMENT,
GO TO FULL VERSION