From bc5153d2a5f28af0f058c8f32e64268346d8dcba Mon Sep 17 00:00:00 2001 From: asmhatre Date: Wed, 30 Sep 2026 06:54:17 +0000 Subject: [PATCH] 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 Claude-Session: https://claude.ai/code/session_01FhzLY5p6okFva3qsnsRyvM --- README.md | 1 + executors/README.md | 72 ++++++++++++ .../output/01-shutdown-vs-shutdownnow.txt | 24 ++++ .../output/02-graceful-shutdown-pattern.txt | 13 +++ .../output/03-autocloseable-executor.txt | 7 ++ .../04-virtual-thread-per-task-executor.txt | 10 ++ executors/output/05-lost-exception-trap.txt | 13 +++ executors/output/06-correctness-tests.txt | 4 + executors/pom.xml | 43 +++++++ executors/scripts/run-all.sh | 33 ++++++ .../executors/AutoCloseableExecutorDemo.java | 36 ++++++ .../GracefulShutdownPatternDemo.java | 70 ++++++++++++ .../executors/LostExceptionTrapDemo.java | 48 ++++++++ .../executors/ShutdownVsShutdownNowDemo.java | 59 ++++++++++ .../VirtualThreadPerTaskExecutorDemo.java | 65 +++++++++++ .../ExecutorShutdownBehaviorTest.java | 108 ++++++++++++++++++ pom.xml | 1 + 17 files changed, 607 insertions(+) create mode 100644 executors/README.md create mode 100644 executors/output/01-shutdown-vs-shutdownnow.txt create mode 100644 executors/output/02-graceful-shutdown-pattern.txt create mode 100644 executors/output/03-autocloseable-executor.txt create mode 100644 executors/output/04-virtual-thread-per-task-executor.txt create mode 100644 executors/output/05-lost-exception-trap.txt create mode 100644 executors/output/06-correctness-tests.txt create mode 100644 executors/pom.xml create mode 100755 executors/scripts/run-all.sh create mode 100644 executors/src/main/java/com/ankurm/executors/AutoCloseableExecutorDemo.java create mode 100644 executors/src/main/java/com/ankurm/executors/GracefulShutdownPatternDemo.java create mode 100644 executors/src/main/java/com/ankurm/executors/LostExceptionTrapDemo.java create mode 100644 executors/src/main/java/com/ankurm/executors/ShutdownVsShutdownNowDemo.java create mode 100644 executors/src/main/java/com/ankurm/executors/VirtualThreadPerTaskExecutorDemo.java create mode 100644 executors/src/test/java/com/ankurm/executors/ExecutorShutdownBehaviorTest.java diff --git a/README.md b/README.md index 7c4f816..426c476 100644 --- a/README.md +++ b/README.md @@ -10,6 +10,7 @@ article; each module's own README has that article's version table, quickstart, | [`atomics`](atomics/) | Java Atomics and VarHandle: CAS, LongAdder, and When Atomics Beat Locks | | [`vt-pinning`](vt-pinning/) | Diagnosing Virtual Thread Pinning in Production: JFR Events, jcmd, and Real Fixes | | [`synchronizers`](synchronizers/) | CountDownLatch vs CyclicBarrier vs Phaser vs Semaphore in Java | +| [`executors`](executors/) | ExecutorService Done Right: Shutdown, try-with-resources and Virtual Thread Executors | ## License diff --git a/executors/README.md b/executors/README.md new file mode 100644 index 0000000..8cbb08e --- /dev/null +++ b/executors/README.md @@ -0,0 +1,72 @@ +# executors + +Companion code for the ankurm.com post *"ExecutorService Done Right: Shutdown, try-with-resources +and Virtual Thread Executors."* Sixth module in `java-core-examples`, the Java-core / concurrency +series. + +Five demos, each isolating one thing people get wrong about `ExecutorService` lifecycle and task +submission: `shutdown()` vs `shutdownNow()`, the two-phase graceful-shutdown pattern, the +`AutoCloseable` try-with-resources shortcut (JDK 19+), what `newVirtualThreadPerTaskExecutor()` +actually costs against a bounded pool doing the same blocking work, and the exception that +`submit()` quietly swallows if nothing ever calls `Future.get()`. + +## Versions this was built and tested against + +| Component | Version | Notes | +|---|---|---| +| JDK | 25.0.4.1+1 (Temurin, LTS) | Every demo here runs on one JDK - no cross-version comparison needed. | +| JUnit Jupiter | 5.11.0 | Correctness tests, 10 repeats for the two timing-sensitive ones. | +| Maven | 3.9.11 | | +| Hardware | 2 vCPU x86-64 VM | Same sandbox as the rest of this series. | + +## Quickstart + +```bash +export JAVA_HOME=/path/to/jdk-25 +mvn package +java -cp target/classes com.ankurm.executors.ShutdownVsShutdownNowDemo +java -cp target/classes com.ankurm.executors.VirtualThreadPerTaskExecutorDemo +``` + +`scripts/run-all.sh` regenerates every file in `output/` (needs `JDK25_HOME`). + +## What's in here + +| File | What it shows | +|---|---| +| `.../ShutdownVsShutdownNowDemo.java` | `shutdown()` finishes everything queued; `shutdownNow()` interrupts what's running and abandons what's still waiting. | +| `.../GracefulShutdownPatternDemo.java` | The two-phase pattern from the `ExecutorService` javadoc: `shutdown()`, bounded `awaitTermination()`, escalate to `shutdownNow()` only if the grace period runs out. | +| `.../AutoCloseableExecutorDemo.java` | Since JDK 19, `ExecutorService` extends `AutoCloseable` - `close()` runs that same two-phase pattern for you, so try-with-resources is a correct one-liner, not a shortcut that skips it. | +| `.../VirtualThreadPerTaskExecutorDemo.java` | 5,000 sleeping tasks against `newVirtualThreadPerTaskExecutor()` vs a 20-thread fixed pool - real wall-clock numbers, not a claim. | +| `.../LostExceptionTrapDemo.java` | `execute()` lets an exception reach the default uncaught-exception handler and print. `submit()` without a later `get()` swallows the identical exception completely. | +| `src/test/.../ExecutorShutdownBehaviorTest.java` | Asserts `shutdown()` drops nothing, `shutdownNow()` abandons exactly the still-queued tasks, try-with-resources blocks until every task is done, the virtual-thread executor really runs on virtual threads, and a `submit()`'d exception only surfaces at `get()`. | +| `output/01`-`05` | Each demo's real run. | +| `output/06` | JUnit correctness run. | + +## Reading the results honestly (2-vCPU sandbox) + +**`shutdownNow()` doesn't guarantee interruption, it requests it.** `output/01` shows both running +tasks stopping mid-`sleep()` because `Thread.sleep()` responds to `interrupt()` - a task that never +checks `Thread.interrupted()` or calls an interruptible blocking method would keep running to +completion regardless of `shutdownNow()`. The 4 queued-but-not-started tasks are handled +differently: `shutdownNow()` returns them directly as a `List` and they never run at all. + +**The virtual-thread number here (`output/04`) is the whole argument, not a summary of it.** 5,000 +tasks sleeping 50ms each finished in 121ms on `newVirtualThreadPerTaskExecutor()` - close to the +cost of one task, because nothing waited for a worker to free up. The identical workload on a +20-thread fixed pool took 12,556ms, in line with the ~250 sequential rounds of 20 that a bounded +pool forces. This isn't "virtual threads are faster" in general - it's that a large number of +*blocking* tasks stop competing for a small thread count once the thread stops being the scarce +resource. + +**The `submit()` silence in `output/05` is real, not a simplification.** The failing task's +`Future.isDone()` reports `true` immediately, and nothing about the program's output changes +whether or not that Future is ever inspected. Compare that block to the `execute()` block right +above it in the same transcript, on the same pool: identical exception, and one prints a full stack +trace unprompted while the other prints nothing at all until something explicitly calls `get()`. +That gap is exactly why a task silently "not running" in production logs is worth checking for an +un-retrieved `Future` before anything else. + +## License + +MIT - see the [repo-wide LICENSE](../LICENSE). diff --git a/executors/output/01-shutdown-vs-shutdownnow.txt b/executors/output/01-shutdown-vs-shutdownnow.txt new file mode 100644 index 0000000..011d82d --- /dev/null +++ b/executors/output/01-shutdown-vs-shutdownnow.txt @@ -0,0 +1,24 @@ +=== shutdown(): lets running AND already-queued tasks finish === +task-0: started +task-1: started +main: called shutdown() - isShutdown=true, isTerminated=false +task-0: finished normally +task-2: started +task-1: finished normally +task-3: started +task-3: finished normally +task-4: started +task-2: finished normally +task-5: started +task-4: finished normally +task-5: finished normally +main: awaitTermination returned true - all 6 tasks ran to completion, none were skipped + +=== shutdownNow(): interrupts running tasks, abandons queued ones === +task-0: started +task-1: started +main: called shutdownNow() - 4 queued tasks abandoned, never started +task-1: interrupted, exiting early +task-0: interrupted, exiting early +main: awaitTermination returned true - the 2 running tasks were interrupted mid-sleep, not allowed to finish +done diff --git a/executors/output/02-graceful-shutdown-pattern.txt b/executors/output/02-graceful-shutdown-pattern.txt new file mode 100644 index 0000000..2051c2c --- /dev/null +++ b/executors/output/02-graceful-shutdown-pattern.txt @@ -0,0 +1,13 @@ +=== Case 1: tasks finish comfortably inside the grace period === +quickPool: shutdown() called, waiting up to 1s for a clean finish +quick-task-1: done (200ms) +quick-task-0: done (200ms) +quickPool: terminated cleanly within the grace period + +=== Case 2: a task outlives the grace period, forcing shutdownNow() === +slowPool: shutdown() called, waiting up to 1s for a clean finish +slow-task: started, will try to sleep 3s +slowPool: grace period expired, escalating to shutdownNow() +slow-task: interrupted during shutdownNow(), stopping early +slowPool: terminated after shutdownNow() +done diff --git a/executors/output/03-autocloseable-executor.txt b/executors/output/03-autocloseable-executor.txt new file mode 100644 index 0000000..7eef580 --- /dev/null +++ b/executors/output/03-autocloseable-executor.txt @@ -0,0 +1,7 @@ +=== try-with-resources: close() runs the graceful pattern automatically === +main: leaving the try-with-resources block now +task-2: finished +task-0: finished +task-1: finished +main: back from the try-with-resources block - the pool is terminated, every task already finished, no separate awaitTermination() call was written +done diff --git a/executors/output/04-virtual-thread-per-task-executor.txt b/executors/output/04-virtual-thread-per-task-executor.txt new file mode 100644 index 0000000..9f26dba --- /dev/null +++ b/executors/output/04-virtual-thread-per-task-executor.txt @@ -0,0 +1,10 @@ +Running 5000 tasks that each sleep 50ms. + +sample task thread: VirtualThread[#22]/runnable@ForkJoinPool-1-worker-1 +newVirtualThreadPerTaskExecutor(): 5000/5000 tasks completed in 121ms +(121ms is close to the 50ms a single task takes - every task ran concurrently, nothing waited for a free worker) + +sample task thread: Thread[#5027,pool-1-thread-1,5,main] +newFixedThreadPool(20): 5000/5000 tasks completed in 12556ms +(12556ms is roughly 250 rounds of 50ms each - only 20 tasks are ever running at once) +done diff --git a/executors/output/05-lost-exception-trap.txt b/executors/output/05-lost-exception-trap.txt new file mode 100644 index 0000000..7f93aff --- /dev/null +++ b/executors/output/05-lost-exception-trap.txt @@ -0,0 +1,13 @@ +=== execute(): uncaught exception reaches the default handler, gets printed === +Exception in thread "pool-1-thread-1" java.lang.RuntimeException: boom from execute() + at com.ankurm.executors.LostExceptionTrapDemo.lambda$main$0(LostExceptionTrapDemo.java:22) + at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1090) + at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:614) + at java.base/java.lang.Thread.run(Thread.java:1474) + +=== submit(): same exception, but nobody calls get() - dead silence === +main: task's Future exists (isDone=true), but its exception was never printed anywhere - it's just sitting inside the Future + +=== submit() + get(): the SAME exception, now surfaced by calling get() === +main: caught ExecutionException, cause = java.lang.RuntimeException: boom from submit(), retrieved via get() +done diff --git a/executors/output/06-correctness-tests.txt b/executors/output/06-correctness-tests.txt new file mode 100644 index 0000000..30b5ed3 --- /dev/null +++ b/executors/output/06-correctness-tests.txt @@ -0,0 +1,4 @@ +------------------------------------------------------------------------------- +Test set: com.ankurm.executors.ExecutorShutdownBehaviorTest +------------------------------------------------------------------------------- +Tests run: 23, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 0.417 s -- in com.ankurm.executors.ExecutorShutdownBehaviorTest diff --git a/executors/pom.xml b/executors/pom.xml new file mode 100644 index 0000000..71f6644 --- /dev/null +++ b/executors/pom.xml @@ -0,0 +1,43 @@ + + + 4.0.0 + + + com.ankurm + java-core-examples + 1.0 + + + executors + executors + ExecutorService done right: shutdown() vs shutdownNow(), the AutoCloseable try-with-resources pattern (JDK 19+), newVirtualThreadPerTaskExecutor(), and the lost-exception trap with submit(). + + + + org.junit.jupiter + junit-jupiter + 5.11.0 + test + + + + + + + org.apache.maven.plugins + maven-compiler-plugin + 3.13.0 + + 25 + + + + org.apache.maven.plugins + maven-surefire-plugin + 3.2.5 + + + + diff --git a/executors/scripts/run-all.sh b/executors/scripts/run-all.sh new file mode 100755 index 0000000..931e582 --- /dev/null +++ b/executors/scripts/run-all.sh @@ -0,0 +1,33 @@ +#!/usr/bin/env bash +# Regenerates every file in ../output/. Requires JDK25_HOME (and JAVA_HOME pointed at it for the +# Maven build). All demos in this module run on a single modern JDK - there's no cross-version +# comparison needed here, unlike vt-pinning or synchronizers. +set -euo pipefail + +if [[ -z "${JDK25_HOME:-}" ]]; then + echo "JDK25_HOME must be set (e.g. /path/to/jdk-25)" >&2 + exit 1 +fi + +cd "$(dirname "$0")/.." +OUT=output +mkdir -p "$OUT" + +JAVA_HOME="$JDK25_HOME" mvn -q -f ../pom.xml -pl executors -am package + +run() { + local class=$1 + local outfile=$2 + echo "==> $class" + "$JDK25_HOME/bin/java" -cp target/classes "com.ankurm.executors.$class" 2>&1 | grep -v "Picked up" > "$OUT/$outfile" +} + +run ShutdownVsShutdownNowDemo 01-shutdown-vs-shutdownnow.txt +run GracefulShutdownPatternDemo 02-graceful-shutdown-pattern.txt +run AutoCloseableExecutorDemo 03-autocloseable-executor.txt +run VirtualThreadPerTaskExecutorDemo 04-virtual-thread-per-task-executor.txt +run LostExceptionTrapDemo 05-lost-exception-trap.txt + +cp target/surefire-reports/com.ankurm.executors.ExecutorShutdownBehaviorTest.txt "$OUT/06-correctness-tests.txt" + +echo "Done. See $OUT/" diff --git a/executors/src/main/java/com/ankurm/executors/AutoCloseableExecutorDemo.java b/executors/src/main/java/com/ankurm/executors/AutoCloseableExecutorDemo.java new file mode 100644 index 0000000..6223cdb --- /dev/null +++ b/executors/src/main/java/com/ankurm/executors/AutoCloseableExecutorDemo.java @@ -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"); + } +} diff --git a/executors/src/main/java/com/ankurm/executors/GracefulShutdownPatternDemo.java b/executors/src/main/java/com/ankurm/executors/GracefulShutdownPatternDemo.java new file mode 100644 index 0000000..14340b5 --- /dev/null +++ b/executors/src/main/java/com/ankurm/executors/GracefulShutdownPatternDemo.java @@ -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"); + } +} diff --git a/executors/src/main/java/com/ankurm/executors/LostExceptionTrapDemo.java b/executors/src/main/java/com/ankurm/executors/LostExceptionTrapDemo.java new file mode 100644 index 0000000..aa6a71f --- /dev/null +++ b/executors/src/main/java/com/ankurm/executors/LostExceptionTrapDemo.java @@ -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"); + } +} diff --git a/executors/src/main/java/com/ankurm/executors/ShutdownVsShutdownNowDemo.java b/executors/src/main/java/com/ankurm/executors/ShutdownVsShutdownNowDemo.java new file mode 100644 index 0000000..1f1698d --- /dev/null +++ b/executors/src/main/java/com/ankurm/executors/ShutdownVsShutdownNowDemo.java @@ -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 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"); + } +} diff --git a/executors/src/main/java/com/ankurm/executors/VirtualThreadPerTaskExecutorDemo.java b/executors/src/main/java/com/ankurm/executors/VirtualThreadPerTaskExecutorDemo.java new file mode 100644 index 0000000..de3d8c9 --- /dev/null +++ b/executors/src/main/java/com/ankurm/executors/VirtualThreadPerTaskExecutorDemo.java @@ -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"); + } +} diff --git a/executors/src/test/java/com/ankurm/executors/ExecutorShutdownBehaviorTest.java b/executors/src/test/java/com/ankurm/executors/ExecutorShutdownBehaviorTest.java new file mode 100644 index 0000000..45f5581 --- /dev/null +++ b/executors/src/test/java/com/ankurm/executors/ExecutorShutdownBehaviorTest.java @@ -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 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()); + } + } +} diff --git a/pom.xml b/pom.xml index ecdec4f..bcefd67 100644 --- a/pom.xml +++ b/pom.xml @@ -18,6 +18,7 @@ atomics vt-pinning synchronizers + executors