Skip to content
elephantoo

ExecutorService & CompletableFuture

Lesson 37 of 43 18 min read

Thread pools, Callable and Future, shutting down cleanly, async pipelines and virtual-thread executors.


Creating a new Thread for every task is wasteful and hard to manage: there's no limit on how many threads run, no easy way to get results back, and exceptions vanish. The java.util.concurrent package offers higher-level tools. ExecutorService manages pools of threads for you, Future represents a pending result, and CompletableFuture lets you build asynchronous pipelines. Finally, Java 21's virtual threads make "one thread per task" cheap again.

ExecutorService: submit tasks, not threads#

FirstExecutor.java
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;

public class FirstExecutor {
    public static void main(String[] args) throws InterruptedException {
        ExecutorService pool = Executors.newFixedThreadPool(3);   // 3 reusable worker threads

        for (int i = 1; i <= 6; i++) {
            int id = i;
            pool.submit(() -> {
                String name = Thread.currentThread().getName();
                System.out.println("task " + id + " on " + name.substring(name.indexOf("thread")));
            });
        }

        pool.shutdown();                                   // no new tasks; finish queued ones
        boolean finished = pool.awaitTermination(5, TimeUnit.SECONDS);
        System.out.println("all done: " + finished);
    }
}

Possible output (task order and thread assignment vary):

Output
task 1 on thread-1
task 3 on thread-3
task 2 on thread-2
task 4 on thread-1
task 5 on thread-3
task 6 on thread-2
all done: true

Six tasks ran on just three threads: each worker takes the next task from the pool's queue when it's free. You describe what to run; the executor decides where and when.

Common executors

FactoryBehaviourTypical use
newFixedThreadPool(n)exactly n threads, unbounded queueCPU-bound work (n ≈ number of cores)
newCachedThreadPool()grows and shrinks on demandmany short-lived tasks (use with care)
newSingleThreadExecutor()one thread, tasks run in orderserialising access to a resource
newScheduledThreadPool(n)delayed and periodic taskstimers, polling
newVirtualThreadPerTaskExecutor() (Java 21)a new virtual thread per taskI/O-bound work: HTTP calls, database queries

Runtime.getRuntime().availableProcessors() tells you the number of cores.

Shutting down

  • shutdown(): stop accepting tasks, let queued ones finish.
  • awaitTermination(timeout): wait for that to happen.
  • shutdownNow(): interrupt running tasks and return the ones still queued.
  • Since Java 19, ExecutorService is AutoCloseable: try (var pool = Executors.newFixedThreadPool(4)) { ... } shuts down and waits automatically at the end of the block.

Callable and Future: getting results back#

A Callable<V> is like a Runnable that returns a value (and may throw checked exceptions). Submitting one gives you a Future<V>, a handle to the result that will exist later:

Futures.java
import java.util.List;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;

public class Futures {
    static long sumRange(long from, long to) {
        long s = 0;
        for (long i = from; i <= to; i++) s += i;
        return s;
    }

    public static void main(String[] args) throws InterruptedException, ExecutionException {
        try (ExecutorService pool = Executors.newFixedThreadPool(4)) {
            Future<Long> f1 = pool.submit(() -> sumRange(1, 50_000_000));
            Future<Long> f2 = pool.submit(() -> sumRange(50_000_001, 100_000_000));
            System.out.println("submitted; main is free to do other work meanwhile");
            System.out.println("total = " + (f1.get() + f2.get()));   // get() blocks until ready

            List<Callable<String>> jobs = List.of(() -> "alpha", () -> "beta", () -> "gamma");
            for (Future<String> f : pool.invokeAll(jobs)) {           // run all, wait for all
                System.out.print(f.get() + " ");
            }
            System.out.println();

            Future<Integer> failing = pool.submit(() -> Integer.parseInt("oops"));
            try {
                failing.get();
            } catch (ExecutionException e) {
                System.out.println("task failed: " + e.getCause().getClass().getSimpleName());
            }
        }
    }
}
Output
submitted; main is free to do other work meanwhile
total = 5000000050000000
alpha beta gamma
task failed: NumberFormatException
  • future.get() waits for the result. get(2, TimeUnit.SECONDS) gives up with a TimeoutException.
  • An exception thrown inside the task is wrapped in an ExecutionException; the original is getCause().
  • invokeAll returns futures in the same order as the tasks; invokeAny returns the first successful result.
  • future.cancel(true) interrupts a running task.

The weakness of Future: you can only block and wait. You can't say "when this finishes, do that".

Scheduled tasks#

Scheduling.java
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;

public class Scheduling {
    public static void main(String[] args) throws InterruptedException {
        ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1);
        AtomicInteger ticks = new AtomicInteger();

        scheduler.schedule(() -> System.out.println("one-off after 100 ms"), 100, TimeUnit.MILLISECONDS);
        scheduler.scheduleAtFixedRate(() -> ticks.incrementAndGet(), 0, 50, TimeUnit.MILLISECONDS);

        Thread.sleep(330);
        scheduler.shutdown();
        scheduler.awaitTermination(1, TimeUnit.SECONDS);
        System.out.println("ticked several times: " + (ticks.get() >= 5));
    }
}
Output
one-off after 100 ms
ticked several times: true

CompletableFuture: asynchronous pipelines#

CompletableFuture<T> (Java 8+) is a Future you can chain and combine without blocking, much like promises in JavaScript:

Pipeline.java
import java.util.concurrent.CompletableFuture;

public class Pipeline {
    static void pause(int ms) {
        try { Thread.sleep(ms); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
    }

    static String fetchUser(int id) { pause(200); return "user" + id; }
    static double fetchBalance(String user) { pause(200); return 1500.0; }
    static double fetchFxRate() { pause(200); return 0.012; }   // INR -> USD

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

        CompletableFuture<Double> balance = CompletableFuture
                .supplyAsync(() -> fetchUser(42))               // runs on a pool thread
                .thenApply(String::toUpperCase)                 // transform the result
                .thenApply(Pipeline::fetchBalance);             // dependent step

        CompletableFuture<Double> rate = CompletableFuture.supplyAsync(Pipeline::fetchFxRate);  // in parallel

        CompletableFuture<String> report = balance
                .thenCombine(rate, (inr, fx) -> String.format("Balance: Rs %.0f = $%.2f", inr, inr * fx));

        System.out.println(report.join());                      // wait for the final result
        long elapsed = System.currentTimeMillis() - t0;
        System.out.println("parallel parts overlapped: " + (elapsed < 550));
    }
}
Output
Balance: Rs 1500 = $18.00
parallel parts overlapped: true

Key methods:

MethodUse
supplyAsync(supplier) / runAsync(runnable)start async work (common pool, or pass an executor)
thenApply(fn)transform the result (like map)
thenCompose(fn)chain another async call that returns a CompletableFuture (like flatMap)
thenCombine(other, fn)merge two independent results
thenAccept(consumer) / thenRun(runnable)consume the result / run afterwards
allOf(...) / anyOf(...)wait for all / the first of many futures
exceptionally(fn) / handle(fn)recover from failures
orTimeout(...) / completeOnTimeout(value, ...)timeouts (Java 9+)
join() / get()block for the result (join throws unchecked exceptions)

Fan-out with allOf and handling errors

FanOut.java
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class FanOut {
    static int priceFrom(String shop) {
        if (shop.equals("BrokenMart")) throw new IllegalStateException(shop + " is down");
        return switch (shop) {
            case "ShopA" -> 999;
            case "ShopB" -> 949;
            default -> 1010;
        };
    }

    public static void main(String[] args) {
        List<String> shops = List.of("ShopA", "ShopB", "BrokenMart", "ShopC");

        try (ExecutorService io = Executors.newFixedThreadPool(4)) {
            List<CompletableFuture<Integer>> quotes = shops.stream()
                    .map(s -> CompletableFuture.supplyAsync(() -> priceFrom(s), io)
                            .exceptionally(ex -> {
                                System.out.println("skipping: " + ex.getCause().getMessage());
                                return Integer.MAX_VALUE;               // fallback value
                            }))
                    .toList();

            CompletableFuture.allOf(quotes.toArray(new CompletableFuture[0])).join();

            int best = quotes.stream().mapToInt(CompletableFuture::join).min().orElseThrow();
            System.out.println("best price: " + best);
        }
    }
}
Output
skipping: BrokenMart is down
best price: 949

Without a handler, an exception inside a stage propagates down the chain, and join() throws a CompletionException wrapping it.

By default, supplyAsync uses the shared ForkJoinPool.commonPool(), sized for CPU work. For blocking I/O (HTTP, database), pass your own executor as the second argument, as above, or use virtual threads.

Virtual threads (Java 21)#

A traditional platform thread maps one-to-one onto an operating-system thread: it is expensive to create and has a large stack, so a server can run only a few thousand of them. Virtual threads are lightweight threads managed by the JVM. When a virtual thread blocks on I/O or sleep, the JVM parks it and reuses the underlying OS thread (a carrier) for other work. You can run millions of them.

Virtual.java
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.ArrayList;
import java.util.List;

public class Virtual {
    public static void main(String[] args) throws Exception {
        long t0 = System.currentTimeMillis();
        List<Future<Integer>> results = new ArrayList<>();

        try (ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor()) {
            for (int i = 0; i < 10_000; i++) {
                int id = i;
                results.add(executor.submit(() -> {
                    Thread.sleep(1000);            // simulate a slow network call
                    return id;
                }));
            }
        }   // close() waits for all tasks

        long sum = 0;
        for (Future<Integer> f : results) sum += f.get();
        long secs = (System.currentTimeMillis() - t0) / 1000;
        System.out.println("10,000 blocking tasks, sum " + sum + ", in about " + (secs <= 3 ? "1-3" : "many") + " seconds");

        Thread vt = Thread.ofVirtual().start(() -> {});
        vt.join();
        System.out.println("virtual? " + vt.isVirtual());
    }
}
Output
10,000 blocking tasks, sum 49995000, in about 1-3 seconds
virtual? true

With a fixed pool of 100 platform threads, the same work would take about 100 seconds. Guidelines:

  • Use virtual threads for I/O-bound tasks: web requests, database calls, file and network access. Spring Boot 3.2+ can serve every request on a virtual thread with spring.threads.virtual.enabled=true.
  • They don't make CPU-bound work faster; use a fixed pool sized to your cores for that.
  • Don't pool virtual threads; create one per task.
  • Keep using normal synchronisation tools. In Java 21, blocking while inside a synchronized block pins the carrier thread, so prefer ReentrantLock around long blocking operations (this limitation was removed in Java 24).

Choosing the right tool#

SituationTool
Run many independent I/O tasksvirtual-thread executor
Parallel CPU work on chunksfixed pool + Callable/invokeAll (or parallel streams)
Async steps that depend on each other or combineCompletableFuture
Periodic or delayed jobsScheduledExecutorService
Hand off work between threadsBlockingQueue (previous lesson)

Common mistakes#

  • Forgetting to shut down executors.
  • Calling future.get() immediately after submitting, which makes the work effectively sequential.
  • Ignoring exceptions inside tasks: a failed submit(Runnable) task fails silently unless you call get().
  • Running blocking I/O in the common fork/join pool.
  • Unbounded newCachedThreadPool under heavy load, which can create thousands of platform threads.
  • Pooling virtual threads, or expecting them to speed up CPU-bound code.

What's next#

Threads, stacks, heaps and pools all live inside the JVM. Next we look under the hood: JVM memory and garbage collection.

Check your understanding

Quick quiz

0/3 answered
  1. 1.What is the difference between Runnable and Callable<V>?

  2. 2.What happens if you never call shutdown() on an ExecutorService created with Executors.newFixedThreadPool(4)?

  3. 3.Which CompletableFuture method combines the results of two independent futures?

Finished reading?

Mark this lesson complete to track your progress.