CompletableFuture: Async Pipelines Without Callback Hell
The order-summary page went live on a Friday. By Monday it was the slowest page on the site. By Black Friday, it was the reason the site fell over.
The page needed three things: the customer's profile (user service, ~800 ms), their recent orders (order service, ~1200 ms), and live pricing. The first version called them one after another on the request thread: 800 ms, then 1200 ms, then pricing — over two seconds of pure waiting, most of it with the CPU idle. The "fix" shipped the next sprint: wrap each call in CompletableFuture.supplyAsync, combine the results, done. Latency dropped, everyone celebrated — until Black Friday traffic arrived and the whole service wedged. Not just the summary page: unrelated endpoints stalled too, including batch jobs that used parallel streams. The post-mortem found two compounding mistakes: a join() inside a stage that waited on another future (blocking a shared worker thread), and every stage running on the default thread pool — the one shared with the rest of the JVM.
This post is the CompletableFuture toolkit that would have prevented both the slow page and the outage: starting async work, chaining stages without nested callbacks, combining and racing futures, handling failures, enforcing deadlines, and — the part most tutorials skip — choosing who runs your stages. Every example below was compiled with javac and run on JDK 21, and every output block is real. One honesty note up front: the sandbox these ran on was heavily loaded (load average ~8 on 2 CPUs), so absolute timings overshoot the Thread.sleep values — the comparisons between runs are the evidence, not the raw milliseconds.
Starting async work: supplyAsync and runAsync
Two factory methods start everything. supplyAsync runs a Supplier that returns a value; runAsync runs a Runnable that returns nothing. Both return a CompletableFuture immediately, and both run the task on another thread while your thread continues:
import java.util.concurrent.*;
public class P7_Start {
public static void main(String[] args) {
// supplyAsync: runs a task that RETURNS a value
CompletableFuture<Integer> value = CompletableFuture.supplyAsync(() -> {
System.out.println(" supplyAsync runs on " + Thread.currentThread().getName());
return 40 + 2;
});
// runAsync: runs a task with NO return value
CompletableFuture<Void> done = CompletableFuture.runAsync(() ->
System.out.println(" runAsync runs on " + Thread.currentThread().getName()));
System.out.println("answer: " + value.join());
done.join();
System.out.println("main thread: " + Thread.currentThread().getName());
}
}
runAsync runs on Thread-1
supplyAsync runs on Thread-0
answer: 42
main thread: main
Wait — Thread-0? Not a pool thread with a sensible name? Here is the first thing tutorials gloss over: when you don't pass an executor, the default is the shared ForkJoinPool.commonPool() — but only if that pool is big enough to be useful. This sandbox has 2 CPUs, so the common pool's parallelism is 1, and the JDK falls back to spawning a fresh thread per async stage instead. On a typical machine (parallelism > 1), the same program runs its stages on pool workers:
$ java -Djava.util.concurrent.ForkJoinPool.common.parallelism=4 P7_Start
runAsync runs on ForkJoinPool.commonPool-worker-2
supplyAsync runs on ForkJoinPool.commonPool-worker-1
answer: 42
main thread: main
Both outputs are real runs of the same class — the only difference is the system property. Remember this: the default executor is a shared, JVM-wide resource, sized for CPU-bound work (availableProcessors() - 1 workers). Treat it like the office kitchen — fine for everyone to use briefly, catastrophic if someone camps in it. We'll come back to exactly how that camping happens in the executors section.
A note on join() vs get(): join() waits for the result and wraps any failure in an unchecked CompletionException; get() throws checked InterruptedException and ExecutionException. In application code join() keeps pipelines readable — you'll see it throughout this post.
Chaining without callback hell: thenApply vs thenCompose
A stage transforms the result of the previous one. The two chaining methods look interchangeable until they bite you:
thenApply(fn)—fntakes the value and returns a plain value. Use it for pure transformations.thenCompose(fn)—fntakes the value and returns another CompletableFuture. The futures are flattened into one. Use it when the next step is itself async (another service call).
import java.util.concurrent.*;
public class P0_Chain {
static String userName(String id) { return "alice"; }
static CompletableFuture<String> userEmail(String name) {
return CompletableFuture.supplyAsync(() -> name + "@example.com");
}
public static void main(String[] args) {
// thenApply: the function returns a plain VALUE
String upper = CompletableFuture.supplyAsync(() -> userName("u-42"))
.thenApply(String::toUpperCase)
.join();
System.out.println("thenApply: " + upper);
// thenCompose: the function returns a CompletableFuture -> flattened
String email = CompletableFuture.supplyAsync(() -> userName("u-42"))
.thenCompose(P0_Chain::userEmail)
.join();
System.out.println("thenCompose: " + email);
// thenApply with a future-returning function NESTS instead of flattening
CompletableFuture<CompletableFuture<String>> nested =
CompletableFuture.supplyAsync(() -> userName("u-42"))
.thenApply(P0_Chain::userEmail);
System.out.println("nested type: " + nested.getClass().getSimpleName()
+ " -> join().join() = " + nested.join().join());
}
}
thenApply: ALICE
thenCompose: alice@example.com
nested type: CompletableFuture -> join().join() = alice@example.com
The third case is the classic bug: someone writes thenApply around a method that returns a future, gets a CompletableFuture<CompletableFuture<String>>, and "fixes" it with join().join() — blocking a thread to unwrap what thenCompose would have flattened for free. Principle: a stage that returns a value needs thenApply; a stage that returns a future needs thenCompose. And notice what neither example needed: no nested callbacks, no manual latches, no shared mutable state. The pipeline reads top to bottom, and each stage's input and output types are checked by the compiler.
Combining independent results: thenCombine
Back to the order-summary page. The profile fetch and the orders fetch don't depend on each other — they're independent — so they should overlap. thenCombine waits for two futures and merges their results with a BiFunction:
import java.util.Arrays;
import java.util.concurrent.*;
public class P1_Pipeline {
record User(String id, String name) {}
record Order(String id, double total) {}
static User fetchUser() throws InterruptedException {
System.out.println(" fetchUser on " + Thread.currentThread().getName());
Thread.sleep(800);
return new User("u-42", "Alice");
}
static Order[] fetchOrders() throws InterruptedException {
System.out.println(" fetchOrders on " + Thread.currentThread().getName());
Thread.sleep(1200);
return new Order[]{ new Order("o-1", 39.98), new Order("o-2", 12.50) };
}
static <T> T sneaky(CallableThrows<T> c) {
try { return c.call(); }
catch (InterruptedException e) { throw new RuntimeException(e); }
}
interface CallableThrows<T> { T call() throws InterruptedException; }
public static void main(String[] args) {
// 1. Sequential: one blocking call after the other
long t0 = System.nanoTime();
User u = sneaky(P1_Pipeline::fetchUser);
Order[] os = sneaky(P1_Pipeline::fetchOrders);
double total = Arrays.stream(os).mapToDouble(Order::total).sum();
long seq = (System.nanoTime() - t0) / 1_000_000;
System.out.println("sequential: " + u.name() + ", " + os.length
+ " orders, total $" + total + ", wall = " + seq + " ms\n");
// 2. Pipeline: both fetches run concurrently, thenCombine merges them
long t1 = System.nanoTime();
CompletableFuture<User> uf =
CompletableFuture.supplyAsync(() -> sneaky(P1_Pipeline::fetchUser));
CompletableFuture<Order[]> of =
CompletableFuture.supplyAsync(() -> sneaky(P1_Pipeline::fetchOrders));
String summary = uf.thenCombine(of, (user, orders) -> {
double t = Arrays.stream(orders).mapToDouble(Order::total).sum();
return user.name() + " has " + orders.length + " orders, total $" + t;
}).join();
long par = (System.nanoTime() - t1) / 1_000_000;
System.out.println("pipeline: " + summary + ", wall = " + par + " ms");
}
}
fetchUser on main
fetchOrders on main
sequential: Alice, 2 orders, total $52.48, wall = 2216 ms
fetchUser on Thread-0
fetchOrders on Thread-1
pipeline: Alice has 2 orders, total $52.48, wall = 1327 ms
The sequential version pays 800 ms + 1200 ms (plus scheduling overhead on this loaded box); the pipeline pays roughly the slower of the two, because the fetches overlap — 2216 ms down to 1327 ms for the same answer. The thread names confirm the overlap is real: the two fetches ran on different threads while main waited exactly once, at the final join().
thenCombine is for independent futures merged by a function. Its siblings cover the other shapes: thenAcceptBoth when you need both results but return nothing, and runAfterBoth when you only need to know both finished. The decision is always the same: dependent steps chain with thenCompose; independent steps combine with thenCombine. Getting this wrong — chaining two independent calls with thenCompose — silently serializes them and throws away the overlap you just measured.
Fan-out and fan-in: allOf and anyOf
Two futures are a special case. When you fan out to N parallel calls — inventory checks across five regions, say — you need allOf (wait for every one) or anyOf (first one wins):
import java.util.List;
import java.util.concurrent.*;
import java.util.stream.Stream;
public class P3_FanOut {
static String fetch(String region, long ms) throws InterruptedException {
Thread.sleep(ms);
return region + "-inventory";
}
@SuppressWarnings("unchecked")
public static void main(String[] args) {
String[] regions = {"us-east", "us-west", "eu", "apac", "sa"};
long[] lat = {1500, 900, 1800, 600, 1200};
// allOf: wait for EVERY fetch, then collect results in original order
long t0 = System.nanoTime();
CompletableFuture<String>[] futures = new CompletableFuture[regions.length];
for (int i = 0; i < regions.length; i++) {
final int k = i;
futures[i] = CompletableFuture.supplyAsync(() -> {
try { return fetch(regions[k], lat[k]); }
catch (InterruptedException e) { throw new RuntimeException(e); }
});
}
List<String> results = CompletableFuture.allOf(futures)
.thenApply(v -> Stream.of(futures).map(CompletableFuture::join).toList())
.join();
long wall = (System.nanoTime() - t0) / 1_000_000;
System.out.println("allOf: " + results);
System.out.println("allOf wall = " + wall + " ms (slowest fetch was 1800 ms)\n");
// anyOf: FIRST completed future wins
long t1 = System.nanoTime();
CompletableFuture<String>[] racers = new CompletableFuture[regions.length];
for (int i = 0; i < regions.length; i++) {
final int k = i;
racers[i] = CompletableFuture.supplyAsync(() -> {
try { return fetch(regions[k], lat[k]); }
catch (InterruptedException e) { throw new RuntimeException(e); }
});
}
Object winner = CompletableFuture.anyOf(racers).join();
long wall2 = (System.nanoTime() - t1) / 1_000_000;
System.out.println("anyOf winner: " + winner + ", wall = " + wall2 + " ms");
}
}
allOf: [us-east-inventory, us-west-inventory, eu-inventory, apac-inventory, sa-inventory]
allOf wall = 1924 ms (slowest fetch was 1800 ms)
anyOf winner: apac-inventory, wall = 617 ms
Two details worth memorizing. First, allOf returns CompletableFuture<Void> — it tells you when, not what. The idiom above (thenApply + join each future in order) is how you get the results back, in the original order, with no extra waiting — every future is already complete by then. Second, anyOf returns CompletableFuture<Object>, so you cast the winner. The measured walls tell the story: allOf waits for the slowest (1924 ms ≈ the 1800 ms fetch), anyOf returns with the fastest (617 ms ≈ the 600 ms fetch). And the losers of an anyOf race keep running in the background — anyOf doesn't cancel them. If the losers hold connections or cost money per call, cancel them explicitly with future.cancel(true).
When a stage fails: exceptionally, handle, and whenComplete
Exceptions in a pipeline don't propagate to the caller — they're captured in the future and rethrown wrapped in CompletionException when you join(). Three methods deal with that, and they are not interchangeable:
exceptionally(fn)— failure becomes a fallback value; success passes through untouched. You never see the success value insidefn.handle(fn)—fnreceives both the value and the exception, and decides the result either way. This is the only one that can distinguish success from failure and transform the value.whenComplete(action)— observes the outcome (logging, metrics) but cannot change the result or swallow the error. The same value/exception continues downstream.
import java.util.concurrent.*;
public class P2_Errors {
static String fetchPricing(String sku) {
if (sku.equals("BAD-SKU")) throw new IllegalArgumentException("no price for " + sku);
return sku + ": $49.99";
}
public static void main(String[] args) {
// exceptionally: failure -> fallback VALUE; success passes through untouched
String r1 = CompletableFuture.supplyAsync(() -> fetchPricing("BAD-SKU"))
.exceptionally(ex -> {
System.out.println(" exceptionally saw "
+ ex.getClass().getSimpleName() + ": " + ex.getMessage());
return "BAD-SKU: price unavailable (fallback)";
})
.join();
System.out.println("recovered: " + r1 + "\n");
// handle: sees BOTH the value and the error; decides the result either way
String ok = CompletableFuture.supplyAsync(() -> fetchPricing("BOOK-007"))
.handle((price, ex) -> ex == null ? "OK -> " + price : "FAILED -> " + ex.getMessage())
.join();
String bad = CompletableFuture.supplyAsync(() -> fetchPricing("BAD-SKU"))
.handle((price, ex) -> ex == null ? "OK -> " + price : "FAILED -> " + ex.getMessage())
.join();
System.out.println(ok);
System.out.println(bad + "\n");
// whenComplete: observe the outcome, CANNOT change the result
String r4 = CompletableFuture.supplyAsync(() -> fetchPricing("BOOK-007"))
.whenComplete((price, ex) ->
System.out.println(" whenComplete: value=" + price + ", error=" + ex))
.join();
System.out.println("passthrough: " + r4);
}
}
exceptionally saw CompletionException: java.lang.IllegalArgumentException: no price for BAD-SKU
recovered: BAD-SKU: price unavailable (fallback)
OK -> BOOK-007: $49.99
FAILED -> java.lang.IllegalArgumentException: no price for BAD-SKU
whenComplete: value=BOOK-007: $49.99, error=null
passthrough: BOOK-007: $49.99
Notice the wrapping: exceptionally and handle receive the IllegalArgumentException wrapped in a CompletionException — always unwrap with getCause() before logging or matching on the exception type, or your logs will blame CompletionException for everything. Also notice whenComplete's output: it saw the value, printed it, and the original "BOOK-007: $49.99" flowed through unchanged.
Decision rule: use exceptionally for a simple fallback ("price unavailable, use cached"), handle when the recovery depends on which failure happened or you must record success and failure uniformly, and whenComplete for side effects only — logging, metrics, closing resources. A common production pattern is all three in sequence: whenComplete to record the outcome metric, exceptionally to substitute the fallback, and the caller never sees the exception at all.
Deadlines: orTimeout and completeOnTimeout
An async pipeline without deadlines is a pipeline that can hang forever — one slow downstream service and your join() never returns. Two methods set deadlines, with opposite philosophies:
orTimeout(duration)— slowness is a failure: the future completes exceptionally withTimeoutException(wrapped inCompletionExceptionatjoin()).completeOnTimeout(fallback, duration)— slowness gets a fallback: the future completes normally with your default value instead.
import java.util.concurrent.*;
public class P4_Timeouts {
static String slowService() {
try { Thread.sleep(2000); }
catch (InterruptedException e) { throw new RuntimeException(e); }
return "slow-result";
}
public static void main(String[] args) {
// orTimeout: the stage FAILS if it takes too long
try {
String r = CompletableFuture.supplyAsync(P4_Timeouts::slowService)
.orTimeout(300, TimeUnit.MILLISECONDS)
.join();
System.out.println("unexpected: " + r);
} catch (CompletionException ce) {
System.out.println("orTimeout -> " + ce.getCause().getClass().getName());
}
// completeOnTimeout: too slow -> fall back to a DEFAULT value
String r2 = CompletableFuture.supplyAsync(P4_Timeouts::slowService)
.completeOnTimeout("cached-result", 300, TimeUnit.MILLISECONDS)
.join();
System.out.println("completeOnTimeout -> " + r2);
}
}
orTimeout -> java.util.concurrent.TimeoutException
completeOnTimeout -> cached-result
Both fired at 300 ms against a 2000 ms stage — the deadlines work. Two caveats that matter in production. First, the timed-out task keeps running in the background: neither method cancels the underlying work (that 2-second sleep still finished on its thread). If the slow task holds a connection, cancel it explicitly or design the task to be interruptible. Second, timeouts compose per stage, not per pipeline: put the deadline on the stage that talks to the slow dependency, and choose the duration from that dependency's latency budget — not from a round number that felt right. Decision rule: orTimeout when a slow dependency is an error the caller must know about; completeOnTimeout when you have a safe degraded answer (cached value, default, empty list).
Who runs your stages: executors and the common-pool hazard
Every supplyAsync, runAsync, thenApplyAsync and friends accept an optional Executor. When you omit it, you get the shared common pool (or the JDK's fallback, as we saw). When you pass one, you get isolation: your blocking I/O stages can never starve someone else's CPU-bound work. Here's the same pipeline on a named pool:
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
public class P5_Executor {
static ExecutorService checkoutPool() {
AtomicInteger n = new AtomicInteger(1);
return Executors.newFixedThreadPool(4, r -> {
Thread t = new Thread(r, "checkout-pool-" + n.getAndIncrement());
t.setDaemon(true);
return t;
});
}
public static void main(String[] args) {
ExecutorService pool = checkoutPool();
// No executor given -> the shared common pool (or the JDK fallback)
String d = CompletableFuture.supplyAsync(
() -> "ran on " + Thread.currentThread().getName()).join();
System.out.println("default executor: " + d);
// Explicit executor -> your own named pool
String c = CompletableFuture.supplyAsync(
() -> "ran on " + Thread.currentThread().getName(), pool).join();
System.out.println("custom executor: " + c);
System.out.println("commonPool parallelism: "
+ ForkJoinPool.getCommonPoolParallelism()
+ " (CPUs on this machine: "
+ Runtime.getRuntime().availableProcessors() + ")");
pool.shutdown();
}
}
default executor: ran on Thread-0
custom executor: ran on checkout-pool-1
commonPool parallelism: 1 (CPUs on this machine: 2)
The named ThreadFactory is not vanity — when a Black Friday post-mortem starts, thread names are how you tell your pool's threads from everyone else's in a thread dump. Name every pool you create.
Now the hazard — the one from the opening story. The common pool is sized for CPU-bound work: availableProcessors() - 1 workers. If your stages block (waiting on I/O, on a latch, or join()ing another future), each blocked stage parks a worker. Park all of them, and nothing that needs the pool can ever run again — including the stage that would have unblocked them. I reproduced it deterministically: park every common-pool worker in a blocking await(), then submit the stage that would count the latch down:
import java.util.concurrent.*;
/**
* Park EVERY common-pool worker in a blocking call; the stage that would
* release them then has no thread to run on. Run with parallelism=4 so the
* pool is actually used (see text for why this matters).
*/
public class P6_PoolHazard {
public static void main(String[] args) {
int workers = ForkJoinPool.getCommonPoolParallelism();
System.out.println("common pool workers: " + workers);
CountDownLatch gate = new CountDownLatch(1);
// Park every common-pool worker inside a blocking await()...
for (int i = 0; i < workers; i++) {
CompletableFuture.runAsync(() -> {
try { gate.await(); }
catch (InterruptedException e) { Thread.currentThread().interrupt(); }
}, ForkJoinPool.commonPool());
}
// ...give the parked tasks a moment to grab their threads
try { Thread.sleep(500); }
catch (InterruptedException e) { Thread.currentThread().interrupt(); }
// This stage is pinned to the SAME shared pool -- but every worker is parked.
CompletableFuture<String> rescuer = CompletableFuture.supplyAsync(() -> {
gate.countDown(); // would release the parked workers...
return "released"; // ...but it never gets to run
}, ForkJoinPool.commonPool());
try {
rescuer.orTimeout(2, TimeUnit.SECONDS).join();
System.out.println("unexpected success");
} catch (CompletionException ce) {
System.out.println("rescuer failed: " + ce.getCause().getClass().getName()
+ " (no free worker to run the releasing stage)");
}
}
}
$ java -Djava.util.concurrent.ForkJoinPool.common.parallelism=4 P6_PoolHazard
common pool workers: 4
rescuer failed: java.util.concurrent.TimeoutException (no free worker to run the releasing stage)
Deterministic, no timing luck: every worker parked, the releasing stage queued behind them, the deadline the only thing that saved the join() from hanging forever. In the real Black Friday incident the "parked" stages were HTTP calls made on default-pool threads, and the victims included parallel streams elsewhere in the JVM that silently shared the same pool.
Two things this demo taught me that I didn't expect — both verified, not assumed. First, the -D…parallelism=4 flag isn't decoration: on this 2-CPU box the common pool's parallelism is 1, and I found that CompletableFuture won't even use a common pool that small — stages submitted with no executor and stages explicitly handed ForkJoinPool.commonPool() both ran on fresh Thread-N threads (verified by thread name), while a direct pool.submit() ran on a real pool worker. The JDK reroutes async stages away from a common pool it considers unusable. Second, orTimeout still fired while the pool was fully blocked, because its timer runs on a separate scheduler — deadlines are enforced outside the pool they protect.
Decision rule: the default pool is for short, non-blocking, CPU-bound stages only. Anything that waits — HTTP calls, database queries, join()/get() on another future, latches — gets a dedicated executor with a meaningful name and a bounded size, or virtual threads (a later post in this track). If you take one sentence from this section: never block a thread you don't own.
Cheat sheet: which method when
start work supplyAsync(() -> value) // returns a value
runAsync(() -> { ... }) // returns nothing
next step thenApply(v -> plainValue) // sync transformation
thenCompose(v -> future) // next step is itself async
combine thenCombine(a, b, (x, y) -> merged) // independent futures
allOf(f1, f2, ...) // wait for ALL, returns Void
anyOf(f1, f2, ...) // first finished wins
failures exceptionally(ex -> fallback) // failure -> value
handle((v, ex) -> result) // decide from value OR error
whenComplete((v, ex) -> { log(); }) // observe only, changes nothing
deadlines orTimeout(300, TimeUnit.MILLISECONDS) // slow = TimeoutException
completeOnTimeout(fallback, 300, TimeUnit.MILLISECONDS) // slow = fallback value
waiting join() // unchecked CompletionException on failure
get() // checked Execution/InterruptedException
threads supplyAsync(task) // shared common pool (or JDK fallback)
supplyAsync(task, myPool) // YOUR pool: blocking I/O goes here
And the three sentences that would have prevented the opening incident: overlap independent calls, don't serialize them. Every pipeline gets a deadline. Never block a thread you don't own.
What's next
You can now build async pipelines that overlap I/O, survive failures, respect deadlines, and don't wedge the JVM's shared thread pool. But there's a question this post hasn't answered: all those threads — your checkout-pool, the common pool's workers, the request threads — live somewhere. Each one holds a stack, each future holds its result, and every object your stages allocate lands on the heap. When a pipeline leaks futures or a pool's queue grows unbounded, the failure mode isn't a slow page — it's an OutOfMemoryError at 3 a.m. The next post in this track, JVM Memory & Garbage Collection: What Java Developers Need to Know, goes below the code: heap regions, what the garbage collector actually does on a pause, and how to read the memory footprint of the concurrent programs you're now writing.
Field check before you move on: (1) Take the P1_Pipeline program and put a deadline on the orders fetch — orTimeout(1500, MILLISECONDS) — then change the sleep to 3000 ms and confirm the TimeoutException path fires; then swap in completeOnTimeout with an empty order array and verify the summary still renders. (2) Chain three dependent service calls with thenCompose (user → orders → shipping estimate), add an exceptionally fallback at the end, and make the middle stage throw — check the fallback value is what join() returns. (3) Reproduce the pool-starvation demo with -Djava.util.concurrent.ForkJoinPool.common.parallelism=4, then replace the blocking await() stages with a dedicated Executors.newFixedThreadPool and confirm the rescuer completes instantly. (4) Find one CompletableFuture (or parallel stream) in your own codebase that does blocking I/O on the default pool, and move it to a named executor — then check a thread dump to confirm the new thread names appear. Bring the thread dump to the next post; we'll be reading memory, and threads are where the objects live.
Comments
Post a Comment