Java · Lesson 18 of 20
Threads, Executors and CompletableFuture
Learn Java concurrency: threads, race conditions, synchronized and atomics, ExecutorService, CompletableFuture pipelines and Java 21 virtual threads.
- Advanced
- 22 min read
- 4 objectives
Before this lessonLesson 17: Files and I/O
What you will learn
- Start threads and explain why shared state causes race conditions
- Protect shared data with synchronized, atomics and concurrent collections
- Run tasks with ExecutorService and Futures
- Compose async work with CompletableFuture and use virtual threads
Your Progress
0 of 20 lessons 0%
- Lessons0 / 20
- Completed0
- Est. time left~ 5 hours
Create a free account to keep your progress on every device.
Tip: pressing Next marks this lesson complete automatically.
A server handling thousands of requests, an app that downloads files while staying responsive, a batch job that uses every CPU core: all of these need concurrency, doing several things at once. Java has had threads since version 1.0, and it has grown a rich toolbox on top of them. Since Java 21, virtual threads have made the simple "one thread per task" style scale to millions of tasks.
Concurrency is powerful but easy to get subtly wrong, so this lesson starts with the core problem (shared state) before moving to the higher-level tools you should use day to day. Examples use short Thread.sleep calls to stand in for slow network or database calls.
Starting a thread
A Thread runs a Runnable (any lambda with no arguments and no result) alongside the main thread. start() begins execution; join() waits for it to finish. The order in which threads run is decided by the scheduler, so never rely on it; here we wait for each result before printing.
import java.util.*;
public class Main {
public static void main(String[] args) throws InterruptedException {
List<String> results = Collections.synchronizedList(new ArrayList<>());
Thread t1 = new Thread(() -> results.add("report generated"));
Thread t2 = Thread.ofPlatform().name("mailer").start(() -> results.add("emails sent"));
t1.start();
t1.join();
t2.join();
Collections.sort(results);
System.out.println(results);
System.out.println("main thread: " + Thread.currentThread().getName());
}
}[emails sent, report generated] main thread: main
Race conditions
Trouble starts when threads share mutable data. count++ looks like one step but is really three: read, add one, write back. If two threads read the same value at the same moment, one increment is lost. Run the unsafe counter and the total is usually less than expected (and different every run).
import java.util.concurrent.atomic.AtomicInteger;
public class Main {
static int unsafeCount = 0;
static int lockedCount = 0;
static final Object LOCK = new Object();
static final AtomicInteger atomicCount = new AtomicInteger();
public static void main(String[] args) throws InterruptedException {
Runnable work = () -> {
for (int i = 0; i < 100_000; i++) {
unsafeCount++; // race condition
synchronized (LOCK) { lockedCount++; } // one thread at a time
atomicCount.incrementAndGet(); // lock-free and safe
}
};
Thread a = new Thread(work), b = new Thread(work);
a.start(); b.start();
a.join(); b.join();
System.out.println("unsafe correct? " + (unsafeCount == 200_000));
System.out.println("synchronized: " + lockedCount);
System.out.println("atomic: " + atomicCount.get());
}
}unsafe correct? false synchronized: 200000 atomic: 200000
Two fixes are shown. A synchronized block lets only one thread hold the lock at a time. AtomicInteger uses special CPU instructions to update a single value safely without a lock. For maps shared between threads, use ConcurrentHashMap (its merge and computeIfAbsent are atomic).
ExecutorService: a pool of workers
Creating a thread by hand for every task is clumsy. An ExecutorService manages a pool of threads: you submit tasks and get back a Future, a handle to a result that will arrive later. future.get() waits for it. Since Java 19 executors are AutoCloseable, so try-with-resources waits for tasks and shuts the pool down for you.
import java.util.*;
import java.util.concurrent.*;
public class Main {
static int fetchPrice(String product) throws InterruptedException {
Thread.sleep(200); // simulate a slow API
return product.length() * 10;
}
public static void main(String[] args) throws Exception {
List<String> products = List.of("mouse", "keyboard", "monitor", "cable");
long start = System.nanoTime();
try (ExecutorService pool = Executors.newFixedThreadPool(4)) {
List<Future<Integer>> futures = new ArrayList<>();
for (String p : products) {
futures.add(pool.submit(() -> fetchPrice(p))); // Callable returns a value
}
int total = 0;
for (Future<Integer> f : futures) total += f.get();
System.out.println("total = " + total);
}
long ms = (System.nanoTime() - start) / 1_000_000;
System.out.println("parallel? " + (ms < 600));
}
}total = 250 parallel? true
Four 200 ms calls finished in roughly 200 ms, not 800, because they ran at the same time. If a task throws, get() rethrows it wrapped in an ExecutionException.
CompletableFuture: composing async steps
Future.get() blocks, and combining several futures by hand is awkward. CompletableFuture lets you describe a pipeline: start work with supplyAsync, transform results with thenApply, combine two independent results with thenCombine, chain another async call with thenCompose, and recover from errors with exceptionally. Nothing blocks until you finally call join().
import java.util.concurrent.*;
public class Main {
static void pause(int ms) {
try { Thread.sleep(ms); } catch (InterruptedException e) { throw new RuntimeException(e); }
}
static String loadUser(int id) { pause(100); return "Ada"; }
static int loadOrderCount(int id) { pause(150); return 3; }
static CompletableFuture<String> loadPlan(String name) {
return CompletableFuture.supplyAsync(() -> { pause(50); return "pro"; });
}
public static void main(String[] args) {
CompletableFuture<String> user = CompletableFuture.supplyAsync(() -> loadUser(1));
CompletableFuture<Integer> orders = CompletableFuture.supplyAsync(() -> loadOrderCount(1));
String dashboard = user
.thenApply(String::toUpperCase)
.thenCombine(orders, (name, n) -> name + " has " + n + " orders")
.join();
System.out.println(dashboard);
System.out.println(user.thenCompose(Main::loadPlan).join());
String safe = CompletableFuture.supplyAsync(() -> {
if (true) throw new IllegalStateException("payment service down");
return "ok";
})
.exceptionally(ex -> "fallback: " + ex.getCause().getMessage())
.join();
System.out.println(safe);
var all = CompletableFuture.allOf(user, orders);
all.join();
System.out.println("all done: " + all.isDone());
}
}ADA has 3 orders pro fallback: payment service down all done: true
By default supplyAsync runs on the shared ForkJoinPool.commonPool(). For blocking I/O, pass your own executor as a second argument so slow calls do not starve other work.
Virtual threads (Java 21+)
A normal (platform) thread maps to an operating system thread, which costs around a megabyte of memory, so you can only have a few thousand. Virtual threads are lightweight threads managed by the JVM: when one blocks on I/O or sleep, the JVM parks it and reuses the underlying OS thread for another. You can run a million of them, which means plain, readable blocking code now scales like complex async code. Spring Boot, Helidon and Quarkus all support them.
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.IntStream;
public class Main {
public static void main(String[] args) {
AtomicInteger done = new AtomicInteger();
long start = System.nanoTime();
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
IntStream.range(0, 10_000).forEach(i ->
executor.submit(() -> {
Thread.sleep(1000); // 10,000 tasks each "waiting" 1 second
done.incrementAndGet();
return i;
}));
} // waits for all tasks
long seconds = (System.nanoTime() - start) / 1_000_000_000;
System.out.println(done.get() + " tasks finished");
System.out.println("under 5 seconds? " + (seconds < 5));
System.out.println(Thread.ofVirtual().unstarted(() -> {}).isVirtual());
}
}10000 tasks finished under 5 seconds? true true
With a fixed pool of 100 platform threads the same work would take 100 seconds. Use virtual threads for I/O-bound tasks (HTTP calls, database queries). They do not make CPU-heavy work faster: for that, a fixed pool sized to Runtime.getRuntime().availableProcessors() or parallel streams are the right tools.
Recap
- Threads run code concurrently;
join()waits, and execution order is never guaranteed. - Shared mutable state causes race conditions; fix with
synchronized, atomics, concurrent collections or, best, no sharing. ExecutorServicemanages thread pools;submitreturns aFuture; close it with try-with-resources.CompletableFuturecomposes async steps withthenApply,thenCombine,thenComposeandexceptionally.- Virtual threads (
newVirtualThreadPerTaskExecutor) make blocking I/O code scale to huge numbers of tasks.
// Write your solution here
Finished reading? Mark this lesson complete to track your progress.
