executors: ExecutorService lifecycle companion code
Adds the executors module: shutdown() vs shutdownNow(), the graceful two-phase shutdown pattern, AutoCloseable/try-with-resources (JDK 19+), newVirtualThreadPerTaskExecutor() timing vs a bounded pool, and the submit()-without-get() lost-exception trap. 5 runnable demos, 1 JUnit correctness test class (23 tests), 6 captured output transcripts. Co-Authored-By: Claude Sonnet 5 <[email protected]> Claude-Session: https://claude.ai/code/session_01FhzLY5p6okFva3qsnsRyvM
This commit is contained in:
@@ -0,0 +1,36 @@
|
||||
package com.ankurm.executors;
|
||||
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
|
||||
/**
|
||||
* Since JDK 19, ExecutorService extends AutoCloseable. Its default close() does exactly the
|
||||
* two-phase pattern from GracefulShutdownPatternDemo for you: shutdown(), then block awaiting
|
||||
* termination, escalating to shutdownNow() if interrupted while waiting - which means
|
||||
* try-with-resources is now a correct, one-line replacement for that hand-rolled block, not a
|
||||
* shortcut that skips the graceful phase.
|
||||
*/
|
||||
public class AutoCloseableExecutorDemo {
|
||||
|
||||
public static void main(String[] args) {
|
||||
System.out.println("=== try-with-resources: close() runs the graceful pattern automatically ===");
|
||||
try (ExecutorService pool = Executors.newFixedThreadPool(3)) {
|
||||
for (int i = 0; i < 3; i++) {
|
||||
int id = i;
|
||||
pool.submit(() -> {
|
||||
try {
|
||||
Thread.sleep(300);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
System.out.println("task-" + id + ": finished");
|
||||
});
|
||||
}
|
||||
System.out.println("main: leaving the try-with-resources block now");
|
||||
}
|
||||
// pool.close() already ran here, and it already blocked until every task above printed.
|
||||
System.out.println("main: back from the try-with-resources block - the pool is terminated,"
|
||||
+ " every task already finished, no separate awaitTermination() call was written");
|
||||
System.out.println("done");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,70 @@
|
||||
package com.ankurm.executors;
|
||||
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
/**
|
||||
* The canonical shutdown pattern from the ExecutorService javadoc: shutdown() first, give it a
|
||||
* bounded grace period via awaitTermination(), and only escalate to shutdownNow() if that grace
|
||||
* period runs out. Run twice - once against tasks that finish inside the grace period, once
|
||||
* against a task that deliberately outlives it - so both branches of the pattern actually execute
|
||||
* rather than just one being described.
|
||||
*/
|
||||
public class GracefulShutdownPatternDemo {
|
||||
|
||||
private static void shutdownGracefully(ExecutorService pool, String label) {
|
||||
pool.shutdown();
|
||||
try {
|
||||
System.out.println(label + ": shutdown() called, waiting up to 1s for a clean finish");
|
||||
if (!pool.awaitTermination(1, TimeUnit.SECONDS)) {
|
||||
System.out.println(label + ": grace period expired, escalating to shutdownNow()");
|
||||
pool.shutdownNow();
|
||||
if (!pool.awaitTermination(1, TimeUnit.SECONDS)) {
|
||||
System.out.println(label + ": pool still didn't terminate after shutdownNow()");
|
||||
} else {
|
||||
System.out.println(label + ": terminated after shutdownNow()");
|
||||
}
|
||||
} else {
|
||||
System.out.println(label + ": terminated cleanly within the grace period");
|
||||
}
|
||||
} catch (InterruptedException e) {
|
||||
System.out.println(label + ": interrupted while awaiting termination - shutting down now and re-interrupting");
|
||||
pool.shutdownNow();
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
}
|
||||
|
||||
public static void main(String[] args) {
|
||||
System.out.println("=== Case 1: tasks finish comfortably inside the grace period ===");
|
||||
ExecutorService quickPool = Executors.newFixedThreadPool(2);
|
||||
for (int i = 0; i < 2; i++) {
|
||||
int id = i;
|
||||
quickPool.submit(() -> {
|
||||
try {
|
||||
Thread.sleep(200);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
System.out.println("quick-task-" + id + ": done (200ms)");
|
||||
});
|
||||
}
|
||||
shutdownGracefully(quickPool, "quickPool");
|
||||
|
||||
System.out.println();
|
||||
System.out.println("=== Case 2: a task outlives the grace period, forcing shutdownNow() ===");
|
||||
ExecutorService slowPool = Executors.newFixedThreadPool(1);
|
||||
slowPool.submit(() -> {
|
||||
try {
|
||||
System.out.println("slow-task: started, will try to sleep 3s");
|
||||
Thread.sleep(3000);
|
||||
System.out.println("slow-task: this line should never print - shutdownNow() should interrupt it first");
|
||||
} catch (InterruptedException e) {
|
||||
System.out.println("slow-task: interrupted during shutdownNow(), stopping early");
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
});
|
||||
shutdownGracefully(slowPool, "slowPool");
|
||||
System.out.println("done");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,48 @@
|
||||
package com.ankurm.executors;
|
||||
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.Future;
|
||||
|
||||
/**
|
||||
* execute(Runnable) lets an uncaught exception reach the worker thread's UncaughtExceptionHandler,
|
||||
* which by default prints a stack trace - the same as any other unhandled exception. submit(...)
|
||||
* instead captures the exception INSIDE the returned Future, and nothing rethrows it unless
|
||||
* something actually calls Future.get() (or isCancelled()/isDone() plus get()). Skip that call and
|
||||
* the exception is silently discarded - the task appears to have failed with no output at all.
|
||||
*/
|
||||
public class LostExceptionTrapDemo {
|
||||
|
||||
public static void main(String[] args) throws InterruptedException {
|
||||
try (ExecutorService pool = Executors.newFixedThreadPool(2)) {
|
||||
|
||||
System.out.println("=== execute(): uncaught exception reaches the default handler, gets printed ===");
|
||||
pool.execute(() -> {
|
||||
throw new RuntimeException("boom from execute()");
|
||||
});
|
||||
Thread.sleep(300); // let the stack trace print before the next section starts
|
||||
|
||||
System.out.println();
|
||||
System.out.println("=== submit(): same exception, but nobody calls get() - dead silence ===");
|
||||
Future<?> lostFuture = pool.submit(() -> {
|
||||
throw new RuntimeException("boom from submit(), never retrieved");
|
||||
});
|
||||
Thread.sleep(300);
|
||||
System.out.println("main: task's Future exists (isDone=" + lostFuture.isDone()
|
||||
+ "), but its exception was never printed anywhere - it's just sitting inside the Future");
|
||||
|
||||
System.out.println();
|
||||
System.out.println("=== submit() + get(): the SAME exception, now surfaced by calling get() ===");
|
||||
Future<?> checkedFuture = pool.submit(() -> {
|
||||
throw new RuntimeException("boom from submit(), retrieved via get()");
|
||||
});
|
||||
try {
|
||||
checkedFuture.get();
|
||||
} catch (ExecutionException e) {
|
||||
System.out.println("main: caught ExecutionException, cause = " + e.getCause());
|
||||
}
|
||||
}
|
||||
System.out.println("done");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,59 @@
|
||||
package com.ankurm.executors;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
/**
|
||||
* shutdown() lets everything already queued or running finish, and just stops accepting new
|
||||
* work. shutdownNow() interrupts every actively-running task and hands back the tasks that were
|
||||
* still sitting in the queue, un-started. Both demos submit 6 tasks to a 2-thread pool, so at the
|
||||
* moment shutdown is called, 2 are running and 4 are still queued.
|
||||
*/
|
||||
public class ShutdownVsShutdownNowDemo {
|
||||
|
||||
private static void runTask(int id) {
|
||||
try {
|
||||
System.out.println("task-" + id + ": started");
|
||||
Thread.sleep(500);
|
||||
System.out.println("task-" + id + ": finished normally");
|
||||
} catch (InterruptedException e) {
|
||||
System.out.println("task-" + id + ": interrupted, exiting early");
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
}
|
||||
|
||||
private static ExecutorService submitSix() {
|
||||
ExecutorService pool = Executors.newFixedThreadPool(2);
|
||||
for (int i = 0; i < 6; i++) {
|
||||
int id = i;
|
||||
pool.submit(() -> runTask(id));
|
||||
}
|
||||
return pool;
|
||||
}
|
||||
|
||||
public static void main(String[] args) throws InterruptedException {
|
||||
System.out.println("=== shutdown(): lets running AND already-queued tasks finish ===");
|
||||
ExecutorService gracefulPool = submitSix();
|
||||
Thread.sleep(100); // let the first 2 tasks actually start
|
||||
gracefulPool.shutdown();
|
||||
System.out.println("main: called shutdown() - isShutdown=" + gracefulPool.isShutdown()
|
||||
+ ", isTerminated=" + gracefulPool.isTerminated());
|
||||
boolean finished = gracefulPool.awaitTermination(5, TimeUnit.SECONDS);
|
||||
System.out.println("main: awaitTermination returned " + finished
|
||||
+ " - all 6 tasks ran to completion, none were skipped");
|
||||
|
||||
System.out.println();
|
||||
System.out.println("=== shutdownNow(): interrupts running tasks, abandons queued ones ===");
|
||||
ExecutorService abruptPool = submitSix();
|
||||
Thread.sleep(100); // let the first 2 tasks actually start
|
||||
List<Runnable> abandoned = abruptPool.shutdownNow();
|
||||
System.out.println("main: called shutdownNow() - " + abandoned.size()
|
||||
+ " queued tasks abandoned, never started");
|
||||
boolean finished2 = abruptPool.awaitTermination(5, TimeUnit.SECONDS);
|
||||
System.out.println("main: awaitTermination returned " + finished2
|
||||
+ " - the 2 running tasks were interrupted mid-sleep, not allowed to finish");
|
||||
System.out.println("done");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,65 @@
|
||||
package com.ankurm.executors;
|
||||
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
/**
|
||||
* newVirtualThreadPerTaskExecutor() isn't a pool at all - there's no fixed thread count to
|
||||
* exhaust, so it never queues. Every submitted task gets its own virtual thread immediately.
|
||||
* Run 5,000 blocking (sleeping) tasks against it and against a 20-thread fixed pool doing
|
||||
* identical work, on the same 2-vCPU box, to see what that actually costs in wall-clock time.
|
||||
*/
|
||||
public class VirtualThreadPerTaskExecutorDemo {
|
||||
|
||||
private static final int TASK_COUNT = 5_000;
|
||||
private static final int SLEEP_MILLIS = 50;
|
||||
|
||||
private static long runAndTime(ExecutorService pool, String label) throws InterruptedException {
|
||||
AtomicInteger completed = new AtomicInteger(0);
|
||||
long start = System.nanoTime();
|
||||
for (int i = 0; i < TASK_COUNT; i++) {
|
||||
pool.submit(() -> {
|
||||
try {
|
||||
Thread.sleep(SLEEP_MILLIS);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
completed.incrementAndGet();
|
||||
});
|
||||
}
|
||||
pool.shutdown();
|
||||
pool.awaitTermination(2, TimeUnit.MINUTES);
|
||||
long elapsedMs = (System.nanoTime() - start) / 1_000_000;
|
||||
System.out.println(label + ": " + completed.get() + "/" + TASK_COUNT
|
||||
+ " tasks completed in " + elapsedMs + "ms");
|
||||
return elapsedMs;
|
||||
}
|
||||
|
||||
public static void main(String[] args) throws InterruptedException {
|
||||
System.out.println("Running " + TASK_COUNT + " tasks that each sleep " + SLEEP_MILLIS + "ms.");
|
||||
System.out.println();
|
||||
|
||||
try (ExecutorService vtPool = Executors.newVirtualThreadPerTaskExecutor()) {
|
||||
// Peek at one task's thread to show it really is a virtual thread, not a pooled platform one.
|
||||
vtPool.submit(() -> System.out.println("sample task thread: " + Thread.currentThread()));
|
||||
Thread.sleep(50);
|
||||
long vtElapsed = runAndTime(vtPool, "newVirtualThreadPerTaskExecutor()");
|
||||
System.out.println("(" + vtElapsed + "ms is close to the " + SLEEP_MILLIS
|
||||
+ "ms a single task takes - every task ran concurrently, nothing waited for a free worker)");
|
||||
}
|
||||
|
||||
System.out.println();
|
||||
|
||||
try (ExecutorService fixedPool = Executors.newFixedThreadPool(20)) {
|
||||
fixedPool.submit(() -> System.out.println("sample task thread: " + Thread.currentThread()));
|
||||
Thread.sleep(50);
|
||||
long fixedElapsed = runAndTime(fixedPool, "newFixedThreadPool(20)");
|
||||
long expectedRounds = (TASK_COUNT + 19) / 20;
|
||||
System.out.println("(" + fixedElapsed + "ms is roughly " + expectedRounds
|
||||
+ " rounds of " + SLEEP_MILLIS + "ms each - only 20 tasks are ever running at once)");
|
||||
}
|
||||
System.out.println("done");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,108 @@
|
||||
package com.ankurm.executors;
|
||||
|
||||
import org.junit.jupiter.api.RepeatedTest;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.Future;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
class ExecutorShutdownBehaviorTest {
|
||||
|
||||
@RepeatedTest(10)
|
||||
void shutdownLetsEveryQueuedTaskFinish() throws InterruptedException {
|
||||
AtomicInteger completed = new AtomicInteger(0);
|
||||
ExecutorService pool = Executors.newFixedThreadPool(2);
|
||||
int taskCount = 12;
|
||||
for (int i = 0; i < taskCount; i++) {
|
||||
pool.submit(completed::incrementAndGet);
|
||||
}
|
||||
pool.shutdown();
|
||||
boolean terminated = pool.awaitTermination(5, TimeUnit.SECONDS);
|
||||
assertTrue(terminated, "pool should terminate within the timeout");
|
||||
assertEquals(taskCount, completed.get(), "shutdown() must not drop any queued task");
|
||||
}
|
||||
|
||||
@RepeatedTest(10)
|
||||
void shutdownNowAbandonsOnlyTheStillQueuedTasks() throws InterruptedException {
|
||||
ExecutorService pool = Executors.newFixedThreadPool(1);
|
||||
CountDownLatch firstTaskStarted = new CountDownLatch(1);
|
||||
CountDownLatch releaseFirstTask = new CountDownLatch(1);
|
||||
|
||||
pool.submit(() -> {
|
||||
firstTaskStarted.countDown();
|
||||
try {
|
||||
releaseFirstTask.await();
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
});
|
||||
int queuedBehindIt = 3;
|
||||
for (int i = 0; i < queuedBehindIt; i++) {
|
||||
pool.submit(() -> {
|
||||
});
|
||||
}
|
||||
|
||||
assertTrue(firstTaskStarted.await(5, TimeUnit.SECONDS), "first task should be running before shutdownNow()");
|
||||
List<Runnable> abandoned = pool.shutdownNow();
|
||||
assertEquals(queuedBehindIt, abandoned.size(),
|
||||
"shutdownNow() must return exactly the tasks that never started");
|
||||
|
||||
releaseFirstTask.countDown();
|
||||
assertTrue(pool.awaitTermination(5, TimeUnit.SECONDS));
|
||||
}
|
||||
|
||||
@Test
|
||||
void tryWithResourcesBlocksUntilEveryTaskIsDone() {
|
||||
AtomicInteger completed = new AtomicInteger(0);
|
||||
int taskCount = 20;
|
||||
try (ExecutorService pool = Executors.newFixedThreadPool(4)) {
|
||||
for (int i = 0; i < taskCount; i++) {
|
||||
pool.submit(() -> {
|
||||
try {
|
||||
Thread.sleep(10);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
completed.incrementAndGet();
|
||||
});
|
||||
}
|
||||
}
|
||||
// close() already ran (and blocked) by the time we get here.
|
||||
assertEquals(taskCount, completed.get(), "AutoCloseable.close() must wait for every task");
|
||||
}
|
||||
|
||||
@Test
|
||||
void virtualThreadPerTaskExecutorReallyUsesVirtualThreads() throws Exception {
|
||||
boolean[] wasVirtual = new boolean[1];
|
||||
try (ExecutorService pool = Executors.newVirtualThreadPerTaskExecutor()) {
|
||||
Future<?> future = pool.submit(() -> wasVirtual[0] = Thread.currentThread().isVirtual());
|
||||
future.get(5, TimeUnit.SECONDS);
|
||||
}
|
||||
assertTrue(wasVirtual[0], "newVirtualThreadPerTaskExecutor() must run tasks on virtual threads");
|
||||
}
|
||||
|
||||
@Test
|
||||
void submitDoesNotThrowUntilGetIsCalled() throws InterruptedException {
|
||||
try (ExecutorService pool = Executors.newFixedThreadPool(1)) {
|
||||
Future<?> future = pool.submit(() -> {
|
||||
throw new IllegalStateException("deliberate failure");
|
||||
});
|
||||
Thread.sleep(200);
|
||||
assertFalse(future.isCancelled(), "a failed task is not the same as a cancelled one");
|
||||
ExecutionException thrown = assertThrows(ExecutionException.class, future::get,
|
||||
"the same exception only surfaces once get() is actually called");
|
||||
assertEquals("deliberate failure", thrown.getCause().getMessage());
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user