From 55ee850a75d4cd4af18c508c078eb62993679f8b Mon Sep 17 00:00:00 2001 From: asmhatre Date: Wed, 30 Sep 2026 06:43:52 +0000 Subject: [PATCH] synchronizers: CountDownLatch vs CyclicBarrier vs Phaser vs Semaphore One three-round worker pipeline implemented with CountDownLatch (one latch per round), CyclicBarrier (reusable, with a barrier action), and Phaser (reusable, plus dynamic mid-run registration), alongside Semaphore solving the genuinely different problem of bounding concurrent access. Covers timeout behavior (CyclicBarrier permanently breaks on one timeout, Phaser doesn't) and virtual-thread compatibility, including the corrected finding that Object.wait() releases its monitor and never pinned, unlike Thread.sleep() inside synchronized. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01FhzLY5p6okFva3qsnsRyvM --- README.md | 1 + pom.xml | 1 + synchronizers/README.md | 80 ++++++++++ .../output/01-countdownlatch-pipeline.txt | 17 +++ .../output/02-cyclicbarrier-pipeline.txt | 18 +++ synchronizers/output/03-phaser-pipeline.txt | 21 +++ .../output/04-semaphore-resource-guard.txt | 16 ++ synchronizers/output/05-timeout-behavior.txt | 20 +++ .../06-virtual-thread-compatibility.txt | 25 ++++ .../output/07-round-ordering-correctness.txt | 4 + synchronizers/pom.xml | 43 ++++++ synchronizers/scripts/run-all.sh | 69 +++++++++ .../synchronizers/CountDownLatchPipeline.java | 70 +++++++++ .../synchronizers/CyclicBarrierPipeline.java | 54 +++++++ .../ankurm/synchronizers/PhaserPipeline.java | 68 +++++++++ .../synchronizers/SemaphoreResourceGuard.java | 74 ++++++++++ .../SynchronizedWaitPinsDemo.java | 59 ++++++++ .../synchronizers/TimeoutBehaviorDemo.java | 89 ++++++++++++ .../VirtualThreadCompatibilityDemo.java | 96 ++++++++++++ .../RoundOrderingCorrectnessTest.java | 137 ++++++++++++++++++ 20 files changed, 962 insertions(+) create mode 100644 synchronizers/README.md create mode 100644 synchronizers/output/01-countdownlatch-pipeline.txt create mode 100644 synchronizers/output/02-cyclicbarrier-pipeline.txt create mode 100644 synchronizers/output/03-phaser-pipeline.txt create mode 100644 synchronizers/output/04-semaphore-resource-guard.txt create mode 100644 synchronizers/output/05-timeout-behavior.txt create mode 100644 synchronizers/output/06-virtual-thread-compatibility.txt create mode 100644 synchronizers/output/07-round-ordering-correctness.txt create mode 100644 synchronizers/pom.xml create mode 100755 synchronizers/scripts/run-all.sh create mode 100644 synchronizers/src/main/java/com/ankurm/synchronizers/CountDownLatchPipeline.java create mode 100644 synchronizers/src/main/java/com/ankurm/synchronizers/CyclicBarrierPipeline.java create mode 100644 synchronizers/src/main/java/com/ankurm/synchronizers/PhaserPipeline.java create mode 100644 synchronizers/src/main/java/com/ankurm/synchronizers/SemaphoreResourceGuard.java create mode 100644 synchronizers/src/main/java/com/ankurm/synchronizers/SynchronizedWaitPinsDemo.java create mode 100644 synchronizers/src/main/java/com/ankurm/synchronizers/TimeoutBehaviorDemo.java create mode 100644 synchronizers/src/main/java/com/ankurm/synchronizers/VirtualThreadCompatibilityDemo.java create mode 100644 synchronizers/src/test/java/com/ankurm/synchronizers/RoundOrderingCorrectnessTest.java diff --git a/README.md b/README.md index 243cb86..7c4f816 100644 --- a/README.md +++ b/README.md @@ -9,6 +9,7 @@ article; each module's own README has that article's version table, quickstart, | [`locks`](locks/) | synchronized vs ReentrantLock vs StampedLock: Benchmarks and a Decision Table | | [`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 | ## License diff --git a/pom.xml b/pom.xml index 0f00dba..ecdec4f 100644 --- a/pom.xml +++ b/pom.xml @@ -17,6 +17,7 @@ locks atomics vt-pinning + synchronizers diff --git a/synchronizers/README.md b/synchronizers/README.md new file mode 100644 index 0000000..7ea23a9 --- /dev/null +++ b/synchronizers/README.md @@ -0,0 +1,80 @@ +# synchronizers + +Companion code for the ankurm.com post *"CountDownLatch vs CyclicBarrier vs Phaser vs Semaphore +in Java."* Fifth module in `java-core-examples`, the Java-core / concurrency series. + +The same three-round worker pipeline, implemented once per class where it's a natural fit +(`CountDownLatch`, `CyclicBarrier`, `Phaser`), plus `Semaphore` solving the genuinely different +problem it actually solves (bounding concurrent access, not waiting for everyone to arrive) inside +the same pipeline shape. Reuse semantics, timeout behavior, and virtual-thread compatibility are +each demonstrated directly rather than asserted. + +## Versions this was built and tested against + +| Component | Version | Notes | +|---|---|---| +| JDK (primary) | 25.0.4.1+1 (Temurin, LTS) | | +| JDK (comparison) | 21.0.10 (OpenJDK) | For the virtual-thread pinning contrast only. | +| JUnit Jupiter | 5.11.0 | Correctness tests, 20 repeats per barrier test. | +| 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.synchronizers.CyclicBarrierPipeline +java -cp target/classes com.ankurm.synchronizers.PhaserPipeline +java -cp target/classes com.ankurm.synchronizers.SemaphoreResourceGuard +``` + +`scripts/run-all.sh` regenerates every file in `output/` (needs `JDK25_HOME` and `JDK21_HOME`). + +## What's in here + +| File | What it shows | +|---|---| +| `.../CountDownLatchPipeline.java` | A starting-gate latch, plus one fresh latch per round - the workaround `CyclicBarrier` exists to remove. | +| `.../CyclicBarrierPipeline.java` | One reusable barrier across all three rounds, with a barrier action. | +| `.../PhaserPipeline.java` | Same barrier behavior, plus a fifth worker dynamically registering after round 1. | +| `.../SemaphoreResourceGuard.java` | The same pipeline, but `Semaphore` bounding concurrent access to a shared resource - a different concern from the other three. | +| `.../TimeoutBehaviorDemo.java` | What each class does when the thing it's waiting for never arrives. | +| `.../VirtualThreadCompatibilityDemo.java` | None of the four pin a virtual thread, even on JDK 21, because none are built on `synchronized`. | +| `.../SynchronizedWaitPinsDemo.java` | The contrast case, and a correction: `synchronized` + `Object.wait()` does NOT pin (it releases the monitor); `synchronized` + `Thread.sleep()` does. | +| `src/test/.../RoundOrderingCorrectnessTest.java` | Asserts no worker's round N+1 ever starts before every worker's round N finished (`CyclicBarrier`, `Phaser`, 20 repeats each), and `Semaphore` never exceeds its permit count. | +| `output/01`-`04` | Each pipeline's real run. | +| `output/05` | Timeout behavior across all four. | +| `output/06` | Virtual-thread compatibility, JDK 21, both the four synchronizers and the `synchronized`+`wait`/`sleep` contrast. | +| `output/07` | JUnit correctness run. | + +## Reading the results honestly (2-vCPU sandbox) + +**`Semaphore` isn't a fourth way to do what the other three do.** `CountDownLatch`, `CyclicBarrier`, +and `Phaser` all answer "has everyone reached this point?" `Semaphore` answers "how many callers +are allowed through at once?" - a genuinely different question, not a stylistic alternative. +`output/04` shows `Semaphore` bounding a resource to 2 concurrent callers, independent of round +boundaries or worker identity; forcing it into a "wait for the group" role would misrepresent what +it's for. + +**A timeout doesn't do the same thing on all four** (`output/05`). `CountDownLatch.await(timeout)` +returns `false` and stays reusable to check again. `Semaphore.tryAcquire(timeout)` behaves the same +way. `CyclicBarrier.await(timeout)` throws `TimeoutException` and **permanently breaks the barrier +for every other waiter** until an explicit `reset()` - one straggler poisons the whole group. +`Phaser.awaitAdvanceInterruptibly(phase, timeout, unit)` throws `TimeoutException` too, but the +phaser itself is untouched and a later, fully-arrived round on the same instance succeeds normally +with no reset needed. That difference alone is a reason to reach for `Phaser` over `CyclicBarrier` +in anything where a slow round shouldn't take down every other round after it. + +**`Object.wait()` doesn't pin a virtual thread - `Thread.sleep()` inside the same `synchronized` +block does** (`output/06`), and this repo initially assumed otherwise before checking. `wait()` +fully releases the monitor for the duration of the wait, so there's nothing held while parked and +nothing to pin; JEP 491 (JDK 24) only needed to fix the case where a blocking call runs *while +still holding* the monitor. All four `java.util.concurrent` classes in this module were already +safe on virtual threads on JDK 21, years before JEP 491, because none of them are built on +`synchronized` in the first place - a fact worth knowing before assuming every pre-JDK-24 +concurrency primitive needed the same fix. + +## License + +MIT - see the [repo-wide LICENSE](../LICENSE). diff --git a/synchronizers/output/01-countdownlatch-pipeline.txt b/synchronizers/output/01-countdownlatch-pipeline.txt new file mode 100644 index 0000000..13aa77e --- /dev/null +++ b/synchronizers/output/01-countdownlatch-pipeline.txt @@ -0,0 +1,17 @@ +$ java CountDownLatchPipeline + +main: all 4 workers created, still waiting at the start gate +main: releasing the start gate +worker-0: finished round 1 +worker-1: finished round 1 +worker-2: finished round 1 +worker-3: finished round 1 +worker-0: finished round 2 +worker-1: finished round 2 +worker-2: finished round 2 +worker-3: finished round 2 +worker-0: finished round 3 +worker-1: finished round 3 +worker-2: finished round 3 +worker-3: finished round 3 +done diff --git a/synchronizers/output/02-cyclicbarrier-pipeline.txt b/synchronizers/output/02-cyclicbarrier-pipeline.txt new file mode 100644 index 0000000..3f186ca --- /dev/null +++ b/synchronizers/output/02-cyclicbarrier-pipeline.txt @@ -0,0 +1,18 @@ +$ java CyclicBarrierPipeline + +worker-0: finished round 1 +worker-1: finished round 1 +worker-2: finished round 1 +worker-3: finished round 1 +barrier action: round 1 complete, all 4 workers arrived +worker-0: finished round 2 +worker-1: finished round 2 +worker-2: finished round 2 +worker-3: finished round 2 +barrier action: round 2 complete, all 4 workers arrived +worker-0: finished round 3 +worker-1: finished round 3 +worker-2: finished round 3 +worker-3: finished round 3 +barrier action: round 3 complete, all 4 workers arrived +done diff --git a/synchronizers/output/03-phaser-pipeline.txt b/synchronizers/output/03-phaser-pipeline.txt new file mode 100644 index 0000000..c380782 --- /dev/null +++ b/synchronizers/output/03-phaser-pipeline.txt @@ -0,0 +1,21 @@ +$ java PhaserPipeline + +worker-0: finished round 1 +worker-1: finished round 1 +worker-2: finished round 1 +worker-3: finished round 1 +phaser onAdvance: phase 1 complete, 4 registered parties -> advancing +worker-4: registers after round 1 (phaser now has 5 parties) +worker-0: finished round 2 +worker-1: finished round 2 +worker-2: finished round 2 +worker-3: finished round 2 +worker-4: finished round 2 +phaser onAdvance: phase 2 complete, 5 registered parties -> advancing +worker-0: finished round 3 +worker-1: finished round 3 +worker-2: finished round 3 +worker-3: finished round 3 +worker-4: finished round 3 +phaser onAdvance: phase 3 complete, 5 registered parties -> advancing +done diff --git a/synchronizers/output/04-semaphore-resource-guard.txt b/synchronizers/output/04-semaphore-resource-guard.txt new file mode 100644 index 0000000..8ea883f --- /dev/null +++ b/synchronizers/output/04-semaphore-resource-guard.txt @@ -0,0 +1,16 @@ +$ java SemaphoreResourceGuard + +worker-0 round 1: acquired permit, 1 callers concurrently in the resource (availablePermits=0) +worker-1 round 1: acquired permit, 2 callers concurrently in the resource (availablePermits=0) +worker-2 round 1: acquired permit, 2 callers concurrently in the resource (availablePermits=0) +worker-3 round 1: acquired permit, 2 callers concurrently in the resource (availablePermits=0) +worker-3 round 2: acquired permit, 1 callers concurrently in the resource (availablePermits=1) +worker-0 round 2: acquired permit, 2 callers concurrently in the resource (availablePermits=0) +worker-1 round 2: acquired permit, 1 callers concurrently in the resource (availablePermits=1) +worker-2 round 2: acquired permit, 2 callers concurrently in the resource (availablePermits=0) +worker-2 round 3: acquired permit, 1 callers concurrently in the resource (availablePermits=1) +worker-3 round 3: acquired permit, 2 callers concurrently in the resource (availablePermits=0) +worker-0 round 3: acquired permit, 1 callers concurrently in the resource (availablePermits=1) +worker-1 round 3: acquired permit, 2 callers concurrently in the resource (availablePermits=0) +permits=2 maxObservedConcurrency=2 +done diff --git a/synchronizers/output/05-timeout-behavior.txt b/synchronizers/output/05-timeout-behavior.txt new file mode 100644 index 0000000..1135830 --- /dev/null +++ b/synchronizers/output/05-timeout-behavior.txt @@ -0,0 +1,20 @@ +$ java TimeoutBehaviorDemo + +=== CountDownLatch.await(timeout): expects 2, only 1 ever counts down === +await returned: false (false = timed out, latch still usable to re-check later) +getCount() after timeout: 1 + +=== CyclicBarrier.await(timeout): expects 2 parties, only 1 arrives === +await threw TimeoutException, as expected +barrier.isBroken() after the timeout: true +a second, later await() on the SAME barrier: BrokenBarrierException immediately - one timeout poisons the barrier for everyone until reset() + +=== Phaser.awaitAdvanceInterruptibly(phase, timeout): expects 2 parties, only 1 arrives === +awaitAdvanceInterruptibly threw TimeoutException, as expected +phaser.getPhase() after the timeout: 0 (unchanged - still phase 0, not advanced or broken) +a later round on the SAME phaser, once both parties actually arrive: succeeds normally, now at phase 1 + +=== Semaphore.tryAcquire(timeout): 0 permits available === +tryAcquire returned: false (false = timed out, no exception, no broken state) +a later tryAcquire after a release: true - fully independent of the earlier timeout +done diff --git a/synchronizers/output/06-virtual-thread-compatibility.txt b/synchronizers/output/06-virtual-thread-compatibility.txt new file mode 100644 index 0000000..482724f --- /dev/null +++ b/synchronizers/output/06-virtual-thread-compatibility.txt @@ -0,0 +1,25 @@ +$ java -Djdk.tracePinnedThreads=full VirtualThreadCompatibilityDemo (pre-JDK-24 runtime) + +java.version=21.0.10 +jdk.tracePinnedThreads=full +CountDownLatch.await(): completed, no trace above this line for it +CyclicBarrier.await(): completed, no trace above this line for it +Phaser.arriveAndAwaitAdvance(): completed, no trace above this line for it +Semaphore.acquire(): completed, no trace above this line for it +done - no pinning trace printed above for any of the four + +$ java -Djdk.tracePinnedThreads=full SynchronizedWaitPinsDemo (pre-JDK-24 runtime, contrast case) + +java.version=21.0.10 +jdk.tracePinnedThreads=full +synchronized + Object.wait(): completed, no pinning trace above this line for it +VirtualThread[#21]/runnable@ForkJoinPool-1-worker-1 reason:MONITOR + java.base/java.lang.VirtualThread$VThreadContinuation.onPinned(VirtualThread.java:199) + java.base/jdk.internal.vm.Continuation.onPinned0(Continuation.java:393) + java.base/java.lang.VirtualThread.parkNanos(VirtualThread.java:635) + java.base/java.lang.VirtualThread.sleepNanos(VirtualThread.java:807) + java.base/java.lang.Thread.sleep(Thread.java:507) + com.ankurm.synchronizers.SynchronizedWaitPinsDemo.lambda$main$1(SynchronizedWaitPinsDemo.java:48) <== monitors:1 + java.base/java.lang.VirtualThread.run(VirtualThread.java:329) +synchronized + Thread.sleep(): completed, pinning trace (if any) appears above this line for it +done diff --git a/synchronizers/output/07-round-ordering-correctness.txt b/synchronizers/output/07-round-ordering-correctness.txt new file mode 100644 index 0000000..fe52a57 --- /dev/null +++ b/synchronizers/output/07-round-ordering-correctness.txt @@ -0,0 +1,4 @@ +------------------------------------------------------------------------------- +Test set: com.ankurm.synchronizers.RoundOrderingCorrectnessTest +------------------------------------------------------------------------------- +Tests run: 41, Failures: 0, Errors: 0, Skipped: 0, Time elapsed: 0.514 s -- in com.ankurm.synchronizers.RoundOrderingCorrectnessTest diff --git a/synchronizers/pom.xml b/synchronizers/pom.xml new file mode 100644 index 0000000..6d2debc --- /dev/null +++ b/synchronizers/pom.xml @@ -0,0 +1,43 @@ + + + 4.0.0 + + + com.ankurm + java-core-examples + 1.0 + + + synchronizers + synchronizers + CountDownLatch vs CyclicBarrier vs Phaser vs Semaphore: one three-round worker pipeline implemented four ways, reuse semantics, timeout behavior, and virtual-thread compatibility. + + + + 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/synchronizers/scripts/run-all.sh b/synchronizers/scripts/run-all.sh new file mode 100755 index 0000000..d8c4515 --- /dev/null +++ b/synchronizers/scripts/run-all.sh @@ -0,0 +1,69 @@ +#!/usr/bin/env bash +# Regenerates every file in output/. Requires: +# - JDK25_HOME pointing at a JDK 24+ install (this repo used Temurin 25.0.4.1+1) +# - JDK21_HOME pointing at a pre-JDK-24 install (this repo used OpenJDK 21.0.10), for the +# virtual-thread pinning contrast at the end +set -euo pipefail +cd "$(dirname "$0")/.." + +: "${JDK25_HOME:?Set JDK25_HOME to a JDK 24+ install}" +: "${JDK21_HOME:?Set JDK21_HOME to a pre-JDK-24 install}" + +mkdir -p target/classes target/classes-21 +"$JDK25_HOME/bin/javac" --release 25 -d target/classes $(find src/main/java -name "*.java") +"$JDK21_HOME/bin/javac" --release 21 -d target/classes-21 $(find src/main/java -name "*.java") + +echo "--- 01: CountDownLatch pipeline ---" +{ + echo '$ java CountDownLatchPipeline' + echo "" + "$JDK25_HOME/bin/java" -cp target/classes com.ankurm.synchronizers.CountDownLatchPipeline 2>&1 | grep -v "Picked up" +} > output/01-countdownlatch-pipeline.txt + +echo "--- 02: CyclicBarrier pipeline ---" +{ + echo '$ java CyclicBarrierPipeline' + echo "" + "$JDK25_HOME/bin/java" -cp target/classes com.ankurm.synchronizers.CyclicBarrierPipeline 2>&1 | grep -v "Picked up" +} > output/02-cyclicbarrier-pipeline.txt + +echo "--- 03: Phaser pipeline (dynamic registration) ---" +{ + echo '$ java PhaserPipeline' + echo "" + "$JDK25_HOME/bin/java" -cp target/classes com.ankurm.synchronizers.PhaserPipeline 2>&1 | grep -v "Picked up" +} > output/03-phaser-pipeline.txt + +echo "--- 04: Semaphore resource guard ---" +{ + echo '$ java SemaphoreResourceGuard' + echo "" + "$JDK25_HOME/bin/java" -cp target/classes com.ankurm.synchronizers.SemaphoreResourceGuard 2>&1 | grep -v "Picked up" +} > output/04-semaphore-resource-guard.txt + +echo "--- 05: timeout behavior across all four ---" +{ + echo '$ java TimeoutBehaviorDemo' + echo "" + "$JDK25_HOME/bin/java" -cp target/classes com.ankurm.synchronizers.TimeoutBehaviorDemo 2>&1 | grep -v "Picked up" +} > output/05-timeout-behavior.txt + +echo "--- 06: virtual-thread compatibility, all four, on JDK21 (pre-JEP-491) ---" +{ + echo '$ java -Djdk.tracePinnedThreads=full VirtualThreadCompatibilityDemo (pre-JDK-24 runtime)' + echo "" + "$JDK21_HOME/bin/java" -Djdk.tracePinnedThreads=full -cp target/classes-21 \ + com.ankurm.synchronizers.VirtualThreadCompatibilityDemo 2>&1 | grep -v "Picked up" + echo "" + echo '$ java -Djdk.tracePinnedThreads=full SynchronizedWaitPinsDemo (pre-JDK-24 runtime, contrast case)' + echo "" + "$JDK21_HOME/bin/java" -Djdk.tracePinnedThreads=full -cp target/classes-21 \ + com.ankurm.synchronizers.SynchronizedWaitPinsDemo 2>&1 | grep -v "Picked up" +} > output/06-virtual-thread-compatibility.txt + +echo "--- 07: correctness tests ---" +"$JDK25_HOME/bin/java" --version > /dev/null # sanity +(cd "$(pwd)" && JAVA_HOME="$JDK25_HOME" mvn -q -f pom.xml test 2>&1 | grep -v "Picked up\|WARNING" || true) +cp target/surefire-reports/com.ankurm.synchronizers.RoundOrderingCorrectnessTest.txt output/07-round-ordering-correctness.txt + +echo "Done. See output/." diff --git a/synchronizers/src/main/java/com/ankurm/synchronizers/CountDownLatchPipeline.java b/synchronizers/src/main/java/com/ankurm/synchronizers/CountDownLatchPipeline.java new file mode 100644 index 0000000..40c4ac7 --- /dev/null +++ b/synchronizers/src/main/java/com/ankurm/synchronizers/CountDownLatchPipeline.java @@ -0,0 +1,70 @@ +package com.ankurm.synchronizers; + +import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.TimeUnit; + +/** + * The same three-round worker pipeline every class in this module runs, done with + * {@link CountDownLatch}. Two uses, both idiomatic: + * + *
    + *
  1. A single "starting gate" latch releases all workers at once instead of letting them + * race to start as soon as each is spawned.
  2. + *
  3. One NEW latch per round enforces "everyone finishes round N before round N+1 starts" - + * because a {@code CountDownLatch} counts down to zero exactly once and cannot be reset. + * That's the limitation {@link CyclicBarrierPipeline} exists to remove.
  4. + *
+ */ +public final class CountDownLatchPipeline { + + private static final int WORKERS = 4; + private static final int ROUNDS = 3; + + public static void main(String[] args) throws Exception { + List log = new CopyOnWriteArrayList<>(); + CountDownLatch startGate = new CountDownLatch(1); + + Thread[] workers = new Thread[WORKERS]; + for (int w = 0; w < WORKERS; w++) { + int id = w; + workers[w] = new Thread(() -> runWorker(id, startGate, log)); + workers[w].start(); + } + + log.add("main: all " + WORKERS + " workers created, still waiting at the start gate"); + Thread.sleep(50); // let them actually reach await() before we release them + log.add("main: releasing the start gate"); + startGate.countDown(); + + for (Thread t : workers) { + t.join(); + } + log.forEach(System.out::println); + System.out.println("done"); + } + + private static void runWorker(int id, CountDownLatch startGate, List log) { + try { + startGate.await(); + CountDownLatch roundLatch = null; + for (int round = 1; round <= ROUNDS; round++) { + // A fresh latch every round - CountDownLatch cannot be reused once it hits zero. + roundLatch = ROUND_LATCHES[round - 1]; + Thread.sleep(20 + (id * 5)); // simulate uneven round work + log.add("worker-" + id + ": finished round " + round); + roundLatch.countDown(); + roundLatch.await(5, TimeUnit.SECONDS); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + + // One latch per round, sized to the worker count, created up front so every worker sees the + // same instances. In real code this bookkeeping is exactly what CyclicBarrier automates. + private static final CountDownLatch[] ROUND_LATCHES = { + new CountDownLatch(WORKERS), new CountDownLatch(WORKERS), new CountDownLatch(WORKERS) + }; +} diff --git a/synchronizers/src/main/java/com/ankurm/synchronizers/CyclicBarrierPipeline.java b/synchronizers/src/main/java/com/ankurm/synchronizers/CyclicBarrierPipeline.java new file mode 100644 index 0000000..e818418 --- /dev/null +++ b/synchronizers/src/main/java/com/ankurm/synchronizers/CyclicBarrierPipeline.java @@ -0,0 +1,54 @@ +package com.ankurm.synchronizers; + +import java.util.List; +import java.util.concurrent.BrokenBarrierException; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CyclicBarrier; + +/** + * The same three-round pipeline as {@link CountDownLatchPipeline}, but with the one thing a + * {@code CountDownLatch} can't do: the barrier resets itself automatically after every round, so + * one {@code CyclicBarrier} instance - not one per round - carries the whole pipeline. It also + * runs an optional barrier action exactly once per round, on whichever thread happens to be the + * last to arrive, useful for round-boundary bookkeeping that shouldn't run once per worker. + */ +public final class CyclicBarrierPipeline { + + private static final int WORKERS = 4; + private static final int ROUNDS = 3; + + public static void main(String[] args) throws Exception { + List log = new CopyOnWriteArrayList<>(); + int[] roundCounter = {0}; + + CyclicBarrier barrier = new CyclicBarrier(WORKERS, () -> { + roundCounter[0]++; + log.add("barrier action: round " + roundCounter[0] + " complete, all " + WORKERS + " workers arrived"); + }); + + Thread[] workers = new Thread[WORKERS]; + for (int w = 0; w < WORKERS; w++) { + int id = w; + workers[w] = new Thread(() -> runWorker(id, barrier, log)); + workers[w].start(); + } + for (Thread t : workers) { + t.join(); + } + + log.forEach(System.out::println); + System.out.println("done"); + } + + private static void runWorker(int id, CyclicBarrier barrier, List log) { + try { + for (int round = 1; round <= ROUNDS; round++) { + Thread.sleep(20 + (id * 5)); + log.add("worker-" + id + ": finished round " + round); + barrier.await(); // same barrier object, every round - it resets on its own + } + } catch (InterruptedException | BrokenBarrierException e) { + Thread.currentThread().interrupt(); + } + } +} diff --git a/synchronizers/src/main/java/com/ankurm/synchronizers/PhaserPipeline.java b/synchronizers/src/main/java/com/ankurm/synchronizers/PhaserPipeline.java new file mode 100644 index 0000000..c086138 --- /dev/null +++ b/synchronizers/src/main/java/com/ankurm/synchronizers/PhaserPipeline.java @@ -0,0 +1,68 @@ +package com.ankurm.synchronizers; + +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.Phaser; + +/** + * The same pipeline again, with {@link Phaser} - which does everything + * {@link CyclicBarrierPipeline} does (a reusable, self-resetting barrier with a per-phase + * action, via {@link Phaser#onAdvance}), plus the one thing neither {@code CountDownLatch} nor + * {@code CyclicBarrier} can: a party count fixed at construction time. Here, a fifth worker + * {@link Phaser#register()}s itself only after round 1 has already completed, and the barrier + * for round 2 correctly waits for all five - something that requires building a brand new + * {@code CyclicBarrier} (and re-pointing every existing worker at it) to do with that class. + */ +public final class PhaserPipeline { + + private static final int INITIAL_WORKERS = 4; + private static final int ROUNDS = 3; + + public static void main(String[] args) throws Exception { + List log = new CopyOnWriteArrayList<>(); + + Phaser phaser = new Phaser(INITIAL_WORKERS) { + @Override + protected boolean onAdvance(int phase, int registeredParties) { + log.add("phaser onAdvance: phase " + (phase + 1) + " complete, " + + registeredParties + " registered parties -> advancing"); + return phase + 1 >= ROUNDS || registeredParties == 0; + } + }; + + Thread[] workers = new Thread[INITIAL_WORKERS]; + for (int w = 0; w < INITIAL_WORKERS; w++) { + int id = w; + workers[w] = new Thread(() -> runWorker(id, phaser, log, 1)); + workers[w].start(); + } + for (Thread t : workers) { + t.join(); + } + + log.forEach(System.out::println); + System.out.println("done"); + } + + private static void runWorker(int id, Phaser phaser, List log, int startRound) { + try { + for (int round = startRound; round <= ROUNDS; round++) { + Thread.sleep(20 + (id * 5)); + log.add("worker-" + id + ": finished round " + round); + phaser.arriveAndAwaitAdvance(); + + // The late-joining worker: registers itself right after round 1 finishes, then + // runs rounds 2 and 3 like everyone else. Nobody else needed to change. + if (id == 0 && round == 1) { + int lateId = 4; + phaser.register(); + log.add("worker-" + lateId + ": registers after round 1 (phaser now has " + + phaser.getRegisteredParties() + " parties)"); + new Thread(() -> runWorker(lateId, phaser, log, 2)).start(); + } + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } +} diff --git a/synchronizers/src/main/java/com/ankurm/synchronizers/SemaphoreResourceGuard.java b/synchronizers/src/main/java/com/ankurm/synchronizers/SemaphoreResourceGuard.java new file mode 100644 index 0000000..cd7f2c3 --- /dev/null +++ b/synchronizers/src/main/java/com/ankurm/synchronizers/SemaphoreResourceGuard.java @@ -0,0 +1,74 @@ +package com.ankurm.synchronizers; + +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.Semaphore; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * {@link Semaphore} solves a genuinely different problem than the other three classes in this + * module, even though it gets listed alongside them constantly. {@code CountDownLatch}, + * {@code CyclicBarrier}, and {@code Phaser} all answer "has everyone reached this point yet?" - + * {@code Semaphore} answers "how many callers are allowed through this point at once?" Nothing + * here stops multiple threads from acquiring at wildly different times; it just caps how many + * hold a permit simultaneously. + * + *

Same three-round pipeline shape as the other demos, but this time each worker also needs a + * permit from a 2-permit {@code Semaphore} before it can call a simulated expensive shared + * resource (a rate-limited downstream API, a connection pool) once per round - independent of + * which round the barrier logic is in, and independent of worker identity. + */ +public final class SemaphoreResourceGuard { + + private static final int WORKERS = 4; + private static final int ROUNDS = 3; + private static final int PERMITS = 2; + + public static void main(String[] args) throws Exception { + List log = new CopyOnWriteArrayList<>(); + Semaphore resourceGuard = new Semaphore(PERMITS, true); // fair, for a predictable transcript + AtomicInteger concurrentInResource = new AtomicInteger(); + AtomicInteger maxObservedConcurrency = new AtomicInteger(); + CyclicBarrier roundBarrier = new CyclicBarrier(WORKERS); + + Thread[] workers = new Thread[WORKERS]; + for (int w = 0; w < WORKERS; w++) { + int id = w; + workers[w] = new Thread(() -> runWorker( + id, resourceGuard, roundBarrier, concurrentInResource, maxObservedConcurrency, log)); + workers[w].start(); + } + for (Thread t : workers) { + t.join(); + } + + log.forEach(System.out::println); + System.out.println("permits=" + PERMITS + " maxObservedConcurrency=" + maxObservedConcurrency.get()); + System.out.println("done"); + } + + private static void runWorker(int id, Semaphore resourceGuard, CyclicBarrier roundBarrier, + AtomicInteger concurrentInResource, AtomicInteger maxObservedConcurrency, + List log) { + try { + for (int round = 1; round <= ROUNDS; round++) { + resourceGuard.acquire(); + try { + int inFlight = concurrentInResource.incrementAndGet(); + maxObservedConcurrency.updateAndGet(max -> Math.max(max, inFlight)); + log.add("worker-" + id + " round " + round + ": acquired permit, " + + inFlight + " callers concurrently in the resource" + + " (availablePermits=" + resourceGuard.availablePermits() + ")"); + Thread.sleep(15); + } finally { + concurrentInResource.decrementAndGet(); + resourceGuard.release(); + } + roundBarrier.await(); // unrelated barrier concern, same as the other demos + } + } catch (Exception e) { + Thread.currentThread().interrupt(); + } + } +} diff --git a/synchronizers/src/main/java/com/ankurm/synchronizers/SynchronizedWaitPinsDemo.java b/synchronizers/src/main/java/com/ankurm/synchronizers/SynchronizedWaitPinsDemo.java new file mode 100644 index 0000000..a44a2f7 --- /dev/null +++ b/synchronizers/src/main/java/com/ankurm/synchronizers/SynchronizedWaitPinsDemo.java @@ -0,0 +1,59 @@ +package com.ankurm.synchronizers; + +/** + * The contrast case for {@link VirtualThreadCompatibilityDemo} - and a correction to a common + * assumption along the way. The naive expectation is that hand-rolling "wait for a signal" with + * {@code synchronized} + {@code Object.wait()}/{@code notifyAll()} instead of a + * {@code java.util.concurrent} class would pin, the same way {@code synchronized} + a blocking + * call pinned before JEP 491. It doesn't: {@code Object.wait()} fully releases the monitor for + * the duration of the wait, so there's no lock held while parked, and nothing to pin. What + * genuinely pins on a pre-JDK-24 runtime is a blocking call made while still holding the + * monitor - {@code Thread.sleep()} inside {@code synchronized}, exactly like + * {@code MonitorPinningDemo} in the {@code vt-pinning} module. Both cases run here, back to back, + * under the same trace flag, so the difference is visible in one transcript instead of asserted. + */ +public final class SynchronizedWaitPinsDemo { + + private static final Object LOCK = new Object(); + private static boolean signalled = false; + + public static void main(String[] args) throws Exception { + System.out.println("java.version=" + System.getProperty("java.version")); + System.out.println("jdk.tracePinnedThreads=" + System.getProperty("jdk.tracePinnedThreads")); + + // Case 1: synchronized + Object.wait() - releases the monitor while parked, does not pin. + Thread waitThread = Thread.ofVirtual().start(() -> { + synchronized (LOCK) { + try { + while (!signalled) { + LOCK.wait(); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + }); + Thread.sleep(100); + synchronized (LOCK) { + signalled = true; + LOCK.notifyAll(); + } + waitThread.join(); + System.out.println("synchronized + Object.wait(): completed, no pinning trace above this line for it"); + + // Case 2: synchronized + Thread.sleep() - still holds the monitor while parked, DOES pin. + Thread sleepThread = Thread.ofVirtual().start(() -> { + synchronized (LOCK) { + try { + Thread.sleep(150); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + }); + sleepThread.join(); + System.out.println("synchronized + Thread.sleep(): completed, pinning trace (if any) appears above this line for it"); + + System.out.println("done"); + } +} diff --git a/synchronizers/src/main/java/com/ankurm/synchronizers/TimeoutBehaviorDemo.java b/synchronizers/src/main/java/com/ankurm/synchronizers/TimeoutBehaviorDemo.java new file mode 100644 index 0000000..b6712cd --- /dev/null +++ b/synchronizers/src/main/java/com/ankurm/synchronizers/TimeoutBehaviorDemo.java @@ -0,0 +1,89 @@ +package com.ankurm.synchronizers; + +import java.util.concurrent.BrokenBarrierException; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.Phaser; +import java.util.concurrent.Semaphore; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; + +/** + * What each class does when the thing it's waiting for never happens, deliberately forced by + * under-supplying arrivals or permits. The four don't behave identically, which matters more + * than it looks: a timed-out {@code CyclicBarrier} poisons itself for every other waiter, and + * the other three don't. + */ +public final class TimeoutBehaviorDemo { + + public static void main(String[] args) throws Exception { + countDownLatchTimeout(); + cyclicBarrierTimeout(); + phaserTimeout(); + semaphoreTimeout(); + System.out.println("done"); + } + + private static void countDownLatchTimeout() throws InterruptedException { + System.out.println("=== CountDownLatch.await(timeout): expects 2, only 1 ever counts down ==="); + CountDownLatch latch = new CountDownLatch(2); + latch.countDown(); // only one of two - the second never arrives + boolean released = latch.await(200, TimeUnit.MILLISECONDS); + System.out.println("await returned: " + released + " (false = timed out, latch still usable to re-check later)"); + System.out.println("getCount() after timeout: " + latch.getCount()); + } + + private static void cyclicBarrierTimeout() throws InterruptedException { + System.out.println(); + System.out.println("=== CyclicBarrier.await(timeout): expects 2 parties, only 1 arrives ==="); + CyclicBarrier barrier = new CyclicBarrier(2); + try { + barrier.await(200, TimeUnit.MILLISECONDS); + } catch (TimeoutException e) { + System.out.println("await threw TimeoutException, as expected"); + System.out.println("barrier.isBroken() after the timeout: " + barrier.isBroken()); + } catch (BrokenBarrierException e) { + System.out.println("unexpected BrokenBarrierException"); + } + // A second, unrelated thread trying to use the SAME barrier instance afterward: + try { + barrier.await(50, TimeUnit.MILLISECONDS); + } catch (BrokenBarrierException e) { + System.out.println("a second, later await() on the SAME barrier: BrokenBarrierException immediately" + + " - one timeout poisons the barrier for everyone until reset()"); + } catch (TimeoutException e) { + System.out.println("unexpected TimeoutException on the second await"); + } + } + + private static void phaserTimeout() { + System.out.println(); + System.out.println("=== Phaser.awaitAdvanceInterruptibly(phase, timeout): expects 2 parties, only 1 arrives ==="); + Phaser phaser = new Phaser(2); + int phase = phaser.arrive(); // this party arrives; the second one never does + try { + phaser.awaitAdvanceInterruptibly(phase, 200, TimeUnit.MILLISECONDS); + } catch (TimeoutException e) { + System.out.println("awaitAdvanceInterruptibly threw TimeoutException, as expected"); + System.out.println("phaser.getPhase() after the timeout: " + phaser.getPhase() + + " (unchanged - still phase " + phase + ", not advanced or broken)"); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + // A later, correctly-completed round on the SAME phaser instance, no reset needed: + phaser.arriveAndAwaitAdvance(); + System.out.println("a later round on the SAME phaser, once both parties actually arrive: " + + "succeeds normally, now at phase " + phaser.getPhase()); + } + + private static void semaphoreTimeout() throws InterruptedException { + System.out.println(); + System.out.println("=== Semaphore.tryAcquire(timeout): 0 permits available ==="); + Semaphore semaphore = new Semaphore(0); + boolean acquired = semaphore.tryAcquire(200, TimeUnit.MILLISECONDS); + System.out.println("tryAcquire returned: " + acquired + " (false = timed out, no exception, no broken state)"); + semaphore.release(); + boolean acquiredAfterRelease = semaphore.tryAcquire(50, TimeUnit.MILLISECONDS); + System.out.println("a later tryAcquire after a release: " + acquiredAfterRelease + " - fully independent of the earlier timeout"); + } +} diff --git a/synchronizers/src/main/java/com/ankurm/synchronizers/VirtualThreadCompatibilityDemo.java b/synchronizers/src/main/java/com/ankurm/synchronizers/VirtualThreadCompatibilityDemo.java new file mode 100644 index 0000000..7613977 --- /dev/null +++ b/synchronizers/src/main/java/com/ankurm/synchronizers/VirtualThreadCompatibilityDemo.java @@ -0,0 +1,96 @@ +package com.ankurm.synchronizers; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.Phaser; +import java.util.concurrent.Semaphore; + +/** + * All four classes in this module are built on {@link java.util.concurrent.locks.AbstractQueuedSynchronizer} + * (Phaser uses its own compare-and-swap state machine, not AQS directly, but the same + * park/unpark mechanism), not on {@code synchronized}. A virtual thread blocking inside any of + * them parks via {@code LockSupport}, which was always safe to unmount from - none of the four + * ever pinned a carrier, on any JDK version, including JDK 21 well before + * JEP 491 fixed {@code synchronized}. This demo blocks + * a virtual thread inside each one in turn, run with {@code -Djdk.tracePinnedThreads=full}: if + * it printed anything, that would be news. It doesn't. Compare against + * {@link SynchronizedWaitPinsDemo}, the one construct in this module family that genuinely does + * pin on a pre-JDK-24 runtime. + */ +public final class VirtualThreadCompatibilityDemo { + + public static void main(String[] args) throws Exception { + System.out.println("java.version=" + System.getProperty("java.version")); + System.out.println("jdk.tracePinnedThreads=" + System.getProperty("jdk.tracePinnedThreads")); + + runOnVirtualThread("CountDownLatch.await()", () -> { + CountDownLatch latch = new CountDownLatch(1); + Thread.ofPlatform().start(() -> { + sleepQuietly(100); + latch.countDown(); + }); + latch.await(); + }); + + runOnVirtualThread("CyclicBarrier.await()", () -> { + CyclicBarrier barrier = new CyclicBarrier(2); + Thread.ofPlatform().start(() -> { + sleepQuietly(100); + awaitQuietly(barrier); + }); + barrier.await(); + }); + + runOnVirtualThread("Phaser.arriveAndAwaitAdvance()", () -> { + Phaser phaser = new Phaser(2); + Thread.ofPlatform().start(() -> { + sleepQuietly(100); + phaser.arriveAndAwaitAdvance(); + }); + phaser.arriveAndAwaitAdvance(); + }); + + runOnVirtualThread("Semaphore.acquire()", () -> { + Semaphore semaphore = new Semaphore(0); + Thread.ofPlatform().start(() -> { + sleepQuietly(100); + semaphore.release(); + }); + semaphore.acquire(); + }); + + System.out.println("done - no pinning trace printed above for any of the four"); + } + + private interface BlockingAction { + void run() throws Exception; + } + + private static void runOnVirtualThread(String label, BlockingAction action) throws InterruptedException { + Thread t = Thread.ofVirtual().start(() -> { + try { + action.run(); + } catch (Exception e) { + throw new RuntimeException(e); + } + }); + t.join(); + System.out.println(label + ": completed, no trace above this line for it"); + } + + private static void sleepQuietly(long ms) { + try { + Thread.sleep(ms); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + + private static void awaitQuietly(CyclicBarrier barrier) { + try { + barrier.await(); + } catch (Exception e) { + throw new RuntimeException(e); + } + } +} diff --git a/synchronizers/src/test/java/com/ankurm/synchronizers/RoundOrderingCorrectnessTest.java b/synchronizers/src/test/java/com/ankurm/synchronizers/RoundOrderingCorrectnessTest.java new file mode 100644 index 0000000..d6d60fd --- /dev/null +++ b/synchronizers/src/test/java/com/ankurm/synchronizers/RoundOrderingCorrectnessTest.java @@ -0,0 +1,137 @@ +package com.ankurm.synchronizers; + +import org.junit.jupiter.api.RepeatedTest; +import org.junit.jupiter.api.Test; + +import java.util.concurrent.BrokenBarrierException; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.Phaser; +import java.util.concurrent.Semaphore; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Sanity checks, not throughput or pinning proofs - see the module README and the article's + * own quoted transcripts for those. These assert the one property that actually matters for a + * round barrier: no worker's round-N+1 work is ever observed to start before every worker's + * round-N work has finished, repeated to catch a scheduling-dependent race a single run wouldn't. + */ +class RoundOrderingCorrectnessTest { + + private static final int WORKERS = 6; + private static final int ROUNDS = 4; + + @RepeatedTest(20) + void cyclicBarrierNeverStartsARoundEarly() throws Exception { + AtomicInteger currentRound = new AtomicInteger(1); + AtomicInteger[] finishedThisRound = newCounters(); + CyclicBarrier barrier = new CyclicBarrier(WORKERS, () -> currentRound.incrementAndGet()); + + Thread[] workers = new Thread[WORKERS]; + for (int w = 0; w < WORKERS; w++) { + workers[w] = new Thread(() -> { + for (int round = 1; round <= ROUNDS; round++) { + // The assertion: this worker's round MUST equal currentRound at the moment + // it starts working - if some other worker's barrier.await() returned before + // this one arrived, currentRound would already have advanced past `round`. + assertEquals(round, currentRound.get(), + "worker started round " + round + " but currentRound was already " + currentRound.get()); + finishedThisRound[round - 1].incrementAndGet(); + try { + barrier.await(); + } catch (InterruptedException | BrokenBarrierException e) { + throw new RuntimeException(e); + } + } + }); + workers[w].start(); + } + for (Thread t : workers) { + t.join(); + } + for (AtomicInteger count : finishedThisRound) { + assertEquals(WORKERS, count.get(), "every worker must have finished every round exactly once"); + } + } + + @RepeatedTest(20) + void phaserNeverStartsARoundEarly() { + AtomicInteger currentRound = new AtomicInteger(1); + AtomicInteger[] finishedThisRound = newCounters(); + Phaser phaser = new Phaser(WORKERS) { + @Override + protected boolean onAdvance(int phase, int registeredParties) { + currentRound.incrementAndGet(); + return phase + 1 >= ROUNDS || registeredParties == 0; + } + }; + + Thread[] workers = new Thread[WORKERS]; + for (int w = 0; w < WORKERS; w++) { + workers[w] = new Thread(() -> { + for (int round = 1; round <= ROUNDS; round++) { + assertEquals(round, currentRound.get(), + "worker started round " + round + " but currentRound was already " + currentRound.get()); + finishedThisRound[round - 1].incrementAndGet(); + phaser.arriveAndAwaitAdvance(); + } + }); + workers[w].start(); + } + for (Thread t : workers) { + try { + t.join(); + } catch (InterruptedException e) { + throw new RuntimeException(e); + } + } + for (AtomicInteger count : finishedThisRound) { + assertEquals(WORKERS, count.get(), "every worker must have finished every round exactly once"); + } + } + + @Test + void semaphoreNeverExceedsItsPermitCount() throws Exception { + int permits = 3; + int callers = 30; + Semaphore semaphore = new Semaphore(permits); + AtomicInteger concurrent = new AtomicInteger(); + AtomicInteger maxObserved = new AtomicInteger(); + + Thread[] threads = new Thread[callers]; + for (int i = 0; i < callers; i++) { + threads[i] = new Thread(() -> { + try { + semaphore.acquire(); + try { + int inFlight = concurrent.incrementAndGet(); + maxObserved.updateAndGet(max -> Math.max(max, inFlight)); + Thread.sleep(5); + } finally { + concurrent.decrementAndGet(); + semaphore.release(); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + }); + threads[i].start(); + } + for (Thread t : threads) { + t.join(); + } + assertTrue(maxObserved.get() <= permits, + "observed " + maxObserved.get() + " concurrent callers with only " + permits + " permits"); + assertEquals(permits, semaphore.availablePermits(), "all permits must be returned when every caller finishes"); + } + + private static AtomicInteger[] newCounters() { + AtomicInteger[] counters = new AtomicInteger[ROUNDS]; + for (int i = 0; i < ROUNDS; i++) { + counters[i] = new AtomicInteger(); + } + return counters; + } +}