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 <[email protected]>
Claude-Session: https://claude.ai/code/session_01FhzLY5p6okFva3qsnsRyvM
This commit is contained in:
2026-09-30 06:43:52 +00:00
co-authored by Claude Sonnet 5
parent 0819993c51
commit 55ee850a75
20 changed files with 962 additions and 0 deletions
@@ -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:
*
* <ol>
* <li>A single "starting gate" latch releases all workers at once instead of letting them
* race to start as soon as each is spawned.</li>
* <li>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.</li>
* </ol>
*/
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<String> 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<String> 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)
};
}
@@ -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<String> 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<String> 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();
}
}
}
@@ -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<String> 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<String> 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();
}
}
}
@@ -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.
*
* <p>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<String> 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<String> 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();
}
}
}
@@ -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 <em>while still holding</em> 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");
}
}
@@ -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");
}
}
@@ -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
* <a href="https://openjdk.org/jeps/491">JEP 491</a> 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);
}
}
}
@@ -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;
}
}