NOTE
6.15 CompletableFuture
1. What it is. Used for asynchronous programming. In Java, so-called asynchronous programming means putting blocking code into a separate thread for execution and notifying the main thread when a result is available. 2. Future vs CompletableFuture. 3. Usage. 4. Source-code analysis.
This is a historical learning note and may contain outdated or incomplete understanding.
1. What It Is
Used for asynchronous programming. (More precisely, it is intended to describe non-blocking behavior.)
In Java, so-called asynchronous programming means putting blocking code into a separate thread for execution and notifying the main thread when a result is available.
2. Future VS CompletableFuture
| Future | CompletableFuture | |
|---|---|---|
| How results are obtained | Active polling. Use isDone to check whether the call has completed, and get to obtain the execution result |
Asynchronous callback. Use callback functions |
| Exception handling | Not supported | Supported |
| Chained calls | Not supported | Supported |
| Manually complete a task | Not supported | Supported |
3. Usage
3.1. Run a Task That Does Not Return a Result
CompletableFuture<Void> future = CompletableFuture.runAsync(() -> {
try
{
TimeUnit.SECONDS.sleep(5);
}
catch (InterruptedException e)
{
throw new IllegalStateException(e);
}
System.out.println("Background task completed");
});
future.get();
3.2. Run a Task That Returns a Result
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
try
{
TimeUnit.SECONDS.sleep(5);
}
catch (InterruptedException e)
{
throw new IllegalStateException(e);
}
return "Background task completed";
});
String s = future.get();
System.out.println(s);
3.3. Thread Pool
By default, tasks are executed using the thread pool in ForkJoin’s common pool, but an Executor can also be passed as the second parameter to specify the thread pool used for execution.
Executor executor = Executors.newFixedThreadPool(10);
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
try {
TimeUnit.SECONDS.sleep(1);
} catch (InterruptedException e) {
throw new IllegalStateException(e);
}
return "Result of the asynchronous computation";
}, executor);
3.4. Manually Complete a Task
CompletableFuture<String> stringCompletableFuture = new CompletableFuture<>();
new Thread(()->{
try
{
TimeUnit.SECONDS.sleep(5);
}
catch (InterruptedException e)
{
e.printStackTrace();
}
stringCompletableFuture.complete("Manually complete the task");
}).run();
String s = stringCompletableFuture.get();
System.out.println(s);
3.5. Callbacks
thenApply(): accepts the result as a parameter and returns a result.thenAccept(): accepts the result as a parameter and returns nothing.thenRun(): no parameters and no return value.
System.out.println("start");
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
try
{
TimeUnit.SECONDS.sleep(5);
}
catch (InterruptedException e)
{
throw new IllegalStateException(e);
}
return "Background task completed";
});
future.thenAccept(System.out::println);
System.out.println("The main thread continues executing and sleeps for 10s");
TimeUnit.SECONDS.sleep(10);
3.6. Chained Calls
System.out.println("start");
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
try
{
TimeUnit.SECONDS.sleep(5);
}
catch (InterruptedException e)
{
throw new IllegalStateException(e);
}
return "Background task completed";
});
future.thenApply(s->{
System.out.println(Thread.currentThread().getName() + "s");
return s;
}).thenApply(s->{
System.out.println(Thread.currentThread().getName() + "s");
return s;
});
System.out.println("The main thread continues executing and sleeps for 10s");
TimeUnit.SECONDS.sleep(10);
3.7. Combine Multiple CompletableFuture Instances
thenCompose(): two Futures with a dependency.thenCombine(): two Futures without a dependency.CompletableFuture.allOf: all Futures complete.CompletableFuture.anyOf: any one Future completes.
System.out.println("start runnning............");
long start = System.currentTimeMillis();
CompletableFuture<String> future1
= CompletableFuture.supplyAsync(() ->
{
try
{
TimeUnit.SECONDS.sleep(5);
}
catch (InterruptedException e)
{
e.printStackTrace();
}
System.out.println("Hello" + Thread.currentThread().getName());
return "Hello";
}
);
CompletableFuture<String> future2
= CompletableFuture.supplyAsync(() ->
{
try
{
TimeUnit.SECONDS.sleep(8);
}
catch (InterruptedException e)
{
e.printStackTrace();
}
System.out.println("Beautiful" + Thread.currentThread().getName());
return "Beautiful";
}
);
CompletableFuture<String> future3
= CompletableFuture.supplyAsync(() ->
{
try
{
TimeUnit.SECONDS.sleep(10);
}
catch (InterruptedException e)
{
e.printStackTrace();
}
System.out.println("World" + Thread.currentThread().getName());
return "World";
}
);
CompletableFuture<Void> combinedFuture
= CompletableFuture.allOf(future1, future2, future3);
combinedFuture.get();
long end = System.currentTimeMillis();
System.out.println("finish run...time is " + (end-start));
assertTrue(future1.isDone());
assertTrue(future2.isDone());
assertTrue(future3.isDone());
System.out.println(future1.get());
System.out.println(future2.get());
System.out.println(future3.get());
3.8. Exception Handling
exceptionallyis called when an exception occurs.handleis called whether or not an exception occurs.
CompletableFuture<Object> future = CompletableFuture.supplyAsync(() -> {
throw new IllegalArgumentException("Age can not be negative");
}).exceptionally(ex -> {
System.out.println("Oops! We have an exception - " + ex.getMessage());
return "Unknown!";
});
System.out.println(future.get());
4. Source-Code Analysis
4.1. Class Diagram

We can see that CompletableFuture implements the Future interface, so it is also something that can obtain an asynchronously executed result.
4.2. Fields
volatile Object result; // Either the result or boxed AltResult
volatile Completion stack; // Top of Treiber stack of dependent actions
The execution result is stored in Object result. If an exception occurs, it is wrapped in AltResult.
4.2.1. AltResult
static final class AltResult { // See above
final Throwable ex; // null only for NIL
AltResult(Throwable x) { this.ex = x; }
}
/** The encoding of the null value. */
static final AltResult NIL = new AltResult(null);
4.3. runAsync
public static CompletableFuture<Void> runAsync(Runnable runnable) {
return asyncRunStage(asyncPool, runnable);
}
It passes asyncPool and the runnable task to asyncRunStage.
First, let’s see how asyncPool is initialized.
4.3.1. Initialize the Default Thread Pool
// Returns true
private static final boolean useCommonPool =
(ForkJoinPool.getCommonPoolParallelism() > 1);
// ForkJoinPool.commonPool() is used here
private static final Executor asyncPool = useCommonPool ?
ForkJoinPool.commonPool() : new ThreadPerTaskExecutor();
So the default is ForkJoinPool.commonPool().
With the default thread pool available, the next call is asyncRunStage.
asyncRunStage
static CompletableFuture<Void> asyncRunStage(Executor e, Runnable f) {
if (f == null) throw new NullPointerException();
CompletableFuture<Void> d = new CompletableFuture<Void>();
e.execute(new AsyncRun(d, f));
return d;
}
- Line 2: if the task is null, throw an exception.
- Line 3: construct a
CompletableFutureto receive the result. - Line 4: first construct an
AsyncRunfrom theCompletableFutureandRunnable, and then call the thread-poolExecutor.executemethod to execute thisAsyncRun. - Line 5: return the
CompletableFuture.
4.3.2. Wrap the Task to Execute [Runnable] and the Result Receiver [CompletableFuture] into AsyncRun
First look at the AsyncRun class.
AsyncRun
static final class AsyncRun extends ForkJoinTask<Void>
implements Runnable, AsynchronousCompletionTask {
CompletableFuture<Void> dep; Runnable fn;
AsyncRun(CompletableFuture<Void> dep, Runnable fn) {
this.dep = dep; this.fn = fn;
}
public final Void getRawResult() { return null; }
public final void setRawResult(Void v) {}
public final boolean exec() { run(); return true; }
public void run() {
CompletableFuture<Void> d; Runnable f;
if ((d = dep) != null && (f = fn) != null) {
// Clear CompletableFuture and Runnable
dep = null; fn = null;
// If the CompletableFuture result is null
if (d.result == null) {
try {
// Execute the Runnable
f.run();
// CAS-set the CompletableFuture result to AltResult NIL -- see AltResult above
d.completeNull();
} catch (Throwable ex) {
// If an exception is thrown, CAS-set the CompletableFuture result to AltResult(exception) -- see AltResult above
d.completeThrowable(ex);
}
}
d.postComplete();
}
}
}
- Line 2: implements the
Runnableinterface. - Lines 4-6: the constructor simply stores the passed-in
RunnableandCompletableFuture. - Lines 12-26: the thread pool’s
executemethod eventually calls thisrunmethod. See the comments for details.
We can look at the methods used to set a null result and an exception result.
completeNull[null]
final boolean completeNull() {
// CAS-set RESULT to NIL
return UNSAFE.compareAndSwapObject(this, RESULT, null,
NIL);
}
completeThrowable[exception]
static AltResult encodeThrowable(Throwable x) {
return new AltResult((x instanceof CompletionException) ? x :
new CompletionException(x));
}
/** Completes with an exceptional result, unless already completed. */
final boolean completeThrowable(Throwable x) {
// CAS-set RESULT to AltResult(exception)
return UNSAFE.compareAndSwapObject(this, RESULT, null,
encodeThrowable(x));
}
4.3.3. Call the Thread Pool’s execute Method to Execute the AsyncRun Above
Executing AsyncRun eventually calls AsyncRun.run. The analysis is the same as in Wrap the Task to Execute [Runnable] and the Result Receiver [CompletableFuture] into AsyncRun.
4.4. supplyAsync
public static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier) {
return asyncSupplyStage(asyncPool, supplier);
}
It passes the default thread pool asyncPool and the task supplier. This supplier is a Supplier [functional interface], as shown below:

4.4.1. Initialize the Default Thread Pool
Initializing the default thread pool is the same as in runAsync above.
Next, continue tracing the asyncSupplyStage method.
asyncSupplyStage
static <U> CompletableFuture<U> asyncSupplyStage(Executor e,
Supplier<U> f) {
if (f == null) throw new NullPointerException();
CompletableFuture<U> d = new CompletableFuture<U>();
e.execute(new AsyncSupply<U>(d, f));
return d;
}
- Line 2: if the task is null, throw an exception.
- Line 3: construct a
CompletableFutureto receive the result. - Line 4: first construct an
AsyncSupplyfrom theCompletableFutureandSupplier, and then call the thread-poolExecutor.executemethod to execute thisAsyncSupply. - Line 5: return the
CompletableFuture.
4.4.2. Wrap the Task to Execute [Supplier] and the Result Receiver [CompletableFuture] into AsyncSupply
AsyncSupply
static final class AsyncSupply<T> extends ForkJoinTask<Void>
implements Runnable, AsynchronousCompletionTask {
CompletableFuture<T> dep; Supplier<T> fn;
AsyncSupply(CompletableFuture<T> dep, Supplier<T> fn) {
this.dep = dep; this.fn = fn;
}
public final Void getRawResult() { return null; }
public final void setRawResult(Void v) {}
public final boolean exec() { run(); return true; }
public void run() {
CompletableFuture<T> d; Supplier<T> f;
if ((d = dep) != null && (f = fn) != null) {
// Clear CompletableFuture and Runnable
dep = null; fn = null;
// If the CompletableFuture result is null
if (d.result == null) {
try {
// Call Supplier.get to obtain the result,
// then call CompletableFuture.completeValue to store it
d.completeValue(f.get());
} catch (Throwable ex) {
// If an exception is thrown, CAS-set the CompletableFuture result to AltResult(exception) -- see AltResult above
d.completeThrowable(ex);
}
}
d.postComplete();
}
}
}
- Line 2: implements the
Runnableinterface. - Lines 4-6: the constructor simply stores the passed-in
RunnableandCompletableFuture. - Lines 12-26: the thread pool’s
executemethod eventually calls thisrunmethod. See the comments for details.
We can look at the completeValue method used to set the result.
completeValue
final boolean completeValue(T t) {
return UNSAFE.compareAndSwapObject(this, RESULT, null,
(t == null) ? NIL : t);
}
4.4.3. Call the Thread Pool’s execute Method to Execute the AsyncRun Above
When AsyncRun is executed, it eventually calls AsyncRun.run. The analysis is the same as in Wrap the Task to Execute [Supplier] and the Result Receiver [CompletableFuture] into AsyncSupply above.
4.5. complete
public boolean complete(T value) {
boolean triggered = completeValue(value);
postComplete();
return triggered;
}
- Line 2: manually set the result.
- Line 3: execute the hook method.
4.5.1. Manually Set the Result
final boolean completeValue(T t) {
return UNSAFE.compareAndSwapObject(this, RESULT, null,
(t == null) ? NIL : t);
}
4.5.2. Execute the Hook Method
I definitely did not understand what this code is supposed to do.
final void postComplete() {
/*
* On each step, variable f holds current dependents to pop
* and run. It is extended along only one path at a time,
* pushing others to avoid unbounded recursion.
*/
CompletableFuture<?> f = this; Completion h;
while ((h = f.stack) != null ||
(f != this && (h = (f = this).stack) != null)) {
CompletableFuture<?> d; Completion t;
if (f.casStack(h, t = h.next)) {
if (t != null) {
if (f != this) {
pushStack(h);
continue;
}
h.next = null; // detach
}
f = (d = h.tryFire(NESTED)) == null ? this : d;
}
}
}
Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub