cow: CopyOnWriteArrayList - deterministic iterator-snapshot and synchronizedList CME demos, JMH read/write throughput, and the addLast/removeFirst fix for a real concurrent-writer race
This commit is contained in:
@@ -0,0 +1,70 @@
|
||||
package com.ankurm.cow;
|
||||
|
||||
import java.util.Iterator;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
|
||||
/**
|
||||
* The claim every CopyOnWriteArrayList article makes - "an iterator sees a snapshot of the list
|
||||
* as it was when the iterator was created" - deserves a demo that doesn't depend on {@code
|
||||
* Thread.sleep} timing guesses to actually prove it. This one forces the exact ordering with
|
||||
* {@link CountDownLatch} barriers: an iterator is opened, a write happens from another thread
|
||||
* while the first iterator is still mid-walk, and a second iterator is opened only after that
|
||||
* write completes. The first iterator must still only see the pre-write elements; the second
|
||||
* must see all of them; and the test suite ({@code IteratorSnapshotTest}) pins every number here
|
||||
* as an assertion.
|
||||
*/
|
||||
public final class IteratorSnapshotDemo {
|
||||
|
||||
private IteratorSnapshotDemo() {}
|
||||
|
||||
public static void main(String[] args) throws InterruptedException {
|
||||
CopyOnWriteArrayList<String> list = new CopyOnWriteArrayList<>();
|
||||
list.add("Apple");
|
||||
list.add("Banana");
|
||||
list.add("Cherry");
|
||||
|
||||
System.out.println("list before any iterator is taken = " + list);
|
||||
|
||||
Iterator<String> beforeWrite = list.iterator(); // snapshot taken HERE, right now
|
||||
System.out.println("beforeWrite.next() = " + beforeWrite.next()); // consumes "Apple"
|
||||
|
||||
CountDownLatch readerReady = new CountDownLatch(1);
|
||||
CountDownLatch writeDone = new CountDownLatch(1);
|
||||
|
||||
Thread writer = new Thread(() -> {
|
||||
try {
|
||||
readerReady.await();
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
return;
|
||||
}
|
||||
list.add("Date");
|
||||
list.add("Elderberry");
|
||||
System.out.println("[writer] added Date and Elderberry - list is now " + list);
|
||||
writeDone.countDown();
|
||||
}, "writer");
|
||||
writer.start();
|
||||
|
||||
readerReady.countDown();
|
||||
writeDone.await();
|
||||
writer.join();
|
||||
|
||||
System.out.println();
|
||||
System.out.println("=== beforeWrite iterator was created BEFORE the write - draining the rest of it ===");
|
||||
while (beforeWrite.hasNext()) {
|
||||
System.out.println("beforeWrite.next() = " + beforeWrite.next());
|
||||
}
|
||||
System.out.println("beforeWrite never saw Date or Elderberry, and threw no exception.");
|
||||
|
||||
System.out.println();
|
||||
System.out.println("=== afterWrite iterator is created AFTER the write - it sees the current array ===");
|
||||
Iterator<String> afterWrite = list.iterator();
|
||||
while (afterWrite.hasNext()) {
|
||||
System.out.println("afterWrite.next() = " + afterWrite.next());
|
||||
}
|
||||
|
||||
System.out.println();
|
||||
System.out.println("Final list = " + list);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,101 @@
|
||||
package com.ankurm.cow;
|
||||
|
||||
import org.openjdk.jmh.annotations.Benchmark;
|
||||
import org.openjdk.jmh.annotations.BenchmarkMode;
|
||||
import org.openjdk.jmh.annotations.Fork;
|
||||
import org.openjdk.jmh.annotations.Level;
|
||||
import org.openjdk.jmh.annotations.Measurement;
|
||||
import org.openjdk.jmh.annotations.Mode;
|
||||
import org.openjdk.jmh.annotations.OutputTimeUnit;
|
||||
import org.openjdk.jmh.annotations.Param;
|
||||
import org.openjdk.jmh.annotations.Scope;
|
||||
import org.openjdk.jmh.annotations.Setup;
|
||||
import org.openjdk.jmh.annotations.State;
|
||||
import org.openjdk.jmh.annotations.TearDown;
|
||||
import org.openjdk.jmh.annotations.Threads;
|
||||
import org.openjdk.jmh.annotations.Warmup;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
/**
|
||||
* The read side of the "read/write ratio" question: does reading through {@code
|
||||
* CopyOnWriteArrayList} or {@code Collections.synchronizedList} cost more when four threads are
|
||||
* reading at once, and does a concurrent background writer change that answer? Four reader
|
||||
* threads read the list by index and sum it; {@code writer=BACKGROUND} runs a fifth thread
|
||||
* continuously mutating the same list with no delay between writes, for the whole measurement
|
||||
* window. Reading is done with {@code get(index)}, not an {@code Iterator}, deliberately - a
|
||||
* {@code synchronizedList} iterator races with a concurrent structural change and throws
|
||||
* {@code ConcurrentModificationException} (reproduced deterministically in {@link
|
||||
* SynchronizedListCmeDemo}), which would make the SYNC/BACKGROUND combination fail to complete
|
||||
* rather than measure anything. {@code get(index)} never has that problem on either
|
||||
* implementation, so it isolates the cost this benchmark is actually about.
|
||||
*/
|
||||
@BenchmarkMode(Mode.Throughput)
|
||||
@OutputTimeUnit(TimeUnit.MILLISECONDS)
|
||||
@Warmup(iterations = 2, time = 1)
|
||||
@Measurement(iterations = 3, time = 1)
|
||||
@Fork(1)
|
||||
@Threads(4)
|
||||
public class ReadThroughputBenchmark {
|
||||
|
||||
@State(Scope.Benchmark)
|
||||
public static class ListState {
|
||||
@Param({"COW", "SYNC"})
|
||||
String container;
|
||||
|
||||
@Param({"NONE", "BACKGROUND"})
|
||||
String writer;
|
||||
|
||||
List<Integer> list;
|
||||
Thread writerThread;
|
||||
AtomicBoolean stop;
|
||||
|
||||
@Setup(Level.Trial)
|
||||
public void setup() {
|
||||
List<Integer> seed = new ArrayList<>();
|
||||
for (int i = 0; i < 1000; i++) seed.add(i);
|
||||
|
||||
list = "COW".equals(container)
|
||||
? new CopyOnWriteArrayList<>(seed)
|
||||
: Collections.synchronizedList(new ArrayList<>(seed));
|
||||
|
||||
stop = new AtomicBoolean(false);
|
||||
if ("BACKGROUND".equals(writer)) {
|
||||
writerThread = new Thread(() -> {
|
||||
int i = 0;
|
||||
while (!stop.get()) {
|
||||
// add-then-remove, in THIS order, so size only ever grows then shrinks
|
||||
// back - it never dips below 1000, so concurrent get(i) for i < 1000
|
||||
// on readers below never races an out-of-bounds index.
|
||||
list.add(0, i);
|
||||
list.remove(list.size() - 1);
|
||||
i++;
|
||||
}
|
||||
}, "bg-writer");
|
||||
writerThread.setDaemon(true);
|
||||
writerThread.start();
|
||||
}
|
||||
}
|
||||
|
||||
@TearDown(Level.Trial)
|
||||
public void tearDown() throws InterruptedException {
|
||||
stop.set(true);
|
||||
if (writerThread != null) writerThread.join(2000);
|
||||
}
|
||||
}
|
||||
|
||||
@Benchmark
|
||||
public long readSum(ListState s) {
|
||||
long sum = 0;
|
||||
int n = 1000; // the writer thread never lets size drop below this
|
||||
for (int i = 0; i < n; i++) {
|
||||
sum += s.list.get(i);
|
||||
}
|
||||
return sum;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,89 @@
|
||||
package com.ankurm.cow;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
import java.util.Iterator;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
/**
|
||||
* {@code Collections.synchronizedList} synchronizes each individual method call, but an
|
||||
* iteration is a sequence of separate {@code hasNext()}/{@code next()} calls - nothing stops
|
||||
* another thread's {@code add()} from landing in between two of them unless the caller wraps the
|
||||
* whole iteration in its own {@code synchronized (list) { ... }} block, exactly as the class's
|
||||
* own Javadoc warns. This demo forces that exact race with latches instead of hoping for it, so
|
||||
* the {@code ConcurrentModificationException} shows up every run rather than most runs.
|
||||
*/
|
||||
public final class SynchronizedListCmeDemo {
|
||||
|
||||
private SynchronizedListCmeDemo() {}
|
||||
|
||||
public static void main(String[] args) throws InterruptedException {
|
||||
List<String> syncList = Collections.synchronizedList(new ArrayList<>(List.of("Apple", "Banana", "Cherry")));
|
||||
|
||||
CountDownLatch readerTookFirstElement = new CountDownLatch(1);
|
||||
CountDownLatch writerAdded = new CountDownLatch(1);
|
||||
AtomicReference<Exception> caught = new AtomicReference<>();
|
||||
|
||||
Thread writer = new Thread(() -> {
|
||||
try {
|
||||
readerTookFirstElement.await();
|
||||
syncList.add("Date"); // structural modification while the reader is mid-iteration
|
||||
System.out.println("[writer] syncList.add(\"Date\") completed - list is now " + syncList);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
} finally {
|
||||
writerAdded.countDown();
|
||||
}
|
||||
}, "writer");
|
||||
|
||||
Thread reader = new Thread(() -> {
|
||||
// Deliberately NOT wrapped in synchronized(syncList) { ... } - this is the mistake the demo reproduces.
|
||||
Iterator<String> it = syncList.iterator();
|
||||
System.out.println("[reader] it.next() = " + it.next()); // "Apple" - fine, no structural change yet
|
||||
readerTookFirstElement.countDown();
|
||||
try {
|
||||
writerAdded.await();
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
return;
|
||||
}
|
||||
try {
|
||||
it.next(); // the writer's add() already happened - this call must throw
|
||||
System.out.println("[reader] it.next() returned without throwing (unexpected)");
|
||||
} catch (java.util.ConcurrentModificationException e) {
|
||||
caught.set(e);
|
||||
System.out.println("[reader] it.next() threw ConcurrentModificationException, as expected");
|
||||
}
|
||||
}, "reader");
|
||||
|
||||
reader.start();
|
||||
writer.start();
|
||||
reader.join();
|
||||
writer.join();
|
||||
|
||||
System.out.println();
|
||||
System.out.println("Exception actually thrown: "
|
||||
+ (caught.get() != null ? caught.get().getClass().getName() : "none"));
|
||||
|
||||
System.out.println();
|
||||
System.out.println("=== The fix: synchronize the WHOLE iteration yourself, as the Javadoc instructs ===");
|
||||
List<String> syncList2 = Collections.synchronizedList(new ArrayList<>(List.of("Apple", "Banana", "Cherry")));
|
||||
Thread safeWriter = new Thread(() -> syncList2.add("Date"), "safe-writer");
|
||||
synchronized (syncList2) {
|
||||
Iterator<String> it = syncList2.iterator();
|
||||
System.out.println("holding the lock for the whole iteration - starting safeWriter now");
|
||||
safeWriter.start();
|
||||
int count = 0;
|
||||
while (it.hasNext()) {
|
||||
it.next();
|
||||
count++;
|
||||
}
|
||||
System.out.println("iterated " + count + " elements with no exception, because safeWriter");
|
||||
System.out.println("cannot acquire the same monitor until this synchronized block exits");
|
||||
}
|
||||
safeWriter.join();
|
||||
System.out.println("after the block: syncList2 = " + syncList2);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,75 @@
|
||||
package com.ankurm.cow;
|
||||
|
||||
import org.openjdk.jmh.annotations.Benchmark;
|
||||
import org.openjdk.jmh.annotations.BenchmarkMode;
|
||||
import org.openjdk.jmh.annotations.Fork;
|
||||
import org.openjdk.jmh.annotations.Level;
|
||||
import org.openjdk.jmh.annotations.Measurement;
|
||||
import org.openjdk.jmh.annotations.Mode;
|
||||
import org.openjdk.jmh.annotations.OutputTimeUnit;
|
||||
import org.openjdk.jmh.annotations.Param;
|
||||
import org.openjdk.jmh.annotations.Scope;
|
||||
import org.openjdk.jmh.annotations.Setup;
|
||||
import org.openjdk.jmh.annotations.State;
|
||||
import org.openjdk.jmh.annotations.Threads;
|
||||
import org.openjdk.jmh.annotations.Warmup;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
/**
|
||||
* The write side: the cost of a single add-then-remove pair (chosen so the list's size, and
|
||||
* therefore the cost of every subsequent write, stays constant across the whole measurement),
|
||||
* at two sizes and two thread counts. {@code CopyOnWriteArrayList}'s write cost is expected to
|
||||
* scale with list size, since every write copies the entire backing array; {@code
|
||||
* synchronizedList}'s write cost should not.
|
||||
*/
|
||||
@BenchmarkMode(Mode.Throughput)
|
||||
@OutputTimeUnit(TimeUnit.MILLISECONDS)
|
||||
@Warmup(iterations = 2, time = 1)
|
||||
@Measurement(iterations = 3, time = 1)
|
||||
@Fork(1)
|
||||
public class WriteThroughputBenchmark {
|
||||
|
||||
@State(Scope.Benchmark)
|
||||
public static class ListState {
|
||||
@Param({"COW", "SYNC"})
|
||||
String container;
|
||||
|
||||
@Param({"10", "1000"})
|
||||
int size;
|
||||
|
||||
List<Integer> list;
|
||||
|
||||
@Setup(Level.Trial)
|
||||
public void setup() {
|
||||
List<Integer> seed = new ArrayList<>();
|
||||
for (int i = 0; i < size; i++) seed.add(i);
|
||||
list = "COW".equals(container)
|
||||
? new CopyOnWriteArrayList<>(seed)
|
||||
: Collections.synchronizedList(new ArrayList<>(seed));
|
||||
}
|
||||
}
|
||||
|
||||
@Benchmark
|
||||
@Threads(1)
|
||||
public void addRemove_oneWriter(ListState s) {
|
||||
s.list.addLast(-1);
|
||||
s.list.removeFirst();
|
||||
}
|
||||
|
||||
@Benchmark
|
||||
@Threads(4)
|
||||
public void addRemove_fourConcurrentWriters(ListState s) {
|
||||
// addLast()/removeFirst() (JEP 431 default methods) instead of add(-1); remove(size()-1):
|
||||
// the latter reads size() and removes by that index as two separate calls, which races
|
||||
// under concurrent writers - a real IndexOutOfBoundsException from exactly that race is
|
||||
// in this module's output/06-write-throughput-race.txt. addLast/removeFirst are each a
|
||||
// single call with no externally-fetched index, so they stay correct under contention.
|
||||
s.list.addLast(-1);
|
||||
s.list.removeFirst();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,137 @@
|
||||
package com.ankurm.cow;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
import java.util.Iterator;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
|
||||
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 CowClaimsTest {
|
||||
|
||||
@Test
|
||||
void iteratorTakenBeforeAWriteNeverSeesThatWrite() {
|
||||
CopyOnWriteArrayList<String> list = new CopyOnWriteArrayList<>(List.of("a", "b"));
|
||||
Iterator<String> it = list.iterator();
|
||||
list.add("c"); // write happens after the iterator snapshot was taken
|
||||
List<String> seen = new ArrayList<>();
|
||||
while (it.hasNext()) seen.add(it.next());
|
||||
assertEquals(List.of("a", "b"), seen, "the iterator must not see a write that happened after it was created");
|
||||
assertEquals(List.of("a", "b", "c"), list, "but the list itself must reflect the write");
|
||||
}
|
||||
|
||||
@Test
|
||||
void iteratorTakenAfterAWriteSeesIt() {
|
||||
CopyOnWriteArrayList<String> list = new CopyOnWriteArrayList<>(List.of("a", "b"));
|
||||
list.add("c");
|
||||
Iterator<String> it = list.iterator();
|
||||
List<String> seen = new ArrayList<>();
|
||||
while (it.hasNext()) seen.add(it.next());
|
||||
assertEquals(List.of("a", "b", "c"), seen);
|
||||
}
|
||||
|
||||
@Test
|
||||
void copyOnWriteIteratorNeverThrowsConcurrentModificationEvenUnderConcurrentWrites() throws InterruptedException {
|
||||
CopyOnWriteArrayList<Integer> list = new CopyOnWriteArrayList<>(List.of(1, 2, 3));
|
||||
CountDownLatch go = new CountDownLatch(1);
|
||||
Thread writer = new Thread(() -> {
|
||||
try {
|
||||
go.await();
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
for (int i = 0; i < 1000; i++) list.add(i);
|
||||
});
|
||||
writer.start();
|
||||
go.countDown();
|
||||
assertFalse(throwsCme(() -> {
|
||||
Iterator<Integer> it = list.iterator();
|
||||
while (it.hasNext()) it.next();
|
||||
}), "CopyOnWriteArrayList's iterator must never throw ConcurrentModificationException");
|
||||
writer.join();
|
||||
}
|
||||
|
||||
@Test
|
||||
void synchronizedListIteratorThrowsCmeOnConcurrentStructuralChange() throws InterruptedException {
|
||||
List<String> list = Collections.synchronizedList(new ArrayList<>(List.of("a", "b", "c")));
|
||||
CountDownLatch readerTookOne = new CountDownLatch(1);
|
||||
CountDownLatch writerDone = new CountDownLatch(1);
|
||||
|
||||
Thread writer = new Thread(() -> {
|
||||
try {
|
||||
readerTookOne.await();
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
list.add("d");
|
||||
writerDone.countDown();
|
||||
});
|
||||
writer.start();
|
||||
|
||||
Iterator<String> it = list.iterator();
|
||||
it.next(); // consumes "a"
|
||||
readerTookOne.countDown();
|
||||
writerDone.await();
|
||||
|
||||
assertThrows(java.util.ConcurrentModificationException.class, it::next,
|
||||
"iterating a synchronizedList without external synchronization must throw when another thread structurally modifies it mid-iteration");
|
||||
writer.join();
|
||||
}
|
||||
|
||||
@Test
|
||||
void synchronizingTheWholeIterationOnTheListPreventsTheSameCme() throws InterruptedException {
|
||||
List<String> list = Collections.synchronizedList(new ArrayList<>(List.of("a", "b", "c")));
|
||||
Thread writer = new Thread(() -> list.add("d"));
|
||||
int count;
|
||||
synchronized (list) {
|
||||
writer.start();
|
||||
Iterator<String> it = list.iterator();
|
||||
int c = 0;
|
||||
while (it.hasNext()) {
|
||||
it.next();
|
||||
c++;
|
||||
}
|
||||
count = c;
|
||||
}
|
||||
writer.join();
|
||||
assertEquals(3, count, "the writer could not run until the synchronized block released the lock");
|
||||
assertEquals(4, list.size());
|
||||
}
|
||||
|
||||
@Test
|
||||
void copyOnWriteArrayListAllowsDuplicatesAndPreservesInsertionOrder() {
|
||||
CopyOnWriteArrayList<String> list = new CopyOnWriteArrayList<>();
|
||||
list.add("x");
|
||||
list.add("y");
|
||||
list.add("x");
|
||||
assertEquals(List.of("x", "y", "x"), list);
|
||||
}
|
||||
|
||||
@Test
|
||||
void copyOnWriteArrayListImplementsSequencedCollectionLikeAnyOtherList() {
|
||||
CopyOnWriteArrayList<Integer> list = new CopyOnWriteArrayList<>(List.of(1, 2, 3));
|
||||
assertTrue(list instanceof java.util.SequencedCollection);
|
||||
assertEquals(3, list.getLast());
|
||||
assertEquals(List.of(3, 2, 1), list.reversed());
|
||||
}
|
||||
|
||||
private interface ThrowingRunnable {
|
||||
void run();
|
||||
}
|
||||
|
||||
private static boolean throwsCme(ThrowingRunnable r) {
|
||||
try {
|
||||
r.run();
|
||||
return false;
|
||||
} catch (java.util.ConcurrentModificationException e) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user