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.

JavaCreated Updated 3 min readhistorical

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

  • exceptionally is called when an exception occurs.
  • handle is 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 CompletableFuture to receive the result.
  • Line 4: first construct an AsyncRun from the CompletableFuture and Runnable, and then call the thread-pool Executor.execute method to execute this AsyncRun.
  • 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 Runnable interface.
  • Lines 4-6: the constructor simply stores the passed-in Runnable and CompletableFuture.
  • Lines 12-26: the thread pool’s execute method eventually calls this run method. 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 CompletableFuture to receive the result.
  • Line 4: first construct an AsyncSupply from the CompletableFuture and Supplier, and then call the thread-pool Executor.execute method to execute this AsyncSupply.
  • 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 Runnable interface.
  • Lines 4-6: the constructor simply stores the passed-in Runnable and CompletableFuture.
  • Lines 12-26: the thread pool’s execute method eventually calls this run method. 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;
        }
    }
}

5. References

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub