AdvancedJava · Lesson 2 of 9

Concurrency: Executors & CompletableFuture

Run work in parallel with thread pools, collect results with futures, and compose async tasks.

Creating raw Thread objects by hand is error-prone. Use an ExecutorService instead: submit tasks, receive Futures, and let the pool manage threads. Since Java 21, ExecutorService is AutoCloseable, so try-with-resources waits for tasks and shuts the pool down.

CompletableFuture composes asynchronous steps: supplyAsync starts work, thenApply transforms the result, thenCombine joins two results, allOf waits for many, and exceptionally provides a fallback on failure. orTimeout fails a step that takes too long.

Use concurrency for independent work — calling several services, processing many files. Measure first: for small tasks, the coordination overhead can make parallel code slower.

Concurrency.javaJava
import java.util.List;
import java.util.concurrent.*;

public class Concurrency {
    static int slowScoreLookup(String student) {
        sleep(300);                                  // pretend: a remote call
        return Math.abs(student.hashCode()) % 101;
    }

    static void sleep(long ms) {
        try { Thread.sleep(ms); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
    }

    public static void main(String[] args) throws Exception {
        List<String> students = List.of("Amina", "Juma", "Neema", "Ali", "Rehema", "Baraka");

        long start = System.nanoTime();
        try (ExecutorService pool = Executors.newFixedThreadPool(6)) {
            List<Future<Integer>> futures = students.stream()
                .map(s -> pool.submit(() -> slowScoreLookup(s)))
                .toList();
            int total = 0;
            for (Future<Integer> f : futures) total += f.get();
            System.out.printf("Total %d in %d ms (not %d)%n",
                total, (System.nanoTime() - start) / 1_000_000, students.size() * 300);
        }

        CompletableFuture<String> profile = CompletableFuture.supplyAsync(() -> { sleep(200); return "Amina Hassan"; });
        CompletableFuture<Integer> score = CompletableFuture.supplyAsync(() -> { sleep(250); return 88; });

        String report = profile
            .thenCombine(score, (name, s) -> name + " scored " + s)
            .thenApply(String::toUpperCase)
            .get(2, TimeUnit.SECONDS);
        System.out.println(report);

        String fallback = CompletableFuture
            .supplyAsync(() -> { if (true) throw new IllegalStateException("SMS gateway down"); return "sent"; })
            .exceptionally(err -> "queued for retry (" + err.getCause().getMessage() + ")")
            .join();
        System.out.println(fallback);

        try {
            CompletableFuture.supplyAsync(() -> { sleep(1000); return "late"; })
                .orTimeout(100, TimeUnit.MILLISECONDS)
                .join();
        } catch (CompletionException e) {
            System.out.println("Timed out: " + e.getCause().getClass().getSimpleName());
        }
    }
}

Key points

  • Use ExecutorService (in try-with-resources) instead of raw threads.
  • CompletableFuture composes async steps, combines results and handles failures.
  • Always bound waiting with timeouts (get(timeout), orTimeout).

Exercise

Simulate fetching a student's profile, results and fee balance from three "services" (with different sleeps) using CompletableFuture, combine them into one Dashboard record, and fall back to a default fee balance if that service throws.

Show solution

Try the exercise yourself first — then compare your approach with this one.

Each "service" call starts as its own CompletableFuture, so all three run at the same time. thenCombine joins the results into a Dashboard record. The fee service fails, and exceptionally supplies a default balance so the dashboard still loads; orTimeout bounds each call.

DashboardLoader.javaJava
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;

public class DashboardLoader {
    record Dashboard(String name, List<Integer> results, long feeBalance) {}

    static void sleep(long ms) {
        try { Thread.sleep(ms); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
    }

    static String fetchProfile()        { sleep(300); return "Amina Hassan"; }
    static List<Integer> fetchResults() { sleep(400); return List.of(88, 79, 91); }
    static long fetchFeeBalance()       { sleep(200); throw new IllegalStateException("fee service unavailable"); }

    public static void main(String[] args) {
        long start = System.nanoTime();

        var profile = CompletableFuture.supplyAsync(DashboardLoader::fetchProfile).orTimeout(2, TimeUnit.SECONDS);
        var results = CompletableFuture.supplyAsync(DashboardLoader::fetchResults).orTimeout(2, TimeUnit.SECONDS);
        var fees = CompletableFuture.supplyAsync(DashboardLoader::fetchFeeBalance)
            .orTimeout(2, TimeUnit.SECONDS)
            .exceptionally(err -> {
                System.out.println("Fee service failed (" + err.getCause().getMessage() + "), using default");
                return -1L;
            });

        Dashboard dashboard = profile
            .thenCombine(results, (name, scores) -> new Dashboard(name, scores, 0))
            .thenCombine(fees, (d, balance) -> new Dashboard(d.name(), d.results(), balance))
            .join();

        System.out.println(dashboard);
        System.out.printf("Loaded in %d ms (the slowest call took 400 ms)%n", (System.nanoTime() - start) / 1_000_000);
    }
}

Check your understanding

  1. Why prefer an ExecutorService over creating Thread objects by hand?

  2. What does thenCombine do?

  3. What does exceptionally(err -> fallback) provide?

  4. Why put orTimeout (or get(timeout, unit)) on async calls?

Ask AI