locks: synchronized vs ReentrantLock vs StampedLock companion code
JMH throughput sweep (1-64 threads), 9:1 read-heavy @Group benchmark, tryLock(timeout) deadlock-avoidance demo, and a JDK21-vs-JDK25 JEP 491 virtual-thread-pinning comparison. Also promotes LICENSE to repo root now that a second module exists.
This commit is contained in:
@@ -0,0 +1,11 @@
|
||||
package com.ankurm.locks;
|
||||
|
||||
/**
|
||||
* A shared mutable counter, protected some way against concurrent increment().
|
||||
* Every implementation in this module implements exactly this interface so the
|
||||
* JMH benchmarks can swap the locking strategy without changing anything else.
|
||||
*/
|
||||
public interface Counter {
|
||||
void increment();
|
||||
long get();
|
||||
}
|
||||
@@ -0,0 +1,34 @@
|
||||
package com.ankurm.locks;
|
||||
|
||||
import java.util.concurrent.locks.ReentrantLock;
|
||||
|
||||
/**
|
||||
* The same {@link ReentrantLock}, constructed with {@code fair = true}. Fair
|
||||
* mode grants the lock to the longest-waiting thread, which bounds
|
||||
* starvation but costs throughput - this class exists so the benchmark can
|
||||
* put a number on that cost instead of just asserting it.
|
||||
*/
|
||||
public final class FairReentrantLockCounter implements Counter {
|
||||
private final ReentrantLock lock = new ReentrantLock(true);
|
||||
private long count;
|
||||
|
||||
@Override
|
||||
public void increment() {
|
||||
lock.lock();
|
||||
try {
|
||||
count++;
|
||||
} finally {
|
||||
lock.unlock();
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public long get() {
|
||||
lock.lock();
|
||||
try {
|
||||
return count;
|
||||
} finally {
|
||||
lock.unlock();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,50 @@
|
||||
package com.ankurm.locks;
|
||||
|
||||
import org.openjdk.jmh.annotations.*;
|
||||
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
/**
|
||||
* Write-only workload: every thread just calls {@code increment()} as fast
|
||||
* as it can. This is the benchmark that isolates pure lock-acquisition
|
||||
* overhead - there is no read path here to give {@link StampedLockCounter}
|
||||
* its usual advantage, so the honest expectation is that all four come out
|
||||
* close, with {@code synchronized} and unfair {@link java.util.concurrent.locks.ReentrantLock}
|
||||
* at the front and the fair lock and the write-locked {@link java.util.concurrent.locks.StampedLock}
|
||||
* paying a small, measurable tax. Run at a fixed thread count per JVM
|
||||
* invocation via {@code -t N}; {@code scripts/run-all.sh} sweeps 1, 2, 4, 8,
|
||||
* 16, 32 and 64 and captures each into {@code output/}.
|
||||
*/
|
||||
@BenchmarkMode(Mode.Throughput)
|
||||
@OutputTimeUnit(TimeUnit.MILLISECONDS)
|
||||
@State(Scope.Benchmark)
|
||||
@Warmup(iterations = 3, time = 1, timeUnit = TimeUnit.SECONDS)
|
||||
@Measurement(iterations = 5, time = 1, timeUnit = TimeUnit.SECONDS)
|
||||
@Fork(1)
|
||||
public class IncrementBenchmark {
|
||||
|
||||
private final SynchronizedCounter synchronizedCounter = new SynchronizedCounter();
|
||||
private final ReentrantLockCounter reentrantLockCounter = new ReentrantLockCounter();
|
||||
private final FairReentrantLockCounter fairReentrantLockCounter = new FairReentrantLockCounter();
|
||||
private final StampedLockCounter stampedLockCounter = new StampedLockCounter();
|
||||
|
||||
@Benchmark
|
||||
public void synchronized_() {
|
||||
synchronizedCounter.increment();
|
||||
}
|
||||
|
||||
@Benchmark
|
||||
public void reentrantLockUnfair() {
|
||||
reentrantLockCounter.increment();
|
||||
}
|
||||
|
||||
@Benchmark
|
||||
public void reentrantLockFair() {
|
||||
fairReentrantLockCounter.increment();
|
||||
}
|
||||
|
||||
@Benchmark
|
||||
public void stampedLockWrite() {
|
||||
stampedLockCounter.increment();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,68 @@
|
||||
package com.ankurm.locks;
|
||||
|
||||
import org.openjdk.jmh.annotations.*;
|
||||
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
/**
|
||||
* A 9:1 read:write workload using JMH's {@code @Group} feature, which runs
|
||||
* two benchmark methods concurrently at a fixed thread ratio and reports
|
||||
* each side's own throughput. This is the workload {@link StampedLockCounter}
|
||||
* is actually for: nine threads spin on {@code get()} while one thread
|
||||
* spins on {@code increment()}. Fixed at 10 total threads per lock type
|
||||
* (this sandbox has 2 vCPUs, so this is already a 5x-oversubscribed,
|
||||
* contention-heavy point, not a scalability sweep - see {@link IncrementBenchmark}
|
||||
* for the thread-count sweep on the write-only path).
|
||||
*/
|
||||
@BenchmarkMode(Mode.Throughput)
|
||||
@OutputTimeUnit(TimeUnit.MILLISECONDS)
|
||||
@Warmup(iterations = 3, time = 1, timeUnit = TimeUnit.SECONDS)
|
||||
@Measurement(iterations = 5, time = 1, timeUnit = TimeUnit.SECONDS)
|
||||
@Fork(1)
|
||||
public class ReadHeavyBenchmark {
|
||||
|
||||
@State(Scope.Group)
|
||||
public static class SynchronizedState {
|
||||
final SynchronizedCounter counter = new SynchronizedCounter();
|
||||
}
|
||||
|
||||
@State(Scope.Group)
|
||||
public static class ReentrantLockState {
|
||||
final ReentrantLockCounter counter = new ReentrantLockCounter();
|
||||
}
|
||||
|
||||
@State(Scope.Group)
|
||||
public static class StampedLockState {
|
||||
final StampedLockCounter counter = new StampedLockCounter();
|
||||
}
|
||||
|
||||
@Benchmark @Group("synchronizedCounter") @GroupThreads(9)
|
||||
public long synchronizedRead(SynchronizedState s) {
|
||||
return s.counter.get();
|
||||
}
|
||||
|
||||
@Benchmark @Group("synchronizedCounter") @GroupThreads(1)
|
||||
public void synchronizedWrite(SynchronizedState s) {
|
||||
s.counter.increment();
|
||||
}
|
||||
|
||||
@Benchmark @Group("reentrantLock") @GroupThreads(9)
|
||||
public long reentrantLockRead(ReentrantLockState s) {
|
||||
return s.counter.get();
|
||||
}
|
||||
|
||||
@Benchmark @Group("reentrantLock") @GroupThreads(1)
|
||||
public void reentrantLockWrite(ReentrantLockState s) {
|
||||
s.counter.increment();
|
||||
}
|
||||
|
||||
@Benchmark @Group("stampedLock") @GroupThreads(9)
|
||||
public long stampedLockRead(StampedLockState s) {
|
||||
return s.counter.get();
|
||||
}
|
||||
|
||||
@Benchmark @Group("stampedLock") @GroupThreads(1)
|
||||
public void stampedLockWrite(StampedLockState s) {
|
||||
s.counter.increment();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,35 @@
|
||||
package com.ankurm.locks;
|
||||
|
||||
import java.util.concurrent.locks.ReentrantLock;
|
||||
|
||||
/**
|
||||
* {@link ReentrantLock} in its default, unfair mode. Unfair means a thread
|
||||
* that is already running can barge in front of threads that have been
|
||||
* parked waiting longer - which is exactly why it usually out-throughputs
|
||||
* the fair variant: no bookkeeping to enforce arrival order, no forced
|
||||
* context switch to wake the "correct" next thread.
|
||||
*/
|
||||
public final class ReentrantLockCounter implements Counter {
|
||||
private final ReentrantLock lock = new ReentrantLock();
|
||||
private long count;
|
||||
|
||||
@Override
|
||||
public void increment() {
|
||||
lock.lock();
|
||||
try {
|
||||
count++;
|
||||
} finally {
|
||||
lock.unlock();
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public long get() {
|
||||
lock.lock();
|
||||
try {
|
||||
return count;
|
||||
} finally {
|
||||
lock.unlock();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,44 @@
|
||||
package com.ankurm.locks;
|
||||
|
||||
import java.util.concurrent.locks.StampedLock;
|
||||
|
||||
/**
|
||||
* {@link StampedLock} used the way it is meant to be used: writers take the
|
||||
* exclusive write lock, but readers first try an <em>optimistic</em> read -
|
||||
* no lock is acquired at all, the read just checks afterwards whether a
|
||||
* writer slipped in while it was running, via {@link StampedLock#validate}.
|
||||
* If a writer did, the reader falls back to a real (blocking) read lock.
|
||||
* {@code get()} here is that full three-step optimistic-read protocol, not
|
||||
* a simplified version of it.
|
||||
*/
|
||||
public final class StampedLockCounter implements Counter {
|
||||
private final StampedLock lock = new StampedLock();
|
||||
private long count;
|
||||
|
||||
@Override
|
||||
public void increment() {
|
||||
long stamp = lock.writeLock();
|
||||
try {
|
||||
count++;
|
||||
} finally {
|
||||
lock.unlockWrite(stamp);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public long get() {
|
||||
long stamp = lock.tryOptimisticRead();
|
||||
long value = count;
|
||||
if (!lock.validate(stamp)) {
|
||||
// A writer ran between the read above and the validate() call.
|
||||
// Fall back to a real, blocking read lock and read again.
|
||||
stamp = lock.readLock();
|
||||
try {
|
||||
value = count;
|
||||
} finally {
|
||||
lock.unlockRead(stamp);
|
||||
}
|
||||
}
|
||||
return value;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
package com.ankurm.locks;
|
||||
|
||||
/**
|
||||
* The baseline: a plain {@code synchronized} method. One monitor, mutual
|
||||
* exclusion for both the read and the write, no fairness knob, no timeout.
|
||||
*/
|
||||
public final class SynchronizedCounter implements Counter {
|
||||
private long count;
|
||||
|
||||
@Override
|
||||
public synchronized void increment() {
|
||||
count++;
|
||||
}
|
||||
|
||||
@Override
|
||||
public synchronized long get() {
|
||||
return count;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,65 @@
|
||||
package com.ankurm.locks;
|
||||
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.locks.ReentrantLock;
|
||||
|
||||
/**
|
||||
* A real deadlock, avoided in real time by {@link ReentrantLock#tryLock(long, TimeUnit)}.
|
||||
* Two threads acquire two locks in opposite order - the textbook deadlock
|
||||
* setup. {@code lock()} would hang both threads forever. {@code tryLock}
|
||||
* with a timeout gives each thread a way out: back off, release what you
|
||||
* hold, and retry. Run this and it always finishes; comment out the
|
||||
* timeout path and call {@code lock()} instead, and it never does.
|
||||
*/
|
||||
public final class TryLockTimeoutDemo {
|
||||
|
||||
private static final ReentrantLock LOCK_A = new ReentrantLock();
|
||||
private static final ReentrantLock LOCK_B = new ReentrantLock();
|
||||
|
||||
public static void main(String[] args) throws InterruptedException {
|
||||
Thread t1 = new Thread(() -> worker("Thread-1", LOCK_A, LOCK_B), "Thread-1");
|
||||
Thread t2 = new Thread(() -> worker("Thread-2", LOCK_B, LOCK_A), "Thread-2");
|
||||
long start = System.nanoTime();
|
||||
t1.start();
|
||||
t2.start();
|
||||
t1.join();
|
||||
t2.join();
|
||||
long elapsedMs = (System.nanoTime() - start) / 1_000_000;
|
||||
System.out.println("Both threads finished in " + elapsedMs + " ms - no deadlock.");
|
||||
}
|
||||
|
||||
private static void worker(String name, ReentrantLock first, ReentrantLock second) {
|
||||
int attempts = 0;
|
||||
while (true) {
|
||||
attempts++;
|
||||
try {
|
||||
if (first.tryLock(200, TimeUnit.MILLISECONDS)) {
|
||||
try {
|
||||
// Force the interleaving that would deadlock under plain lock():
|
||||
// give the other thread time to grab its own first lock before
|
||||
// this thread tries for the second one.
|
||||
Thread.sleep(50);
|
||||
if (second.tryLock(200, TimeUnit.MILLISECONDS)) {
|
||||
try {
|
||||
System.out.println(name + ": acquired both locks on attempt " + attempts + ".");
|
||||
return;
|
||||
} finally {
|
||||
second.unlock();
|
||||
}
|
||||
} else {
|
||||
System.out.println(name + ": timed out waiting for second lock on attempt "
|
||||
+ attempts + " - backing off and retrying.");
|
||||
}
|
||||
} finally {
|
||||
first.unlock();
|
||||
}
|
||||
} else {
|
||||
System.out.println(name + ": timed out waiting for first lock on attempt " + attempts + ".");
|
||||
}
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,34 @@
|
||||
package com.ankurm.locks;
|
||||
|
||||
/**
|
||||
* The smallest program that shows JEP 491 (Synchronize Virtual Threads
|
||||
* without Pinning, delivered JDK 24) doing its job. A virtual thread enters
|
||||
* a {@code synchronized} block and then blocks (a plain {@code Thread.sleep}).
|
||||
* Run with {@code -Djdk.tracePinnedThreads=full}:
|
||||
* <ul>
|
||||
* <li>On JDK 21 (pre-JEP-491) this prints a pinned-thread trace pointing
|
||||
* straight at the {@code synchronized} block below - the virtual
|
||||
* thread cannot unmount because it is holding a monitor.</li>
|
||||
* <li>On JDK 25 (post-JEP-491) it prints nothing: the virtual thread
|
||||
* unmounts from its carrier for the sleep and remounts afterwards,
|
||||
* monitor and all.</li>
|
||||
* </ul>
|
||||
* Both runs are captured verbatim in {@code output/}; nothing here is
|
||||
* asserted, only observed.
|
||||
*/
|
||||
public final class VirtualThreadPinningDemo {
|
||||
|
||||
public static void main(String[] args) throws InterruptedException {
|
||||
Thread vt = Thread.ofVirtual().start(() -> {
|
||||
synchronized (VirtualThreadPinningDemo.class) {
|
||||
try {
|
||||
Thread.sleep(200);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
}
|
||||
});
|
||||
vt.join();
|
||||
System.out.println("Virtual thread finished. (No output above this line means it did not pin.)");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,73 @@
|
||||
package com.ankurm.locks;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.Timeout;
|
||||
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.stream.IntStream;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* Ordinary correctness checks: N threads each increment M times, the final
|
||||
* count must be exactly N*M. This does NOT test throughput, fairness, or
|
||||
* memory-visibility ordering - it only proves each locking strategy
|
||||
* actually serializes increments (no lost updates). The benchmarks in
|
||||
* this module measure the performance claims; this test just guards
|
||||
* against a broken implementation slipping in.
|
||||
*/
|
||||
class CounterCorrectnessTest {
|
||||
|
||||
private static final int THREADS = 8;
|
||||
private static final int INCREMENTS_PER_THREAD = 50_000;
|
||||
|
||||
@Test
|
||||
@Timeout(30)
|
||||
void synchronizedCounterHasNoLostUpdates() throws InterruptedException {
|
||||
assertNoLostUpdates(new SynchronizedCounter());
|
||||
}
|
||||
|
||||
@Test
|
||||
@Timeout(30)
|
||||
void reentrantLockCounterHasNoLostUpdates() throws InterruptedException {
|
||||
assertNoLostUpdates(new ReentrantLockCounter());
|
||||
}
|
||||
|
||||
@Test
|
||||
@Timeout(30)
|
||||
void fairReentrantLockCounterHasNoLostUpdates() throws InterruptedException {
|
||||
assertNoLostUpdates(new FairReentrantLockCounter());
|
||||
}
|
||||
|
||||
@Test
|
||||
@Timeout(30)
|
||||
void stampedLockCounterHasNoLostUpdates() throws InterruptedException {
|
||||
assertNoLostUpdates(new StampedLockCounter());
|
||||
}
|
||||
|
||||
private void assertNoLostUpdates(Counter counter) throws InterruptedException {
|
||||
CountDownLatch ready = new CountDownLatch(THREADS);
|
||||
CountDownLatch start = new CountDownLatch(1);
|
||||
CountDownLatch done = new CountDownLatch(THREADS);
|
||||
|
||||
IntStream.range(0, THREADS).forEach(i -> new Thread(() -> {
|
||||
ready.countDown();
|
||||
try {
|
||||
start.await();
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
return;
|
||||
}
|
||||
for (int j = 0; j < INCREMENTS_PER_THREAD; j++) {
|
||||
counter.increment();
|
||||
}
|
||||
done.countDown();
|
||||
}).start());
|
||||
|
||||
ready.await();
|
||||
start.countDown();
|
||||
done.await();
|
||||
|
||||
assertEquals((long) THREADS * INCREMENTS_PER_THREAD, counter.get());
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user