Creating threads by hand doesn't scale. Real Java code uses the java.util.concurrent toolkit: thread pools via ExecutorService, async pipelines with CompletableFuture, divide-and-conquer with Fork/Join, coordination with latches and semaphores, smarter locks, producer-consumer queues, and, since Java 21, virtual threads.
Executor Framework & Thread Pools
Why Thread Pools?
// Problem: Creating new Thread for each task is EXPENSIVE
// Each thread: ~1MB stack, OS context switching overhead
// Solution: Thread Pool — reuse pre-created threads
// ── EXECUTOR SERVICE ──
// Fixed pool: exactly N threads
ExecutorService fixed = Executors.newFixedThreadPool(4);
// Cached pool: grow/shrink as needed, idle threads die after 60s
ExecutorService cached = Executors.newCachedThreadPool();
// Single thread: guarantees sequential execution
ExecutorService single = Executors.newSingleThreadExecutor();
// Scheduled: run after delay or periodically
ScheduledExecutorService scheduled = Executors.newScheduledThreadPool(2);
scheduled.schedule(() -> System.out.println("After 2s"), 2, TimeUnit.SECONDS);
scheduled.scheduleAtFixedRate(() -> System.out.println("Every 1s"),
0, 1, TimeUnit.SECONDS);
// Submitting tasks
fixed.submit(() -> System.out.println("Task 1")); // Runnable
Future<String> f = fixed.submit(() -> "Result"); // Callable
String result = f.get(5, TimeUnit.SECONDS); // timeout
// Shutdown (important!)
fixed.shutdown(); // stop accepting, finish existing
fixed.shutdownNow(); // interrupt all threads immediately
fixed.awaitTermination(10, TimeUnit.SECONDS); // wait for finish
ThreadPoolExecutor — Full Control
// All Executors.newXxx() use ThreadPoolExecutor internally
ThreadPoolExecutor executor = new ThreadPoolExecutor(
2, // corePoolSize: min threads always alive
10, // maximumPoolSize: max threads ever
60L, // keepAliveTime: idle extra threads die after 60
TimeUnit.SECONDS,
new LinkedBlockingQueue<>(100), // task queue (capacity 100)
Executors.defaultThreadFactory(),
new ThreadPoolExecutor.CallerRunsPolicy() // rejection policy
);
// Rejection Policies (when queue full + max threads reached):
// AbortPolicy — throws RejectedExecutionException (default)
// CallerRunsPolicy — caller thread runs the task
// DiscardPolicy — silently drops task
// DiscardOldestPolicy — drops oldest queued task, retries
// Flow: task submitted →
// if corePool has idle thread: use it
// else if queue not full: queue it
// else if maxPool not reached: create new thread
// else: rejection policy
Interview Questions
What is the difference between submit() and execute()?
execute() accepts Runnable, returns nothing. submit() accepts Runnable or Callable, returns Future. Future lets you check if done, get result (blocking), cancel, or handle exceptions. Prefer submit() for better error handling.
What is a Future in Java?
Future represents the result of an async computation. Methods: get() (blocking), get(timeout) (with timeout), isDone(), cancel(), isCancelled(). CompletableFuture (Java 8) is an enhanced Future with non-blocking chaining via thenApply(), thenAccept(), etc.
CompletableFuture — Async Programming
import java.util.concurrent.CompletableFuture;
// ── BASIC ASYNC ──
// supplyAsync: returns a value (Supplier)
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
Thread.sleep(1000); // simulate async work
return "Hello from async!";
});
// runAsync: no return value (Runnable)
CompletableFuture<Void> fire = CompletableFuture.runAsync(() -> {
System.out.println("Fire and forget");
});
// ── CHAINING ──
CompletableFuture<Integer> chain = CompletableFuture
.supplyAsync(() -> "Hello World") // async: String
.thenApply(s -> s.length()) // transform: String→Integer
.thenApply(len -> len * 2); // transform: Integer→Integer
System.out.println(chain.get()); // blocks: 22
// thenAccept: consume result (void)
future.thenAccept(result -> System.out.println("Got: " + result));
// thenRun: run after completion (no result)
future.thenRun(() -> System.out.println("Done!"));
// thenCompose: flatten nested CompletableFuture
CompletableFuture<String> composed = future
.thenCompose(result -> CompletableFuture.supplyAsync(
() -> result + " Processed"));
// ── COMBINING ──
CompletableFuture<String> userFuture = CompletableFuture.supplyAsync(() -> fetchUser(1));
CompletableFuture<String> orderFuture = CompletableFuture.supplyAsync(() -> fetchOrders(1));
// thenCombine: combine TWO futures when both complete
CompletableFuture<String> combined = userFuture.thenCombine(
orderFuture, (user, orders) -> user + " has orders: " + orders);
// allOf: wait for ALL futures
CompletableFuture<Void> all = CompletableFuture.allOf(
userFuture, orderFuture);
all.join(); // wait without checked exception
// anyOf: complete when FIRST finishes
CompletableFuture<Object> any = CompletableFuture.anyOf(
userFuture, orderFuture);
// ── ERROR HANDLING ──
future
.exceptionally(ex -> {
System.out.println("Error: " + ex.getMessage());
return "default value"; // fallback
})
.thenAccept(System.out::println);
// handle: always called (success or failure)
future.handle((result, ex) -> {
if (ex != null) return "Error: " + ex.getMessage();
return result.toUpperCase();
});
// whenComplete: callback (cannot transform result)
future.whenComplete((result, ex) -> {
if (ex != null) logError(ex);
else logSuccess(result);
});
// ── CUSTOM EXECUTOR ──
ExecutorService pool = Executors.newFixedThreadPool(4);
CompletableFuture.supplyAsync(() -> heavyWork(), pool);
pool.shutdown();
Fork/Join Framework
import java.util.concurrent.*;
// ForkJoinPool: divide-and-conquer parallel computation
// RecursiveTask<V>: returns a value
// RecursiveAction: no return value
public class ParallelSum extends RecursiveTask<Long> {
private static final int THRESHOLD = 10_000;
private final int[] array;
private final int start, end;
public ParallelSum(int[] array, int start, int end) {
this.array = array; this.start = start; this.end = end;
}
@Override
protected Long compute() {
int length = end - start;
if (length <= THRESHOLD) {
// Base case: compute directly
long sum = 0;
for (int i = start; i < end; i++) sum += array[i];
return sum;
}
// Fork: split into two halves
int mid = start + length / 2;
ParallelSum leftTask = new ParallelSum(array, start, mid);
ParallelSum rightTask = new ParallelSum(array, mid, end);
leftTask.fork(); // execute left asynchronously
long rightResult = rightTask.compute(); // execute right in current thread
long leftResult = leftTask.join(); // wait for left result
return leftResult + rightResult;
}
}
// Usage
int[] array = new int[1_000_000];
Arrays.fill(array, 1); // all 1s
ForkJoinPool pool = ForkJoinPool.commonPool();
Long sum = pool.invoke(new ParallelSum(array, 0, array.length));
System.out.println(sum); // 1000000
Semaphore, CountDownLatch, CyclicBarrier, Phaser
// SEMAPHORE: limit concurrent access (e.g., max 3 DB connections)
Semaphore semaphore = new Semaphore(3); // 3 permits
ExecutorService pool = Executors.newFixedThreadPool(10);
for (int i = 0; i < 10; i++) {
pool.submit(() -> {
try {
semaphore.acquire(); // wait for permit (blocks if 0 permits)
System.out.println("Working: " + Thread.currentThread().getName());
Thread.sleep(1000); // simulate work
} finally {
semaphore.release(); // always release!
}
});
}
// At most 3 threads work simultaneously
// COUNTDOWNLATCH: wait for N events
CountDownLatch latch = new CountDownLatch(3);
for (int i = 0; i < 3; i++) {
new Thread(() -> {
System.out.println("Service ready: " + Thread.currentThread().getName());
latch.countDown(); // decrement count
}).start();
}
latch.await(); // main thread waits until count = 0
System.out.println("All services ready! Starting app...");
// CYCLICBARRIER: all threads wait at barrier, then continue together
CyclicBarrier barrier = new CyclicBarrier(3, () ->
System.out.println("All threads reached barrier — continuing!"));
for (int i = 0; i < 3; i++) {
new Thread(() -> {
System.out.println(Thread.currentThread().getName() + " at barrier");
try { barrier.await(); } catch (Exception e) {}
System.out.println(Thread.currentThread().getName() + " continuing");
}).start();
}
// CyclicBarrier can be RESET and reused (unlike CountDownLatch)
// BLOCKING QUEUE: thread-safe producer-consumer
BlockingQueue<String> queue = new ArrayBlockingQueue<>(100);
// Producer thread
new Thread(() -> {
try {
queue.put("item1"); // blocks if queue full
queue.put("item2");
} catch (InterruptedException e) { Thread.currentThread().interrupt(); }
}).start();
// Consumer thread
new Thread(() -> {
try {
String item = queue.take(); // blocks if queue empty
System.out.println("Consumed: " + item);
} catch (InterruptedException e) { Thread.currentThread().interrupt(); }
}).start();
ReentrantReadWriteLock & StampedLock
import java.util.concurrent.locks.*;
// ReentrantReadWriteLock: multiple readers OR one writer
public class SafeCache<K,V> {
private final Map<K,V> cache = new HashMap<>();
private final ReentrantReadWriteLock lock = new ReentrantReadWriteLock();
private final Lock readLock = lock.readLock();
private final Lock writeLock = lock.writeLock();
public V get(K key) {
readLock.lock(); // multiple threads can read simultaneously
try { return cache.get(key); }
finally { readLock.unlock(); }
}
public void put(K key, V value) {
writeLock.lock(); // exclusive — no readers or writers
try { cache.put(key, value); }
finally { writeLock.unlock(); }
}
}
// StampedLock (Java 8): optimistic reads — fastest!
public class Point {
private double x, y;
private final StampedLock sl = new StampedLock();
public void move(double deltaX, double deltaY) {
long stamp = sl.writeLock();
try { x += deltaX; y += deltaY; }
finally { sl.unlockWrite(stamp); }
}
public double distanceFromOrigin() {
// Optimistic read — no lock acquired!
long stamp = sl.tryOptimisticRead();
double currentX = x, currentY = y;
if (!sl.validate(stamp)) { // was there a write while reading?
stamp = sl.readLock(); // fall back to real read lock
try { currentX = x; currentY = y; }
finally { sl.unlockRead(stamp); }
}
return Math.sqrt(currentX*currentX + currentY*currentY);
}
}
// StampedLock optimistic read: ~3x faster than ReadWriteLock in low-contention
BlockingQueue Hierarchy
// BlockingQueue: thread-safe queue with blocking operations
// Perfect for producer-consumer pattern
// ── ArrayBlockingQueue: bounded, backed by array ──
BlockingQueue<String> bounded = new ArrayBlockingQueue<>(10);
bounded.put("item"); // blocks if full
bounded.take(); // blocks if empty
bounded.offer("item", 5, TimeUnit.SECONDS); // timeout
bounded.poll(5, TimeUnit.SECONDS); // timeout
bounded.peek(); // look without removing
bounded.remainingCapacity(); // space left
// ── LinkedBlockingQueue: optionally bounded ──
BlockingQueue<Integer> lbq = new LinkedBlockingQueue<>(); // unbounded
BlockingQueue<Integer> lbq2 = new LinkedBlockingQueue<>(100); // bounded
// Separate head/tail locks → better concurrent throughput than Array
// ── PriorityBlockingQueue: unbounded priority queue ──
BlockingQueue<Integer> pbq = new PriorityBlockingQueue<>();
pbq.add(5); pbq.add(1); pbq.add(3);
pbq.take(); // 1 (always min first)
// put() never blocks (unbounded), take() blocks if empty
// ── SynchronousQueue: no storage! ──
BlockingQueue<String> sync = new SynchronousQueue<>();
// put() blocks until another thread calls take()
// take() blocks until another thread calls put()
// Used in Executors.newCachedThreadPool()
new Thread(() -> {
try { sync.put("handoff"); } catch(InterruptedException e){}
}).start();
String received = sync.take(); // blocks until put() called
// ── DelayQueue: elements available after delay ──
class DelayedTask implements Delayed {
private final long executeAt;
private final String name;
DelayedTask(String name, long delayMs) {
this.name = name;
this.executeAt = System.currentTimeMillis() + delayMs;
}
public long getDelay(TimeUnit unit) {
return unit.convert(executeAt - System.currentTimeMillis(), TimeUnit.MILLISECONDS);
}
public int compareTo(Delayed other) {
return Long.compare(executeAt, ((DelayedTask)other).executeAt);
}
}
DelayQueue<DelayedTask> dq = new DelayQueue<>();
dq.put(new DelayedTask("task1", 1000)); // available after 1 second
DelayedTask task = dq.take(); // blocks until task is ready
System.out.println(task.name); // task1 (after 1s)
Java 21 — Virtual Threads (Project Loom)
What are Virtual Threads?
Traditional Java threads are OS threads (expensive — ~1MB each). Virtual threads are JVM-managed lightweight threads (few KB each). You can create MILLIONS without exhausting system resources. Perfect for I/O-bound tasks like web servers.
// ── VIRTUAL THREADS (Java 21) ──
// Create virtual thread
Thread vt = Thread.ofVirtual().start(() -> {
System.out.println("Virtual thread: " + Thread.currentThread().isVirtual());
});
vt.join();
// Virtual thread executor
try (ExecutorService exec = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < 1_000_000; i++) {
int taskId = i;
exec.submit(() -> {
// Each task gets its own virtual thread!
Thread.sleep(Duration.ofMillis(100)); // I/O wait
return taskId;
});
}
} // waits for all tasks to complete
// Before: 1M requests → need 1M OS threads → crash
// After: 1M requests → 1M virtual threads → works fine!
// ── STRUCTURED CONCURRENCY (Java 21) ──
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
Future<String> user = scope.fork(() -> fetchUser(userId));
Future<String> orders = scope.fork(() -> fetchOrders(userId));
scope.join(); // wait for both
scope.throwIfFailed(); // propagate any failure
// Both succeeded
return user.resultNow() + orders.resultNow();
}
// WHY VIRTUAL THREADS MATTER:
// Traditional: 1 request → 1 OS thread (while waiting for DB: thread blocked!)
// Virtual: 1 request → 1 virtual thread (while waiting for DB:
// virtual thread unmounts from OS thread,
// OS thread serves OTHER requests!)