Language-Level Concurrency and Multithreading Questions
Per-language and per-runtime concurrency: the threading and async APIs each language and platform provides (goroutines, channels, worker pools and pipelines in Go, threads, executors, ThreadPoolExecutor tuning and CompletableFuture in Java, Kotlin coroutines, dispatchers, Flow and structured concurrency, Swift GCD with queues and QoS, DispatchGroup, OperationQueue, actors and async/await, Android Handler/Looper, HandlerThread and thread pools, Flutter isolates, C++ std::thread and atomics, Python threads versus multiprocessing versus asyncio tasks, gather and graceful shutdown), each language's memory model and visibility guarantees (Java happens-before and volatile, C++11 memory orders such as relaxed, acquire and release, Objective-C and Swift atomic versus nonatomic), mobile main-thread rules and JNI thread attachment, cancellation and shutdown idioms, and the idioms for coordinating shared state safely in that language, including thread-safe caches, singletons and bounded queues. Also covers reproducing, testing and diagnosing races and deadlocks in a specific language or app, and migrating callback, GCD or thread-pool code to structured concurrency. Covers choosing and using a language's concurrency primitives correctly. Boundary: general synchronization theory, deadlock and lock-free algorithm internals, OS scheduling, database isolation levels, and callback or event-loop architecture are covered elsewhere.
Write a concise C++ example in which two threads using relaxed atomic stores and loads reach a state that sequential consistency forbids, as can happen on weakly ordered hardware such as ARM. Then show the minimal change that makes it correct and explain why it works.
Sample Answer
Direct answer
The classic reordering bug is "message passing": a writer stores a payload, then raises a flag; a reader waits for the flag, then reads the payload. With every operation memory_order_relaxed, the reader can see the flag as 1 and the payload as 0, an outcome that sequential consistency (SC), the model where all operations appear to run in one agreed order that respects each thread's program order, forbids. The minimal fix is two arguments: make the flag store memory_order_release and the flag load memory_order_acquire.
Terms used here: a relaxed atomic operation guarantees only that the single operation is indivisible, with no ordering against other memory accesses; release on a store means every earlier read and write of that thread is visible to whoever reads that stored value with acquire; a weakly ordered CPU such as ARM may make two stores from one core visible to other cores in a different order than the program issued them.
The program
#include <atomic>
#include <cstdio>
#include <thread>
#include <vector>
// Message passing: the writer stores data, then raises a flag.
// The reader spins until the flag is 1, then reads data.
// Sequential consistency forbids "flag seen as 1 but data seen as 0".
#ifdef FIXED
constexpr auto kStore = std::memory_order_release;
constexpr auto kLoad = std::memory_order_acquire;
#else
constexpr auto kStore = std::memory_order_relaxed;
constexpr auto kLoad = std::memory_order_relaxed;
#endif
struct Slot {
alignas(64) std::atomic<int> data{0};
alignas(64) std::atomic<int> flag{0};
};
int main() {
constexpr int kSlots = 200000, kPasses = 20;
std::vector<Slot> slots(kSlots);
long bad = 0;
for (int pass = 0; pass < kPasses; ++pass) {
for (auto& s : slots) { s.data.store(0); s.flag.store(0); }
std::thread writer([&] {
for (auto& s : slots) {
s.data.store(1, std::memory_order_relaxed);
s.flag.store(1, kStore);
}
});
long local = 0;
std::thread reader([&] {
for (auto& s : slots) {
while (s.flag.load(kLoad) == 0) {}
if (s.data.load(std::memory_order_relaxed) == 0) ++local;
}
});
writer.join(); reader.join();
bad += local;
}
std::printf("forbidden outcomes (flag==1 but data==0): %ld out of %d\n",
bad, kSlots * kPasses);
}
Each slot has its own data and flag. alignas(64) places each on its own 64-byte cache line (the unit in which cores copy and exchange memory), so the two stores travel through the memory system independently and the reorder has room to show up; if both sat on one line they would be exchanged between cores together and the reorder would be much harder to provoke. The data operations are atomic only so that the program has no undefined behaviour: the bug here is ordering, not tearing.
Compiled in a gcc:14 container (GCC 14.4, aarch64, 14 cores visible) with g++ -O2 -std=c++20 -pthread, and run at several CPU counts:
- Relaxed build: on the 14-core container the count out of 4,000,000 slots was nonzero and a different number each time; the count is not reproducible, only its presence, and with two CPUs a run can still print 0. It needs two cores running at the same time: with the container confined to one CPU (
--cpuset-cpus 0) the same binary printed 0, because the threads then take turns and their stores and loads never overlap. -DFIXEDbuild (release store, acquire load): every run printedforbidden outcomes (flag==1 but data==0): 0 out of 4000000.
Because the machine in that container is ARM (a weakly ordered CPU), the run shows what the hardware really does, not just what the language permits. A compiler may also reorder relaxed operations on any CPU, so the relaxed version is wrong everywhere even when a given machine happens not to show it.
Why the fix works
flag.store(1, release) forbids moving the earlier data.store after it. flag.load(acquire) forbids moving the later data.load before it. When the load reads the 1 that the store wrote, the two form a synchronizes-with edge: everything the writer did before the store, including data = 1, happens-before everything the reader does after the load. Seeing flag == 1 then implies seeing data == 1. On ARM the compiler turns the pair into store-release and load-acquire instructions (or barriers), which is exactly the extra work the relaxed version skipped. Nothing stronger than release/acquire is needed here, so seq_cst would only add cost.
A second reorder that release/acquire does not fix
Store buffering: thread A does x = 1; r1 = y; and thread B does y = 1; r2 = x;. SC forbids r1 == 0 && r2 == 0 because whichever store came first in the global order is visible to the other thread's load. Each thread's store may still be waiting in its core's store buffer (a small queue of pending writes) when the other thread's load runs, which is how both loads can return 0:
#include <atomic>
#include <cstdio>
#include <thread>
#include <vector>
// Store buffering: A does x = 1; r1 = y. B does y = 1; r2 = x.
// Sequential consistency forbids r1 == 0 && r2 == 0.
#ifdef SEQCST
constexpr auto kOrder = std::memory_order_seq_cst;
#else
constexpr auto kOrder = std::memory_order_relaxed;
#endif
struct Slot {
alignas(64) std::atomic<int> x{0};
alignas(64) std::atomic<int> y{0};
int r1 = -1, r2 = -1; // each written by one thread only
};
std::atomic<int> arrived{0}; // lines the two threads up, round by round
int main() {
constexpr int kSlots = 200000;
std::vector<Slot> slots(kSlots);
auto lineUp = [](int round) {
arrived.fetch_add(1);
while (arrived.load() < 2 * (round + 1)) {}
};
std::thread a([&] {
for (int i = 0; i < kSlots; ++i) {
lineUp(i);
slots[i].x.store(1, kOrder);
slots[i].r1 = slots[i].y.load(kOrder);
}
});
std::thread b([&] {
for (int i = 0; i < kSlots; ++i) {
lineUp(i);
slots[i].y.store(1, kOrder);
slots[i].r2 = slots[i].x.load(kOrder);
}
});
a.join(); b.join();
long both_zero = 0;
for (auto& s : slots) if (s.r1 == 0 && s.r2 == 0) ++both_zero;
std::printf("r1==0 && r2==0 in %ld of %d rounds\n", both_zero, kSlots);
}
The lineUp helper makes the two threads start every round together, so their store and load really overlap; without it the threads drift apart and the outcome rarely appears. Compiled the same way in the gcc:14 container (aarch64) and run at several CPU counts per variant:
- All-relaxed: the count of the 200,000 rounds was nonzero on the 14-core container and a different number each time, and also with two CPUs. This also needs at least two CPUs: confined to one, the spin-wait in
lineUpmakes the program crawl instead of finishing promptly. -DSEQCST(seq_cston all four operations): every run printedr1==0 && r2==0 in 0 of 200000 rounds.- Release on the two stores and acquire on the two loads (edit
kOrderat those four operations): also 0 on every run on this aarch64 machine, because AArch64 implements release stores and acquire loads with store-release and load-acquire instructions (stlr/ldar) that are not reordered against each other. The C++ standard does not promise that result (it allows both threads to read 0 in this shape), and x86-64 hardware is allowed to show it: x86 lets a load complete before an earlier store to a different address is visible to other cores, and release/acquire compile to plainmovthere. No x86-64 result is shown here, so the aarch64 zero says nothing about x86. Portable code should useseq_csthere.
Example: a single-producer single-consumer ring buffer
A ring buffer (a fixed array used as a circular queue) with exactly one producer thread and one consumer thread needs release/acquire on its two indices and nothing stronger:
#include <atomic>
#include <cstdio>
#include <cstddef>
#include <thread>
// Single-producer single-consumer ring buffer. head is written only by the
// consumer, tail only by the producer. Capacity N-1 usable slots: one slot is
// always left empty so that head == tail can only mean "empty" and
// (tail + 1) % N == head can only mean "full".
template <typename T, size_t N>
class SpscRing {
T buf_[N];
alignas(64) std::atomic<size_t> head_{0}; // next slot to pop
alignas(64) std::atomic<size_t> tail_{0}; // next slot to push
public:
bool push(const T& v) { // producer thread only
size_t t = tail_.load(std::memory_order_relaxed); // own index
size_t next = (t + 1) % N;
if (next == head_.load(std::memory_order_acquire)) // consumer's progress
return false; // full
buf_[t] = v; // plain write
tail_.store(next, std::memory_order_release); // publish the slot
return true;
}
bool pop(T& out) { // consumer thread only
size_t h = head_.load(std::memory_order_relaxed); // own index
if (h == tail_.load(std::memory_order_acquire)) // producer's progress
return false; // empty
out = buf_[h]; // plain read
head_.store((h + 1) % N, std::memory_order_release);// free the slot
return true;
}
};
int main() {
constexpr int kItems = 1000000;
static SpscRing<int, 1024> ring;
long long sum = 0;
int out_of_order = 0;
std::thread producer([&] {
for (int i = 1; i <= kItems; ++i)
while (!ring.push(i)) {}
});
std::thread consumer([&] {
int expect = 1, v;
while (expect <= kItems) {
if (ring.pop(v)) {
if (v != expect) ++out_of_order;
sum += v; ++expect;
}
}
});
producer.join(); consumer.join();
std::printf("sum=%lld expected=%lld out_of_order=%d\n", sum,
(long long)kItems * (kItems + 1) / 2, out_of_order);
}
% N wraps an index back to 0 after the last slot, which is what makes the array circular. The producer is full when advancing tail_ would land on head_ (it must not overwrite a slot the consumer has not read); the consumer is empty when head_ has caught up with tail_. The producer's release store of tail_ publishes the plain write to buf_[t]; the consumer's acquire load of tail_ makes it visible. The consumer's release store of head_ tells the producer a slot may be reused, and the producer's acquire load of head_ ensures the consumer finished reading it first. Each side reads its own index with relaxed because no other thread writes it. Built with -fsanitize=thread and also with -O2, the program printed sum=500000500000 expected=500000500000 out_of_order=0 both times, with no TSan report. Downgrade either index store (the tail_ store or the head_ store) to relaxed and the plain buf_ access becomes a data race that TSan flags and that can lose or corrupt items on ARM. The design is only valid for exactly one producer and one consumer; with more threads the index updates need compare-and-swap loops (retry loops that read the index, compute the new value, and write it only if no other thread changed it in between) and a different design.
Explain the meaning and use-cases of the volatile keyword in Java and compare it with the volatile qualifier in C++ prior to C++11 and with C++11 atomics. What visibility and reordering guarantees does Java volatile provide?
Sample Answer
Direct answer
The same word means three different things:
- Java
volatileis a concurrency tool. A write to a volatile field happens-before every later read of that field (happens-before is the Java Memory Model's guarantee that one action's effects are visible to another). It gives visibility between threads, an ordering guarantee, and atomic reads and writes even forlonganddouble. It does not make compound operations likecount++atomic. - C++
volatile(before C++11, and still today) is not a concurrency tool. It tells the compiler that every access must really happen (no caching in a register, no removal), which is what you need for memory-mapped hardware registers. It gives no atomicity and no ordering between threads, and a data race on a volatile variable is still undefined behaviour (UB: the language places no requirements on what the program does, so results can be wrong in ways that vary by compiler and machine). - C++11
std::atomic<T>is the C++ equivalent of what Java volatile does: atomic operations with defined ordering rules, selected by amemory_orderargument. The default,seq_cst(sequentially consistent: all threads agree on one global order of atomic operations), is the strongest and the right choice unless you have a measured reason to weaken it; the weaker orders (relaxed,acquire,release) exist for expert tuning.
What Java volatile guarantees
From the Java Language Specification (chapter 17):
- Visibility and happens-before. "A write to a volatile field (§8.3.1.4) happens-before every subsequent read of that field." Everything the writing thread did before the write (including writes to ordinary fields) is visible to a thread that reads the new value.
- Order. All volatile accesses are part of one total synchronization order (a single sequence of all volatile reads and writes that every thread agrees on) that is consistent with each thread's program order (the order written in that thread's code). In practice, volatile accesses are not reordered with each other, and ordinary reads and writes cannot be moved across them in a way that breaks the happens-before edge.
- Atomic reads and writes. Reads and writes of volatile
longanddoubleare always atomic; for non-volatile ones the specification allows them to be treated as two 32-bit halves.
It does not give: atomic read-modify-write (x++ is read, add, write, so two threads can lose an update), mutual exclusion over several fields, or any guarantee about non-volatile fields that are not published through a volatile write.
Demonstration
public class Visibility {
static boolean plainStop = false; // no volatile
static volatile boolean volStop = false;
static int payload; // plain field published via the volatile flag
static volatile boolean ready;
static void spinPlain() throws Exception {
Thread t = new Thread(() -> { long n = 0; while (!plainStop) n++; });
t.start();
Thread.sleep(200);
plainStop = true;
t.join(2000);
System.out.println("plain flag: worker " + (t.isAlive() ? "STILL RUNNING after 2 s" : "stopped"));
if (t.isAlive()) System.exit(0);
}
static void spinVolatile() throws Exception {
Thread t = new Thread(() -> { long n = 0; while (!volStop) n++; });
t.start();
Thread.sleep(200);
volStop = true;
t.join(2000);
System.out.println("volatile flag: worker " + (t.isAlive() ? "STILL RUNNING" : "stopped"));
}
static void publish() throws Exception {
Thread reader = new Thread(() -> {
while (!ready) { }
System.out.println("reader saw payload = " + payload); // guaranteed 42
});
reader.start();
payload = 42; // plain write, ordered before the volatile write below
ready = true; // volatile write: happens-before every later read that sees true
reader.join();
}
public static void main(String[] a) throws Exception {
if (a.length > 0 && a[0].equals("plain")) spinPlain(); else { spinVolatile(); publish(); }
}
}
Run with java Visibility.java plain on OpenJDK 21.0.12.1 (container eclipse-temurin:21-jdk, aarch64), it printed plain flag: worker STILL RUNNING after 2 s at 2, 4 and 8 CPUs. Limited to one CPU the outcome varied from run to run: some runs printed stopped and others STILL RUNNING, because whether the JIT has compiled the loop within the 200 ms depends on how much CPU time its compiler threads get. Once the loop is compiled by the JIT (the JVM's just-in-time compiler, which turns hot bytecode into machine code while the program runs), the non-volatile flag is read once and the loop never sees the update. That is behaviour this JVM exhibited, not a guarantee: the specification permits the worker to see the change or not, and the timing makes it vary. Run with no argument it printed:
volatile flag: worker stopped
reader saw payload = 42
The second line is the use case that matters: payload is a plain field, yet it is guaranteed to be 42 because the write to it comes before the volatile write to ready and the reader saw ready == true.
C++: volatile is not atomic
#include <atomic>
#include <cstdio>
#include <thread>
#include <vector>
volatile int vcount = 0; // volatile int: NOT atomic
std::atomic<int> acount{0};
int main() {
std::vector<std::thread> ts;
for (int t = 0; t < 4; t++)
ts.emplace_back([] {
for (int i = 0; i < 100000; i++) {
vcount = vcount + 1; // load, add, store: races
acount.fetch_add(1, std::memory_order_relaxed);
}
});
for (auto& t : ts) t.join();
std::printf("volatile: %d atomic: %d (expected 400000)\n", (int)vcount, acount.load());
}
Compiled with g++ -std=c++20 -O2 -pthread (GCC 14.4, aarch64 container), the atomic total is 400000 at every CPU count, because fetch_add is atomic. The volatile total depends on how many CPUs the process may use: with the container limited to 2 or more CPUs it usually came out below 400000 and changed from run to run (on 2 CPUs some runs still printed exactly 400000; with 8 CPUs it was typically far lower). Limited to a single CPU it normally printed 400000, because the four threads rarely got interleaved in the middle of an increment, yet that is still a data race. The volatile total is never guaranteed to be low (a run in which the threads happen not to overlap can reach 400000); the point is that nothing makes it reliable. One sample run printed:
volatile: 119392 atomic: 400000 (expected 400000)
Read it as: four threads each tried 100,000 increments, so 400,000 is the correct total. The atomic figure matches it. The volatile figure is the number of increments that survived; every missing one is a case where two threads loaded the same old value and both stored old+1. Compiled with -fsanitize=thread, ThreadSanitizer (a compiler-inserted checker that reports data races while the program runs) printed WARNING: ThreadSanitizer: data race pointing at the volatile line, including on a run limited to one CPU whose volatile total happened to be correct. (The atomic uses memory_order_relaxed because a pure counter needs atomicity but no ordering; that is safe only because nothing else is published through it. For a flag that publishes data, use release on the store and acquire on the load (a pair that makes everything written before the store visible after the load), or simply leave the default seq_cst, which is what a beginner should write.)
Comparison
Java volatile | C++ volatile | C++ std::atomic | |
|---|---|---|---|
| Atomic read and write | yes (even long, double) | no | yes |
Atomic x++ | no | no | yes (fetch_add, ++) |
| Visibility between threads | yes | no | yes |
| Ordering of other memory | yes (happens-before) | none between threads | per memory_order |
| Concurrent unsynchronized access | not a data race; behaviour is defined | a data race: undefined behaviour | not a data race; behaviour is defined |
| Typical use | stop flag, publishing an object | memory-mapped hardware registers (device control addresses where every read and write must really happen) | counters, flags, lock-free structures (shared data structures built without mutexes) |
Use cases for Java volatile: a stop or status flag written by one thread and read by others; safely publishing an immutable or effectively immutable object; the one reference field in double-checked locking. When you need atomic read-modify-write or compare-and-set, use AtomicInteger and friends or VarHandle (the java.lang.invoke API that gives atomic and ordered operations on a field or array element); when you need to guard several variables together, use synchronized or a Lock.
In Go, write a function that runs multiple independent checks in parallel, returns all failures in a single error value, and preserves enough context for an operator to diagnose the problem quickly. The function should respect context cancellation by stopping new work and allowing in-flight checks to exit promptly.
Sample Answer
Design in one paragraph
Run each check in its own goroutine (Go's lightweight thread), at most limit at a time (enforced with a buffered channel used as a semaphore: sending into it takes one of limit slots and blocks when none is free, receiving gives a slot back), all sharing the caller's context.Context (a value that carries a cancellation signal and deadline; ctx.Done() is a channel that closes when it is cancelled). Each goroutine writes its failure into its own slot of a slice indexed by check position, so no mutex is needed and the result order is deterministic. After a sync.WaitGroup (a counter that Wait blocks on until every started goroutine has called Done) says every started goroutine has finished, errors.Join merges the non-nil slots into one error (it returns nil when there are none). The reason a skipped check never started is wrapped with %w in fmt.Errorf, which keeps the original error inside the new one, and CheckError has an Unwrap method, so errors.Is (does the chain contain this error?) and errors.As (find an error of this type in it) can still see through both layers. A recovered panic is recorded with %v, so it is text only. context.Cause(ctx) returns the reason recorded when the context ended; for a deadline that is context.DeadlineExceeded. Each failure is wrapped in a CheckError that names the check, how long it ran and the cause, which is the context an operator needs. Cancellation stops the loop that launches checks (with one narrow exception, described under "Cancellation checked twice"), and the running checks exit because they watch the same ctx.
Why not errgroup
golang.org/x/sync/errgroup with WithContext cancels the derived context on the first error, and its Wait returns only the first error, which is the opposite of "report all failures". The standard library pieces above give the required behaviour with no dependency.
The code
package main
import (
"context"
"errors"
"fmt"
"runtime"
"sort"
"strings"
"sync"
"sync/atomic"
"time"
)
// Check is one independent health check. It must return promptly when ctx is done.
type Check struct {
Name string
Run func(ctx context.Context) error
}
// CheckError records which check failed, how long it ran and why.
type CheckError struct {
Name string
Elapsed time.Duration
Err error
}
func (e *CheckError) Error() string {
if e.Elapsed == 0 {
return fmt.Sprintf("check %q: %v", e.Name, e.Err)
}
return fmt.Sprintf("check %q failed after %v: %v", e.Name, e.Elapsed.Round(time.Millisecond), e.Err)
}
func (e *CheckError) Unwrap() error { return e.Err }
// RunChecks runs checks concurrently (at most limit at a time) and returns nil
// or one error that joins every failure. Once the launcher sees ctx cancelled it
// starts no new check; checks not started are reported as skipped. (A check can
// still start in the instant between the ctx.Err() test and the select; it then
// sees the cancelled ctx.)
func RunChecks(ctx context.Context, limit int, checks []Check) error {
if limit < 1 {
limit = 1
}
sem := make(chan struct{}, limit)
errs := make([]error, len(checks)) // one slot per check: no lock needed, each goroutine writes its own index
var wg sync.WaitGroup
for i, c := range checks {
skipRest := ctx.Err() != nil // a cancelled context wins even if a slot is free
if !skipRest {
select {
case sem <- struct{}{}: // wait for a free slot, unless cancelled
case <-ctx.Done():
skipRest = true
}
}
if skipRest {
for j := i; j < len(checks); j++ {
errs[j] = &CheckError{Name: checks[j].Name, Err: fmt.Errorf("not started: %w", context.Cause(ctx))}
}
break
}
wg.Add(1)
go func(i int, c Check) {
defer wg.Done()
defer func() { <-sem }()
start := time.Now()
defer func() {
if r := recover(); r != nil { // a panicking check must not kill the process or hide the others
errs[i] = &CheckError{Name: c.Name, Elapsed: time.Since(start), Err: fmt.Errorf("panic: %v", r)}
}
}()
if err := c.Run(ctx); err != nil {
errs[i] = &CheckError{Name: c.Name, Elapsed: time.Since(start), Err: err}
}
}(i, c)
}
wg.Wait() // in-flight checks see the same ctx and exit on their own
return errors.Join(errs...) // errors.Join drops nil entries and returns nil if all are nil
}
func sleepCheck(name string, d time.Duration, fail error) Check {
return Check{name, func(ctx context.Context) error {
select {
case <-time.After(d):
return fail
case <-ctx.Done():
return ctx.Err()
}
}}
}
func lines(err error) string {
s := strings.Split(err.Error(), "\n")
sort.Strings(s)
return strings.Join(s, "\n ")
}
func main() {
base := runtime.NumGoroutine()
fmt.Println("1. all pass:", RunChecks(context.Background(), 4, []Check{
sleepCheck("db", 10*time.Millisecond, nil), sleepCheck("cache", 5*time.Millisecond, nil)}))
err := RunChecks(context.Background(), 4, []Check{
sleepCheck("db", 10*time.Millisecond, errors.New("connection refused")),
sleepCheck("cache", 5*time.Millisecond, nil),
sleepCheck("queue", 20*time.Millisecond, errors.New("lag 9000")),
{"flaky", func(context.Context) error { panic("nil map write") }},
})
var ce *CheckError
fmt.Printf("2. failures (%d lines):\n %s\n", strings.Count(err.Error(), "\n")+1, lines(err))
fmt.Println(" errors.As finds a CheckError:", errors.As(err, &ce))
ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond)
defer cancel()
slow := make([]Check, 6)
for i := range slow {
slow[i] = sleepCheck(fmt.Sprintf("slow-%d", i), 2*time.Second, nil)
}
t0 := time.Now()
err = RunChecks(ctx, 2, slow)
fmt.Printf("3. cancelled run returned after %v, deadline exceeded: %v\n %s\n",
time.Since(t0).Round(10*time.Millisecond), errors.Is(err, context.DeadlineExceeded), lines(err))
dead, stop := context.WithCancel(context.Background())
stop() // already cancelled before the call
var started int32
probe := Check{"probe", func(context.Context) error { atomic.AddInt32(&started, 1); return nil }}
err = RunChecks(dead, 4, []Check{probe, probe, probe})
fmt.Println("4. pre-cancelled: checks started =", atomic.LoadInt32(&started), "canceled:", errors.Is(err, context.Canceled))
time.Sleep(50 * time.Millisecond)
fmt.Println("5. goroutines before/after:", base, runtime.NumGoroutine())
}
Run in a golang:1.23 container (Go 1.23.12, linux/arm64), the program passes go vet; built with -race and run repeatedly, including with GOMAXPROCS=1 and GOMAXPROCS=2, it exits 0, the race detector reports nothing, and the output has the same shape each time apart from the millisecond figures. The program needs Go 1.20 or later, because errors.Join and context.Cause were added in 1.20. One run printed:
1. all pass: <nil>
2. failures (3 lines):
check "db" failed after 10ms: connection refused
check "flaky" failed after 0s: panic: nil map write
check "queue" failed after 21ms: lag 9000
errors.As finds a CheckError: true
3. cancelled run returned after 50ms, deadline exceeded: true
check "slow-0" failed after 55ms: context deadline exceeded
check "slow-1" failed after 55ms: context deadline exceeded
check "slow-2": not started: context deadline exceeded
check "slow-3": not started: context deadline exceeded
check "slow-4": not started: context deadline exceeded
check "slow-5": not started: context deadline exceeded
4. pre-cancelled: checks started = 0 canceled: true
5. goroutines before/after: 1 1
Reading the output
- Case 2. The
flakycheck panics at once.recover()(callable only inside a deferred function, it stops the panic in that goroutine and returns the panic value) turns it into aCheckError. Its elapsed time is a few microseconds, whichRound(time.Millisecond)normally shows as0s;Error()only drops the "failed after" part whenElapsedis exactly zero, as for checks that never started.dbandqueueshow their own durations, a little above their 10 ms and 20 ms sleeps.linessorts the lines so the output does not depend on goroutine finish order (the slots are already in input order, but the sort keeps the transcript stable). - Case 3, step by step (limit 2, six checks of 2 seconds each, deadline 50 ms). At 0 ms, check 0: the context is live, a slot is free,
sem <- struct{}{}succeeds (struct{}{}is the empty value, zero bytes, used only as a token) andslow-0starts. Check 1 does the same, so both slots are now taken. Check 2:selectfinds the semaphore full and the context live, so the launching loop blocks there. At 50 ms the deadline closesctx.Done().slow-0andslow-1wake and returnctx.Err(); the blockedselectin the loop also wakes through itsctx.Done()case, setsskipRest, and the inner loop fills slots 2 to 5 with "not started: context deadline exceeded" beforebreak.wg.Wait()then waits only for the two goroutines that exist. Checks 2 to 5 never ran at all. - Why 55 ms and not 50. The deadline fires at 50 ms; the figure is measured from the goroutine's start until it returns, so it also includes the timer firing and the goroutine being scheduled. The sample run above shows 55 ms; other runs show anything from 50 ms to several milliseconds more, with the excess larger when the machine is busy or the race detector is on, while the call as a whole returns shortly after the 50 ms deadline. The extra milliseconds are wake-up latency, not extra work, so a test should never assert an exact figure.
- The passed-in
i, c. The goroutine closure takesiandcas parameters so each goroutine gets its own copy of the loop's values instead of sharing the loop variables. Since Go 1.22, in a module whosego.moddeclaresgo 1.22or later, every loop iteration already has its own variables, so the parameters are belt and braces there; before 1.22 they are required, because the closure would otherwise seeiandcchange under it.
Choices worth defending
- Per-index result slots, not a shared slice with append. Appending from many goroutines is a data race (two goroutines writing the same memory without ordering). Writing distinct elements is safe, and
wg.Wait()establishes a happens-before edge (a guarantee in Go's memory model that everything one goroutine did beforeDoneis visible to the code afterWait), which makes the writes visible to the reader. - Bounded concurrency (
limit). Without it, 5,000 checks means 5,000 simultaneous connections. The semaphore (a buffered channel used as a counter of free slots) also makes "stop new work" real: withlimit2 and a 50 ms deadline, checks 2 to 5 never started (output case 3), where unbounded fan-out would have launched all of them first. - Cancellation checked twice.
ctx.Err()is tested before theselect, because when a slot is free and the context is done,select(which waits on several channel operations and runs one that is ready; if several are ready at once it picks one at random) picks either ready case at random. Case 4 shows a pre-cancelled context starts zero checks. The test narrows the window but cannot close it: if the context is cancelled after thectx.Err()test and before theselect, and a slot is free, theselectcan still pick the send and start one more check. That check sees the cancelled context and returns at once, so the cost is small, but "no new check starts after cancellation" holds only for cancellation the launcher has already observed. - Skipped checks are reported, not hidden. They appear as "not started: context deadline exceeded", with
context.Causepreserving why, so the operator can tell "never ran" from "ran and failed", anderrors.Is(err, context.DeadlineExceeded)still works through the joined error. recoverper goroutine. A panic in a goroutine you did not recover crashes the whole process; here it becomes one more failure line (flakyabove).- No leak. Goroutine count is 1 before and after (case 5), because every goroutine exits when its check returns.
The caveat that matters
Go cannot kill a goroutine, so cancellation is cooperative:
RunCheckswaits for in-flight checks, so it returns promptly only if eachCheck.Runhonoursctx. A check that ignores the context (a bare blocking call) makesRunCheckswait for it.- The fix is at the call sites: use context-aware clients (
http.NewRequestWithContext,db.QueryContext, a dialer with a deadline). - Returning early without waiting would hide the problem and let goroutines write into
errsafter the caller has read it, which is a data race. - If a bounded wait is mandatory anyway, call
RunChecksin its own goroutine andselecton its result channel against a timer, accepting that the stragglers keep running until they return, and log that they did.
Explain the practical differences between concurrency models in Python: threading, multiprocessing, and asyncio. For each model, state when it is appropriate (network-bound, CPU-bound, I/O-bound), how the Global Interpreter Lock (GIL) affects behavior, and how you would pick one for a systems scripting task that polls thousands of sockets.
Sample Answer
Direct answer
For a script that polls thousands of sockets (waiting on the network, almost no CPU per socket), use asyncio: one thread runs an event loop (a loop that waits for any of many sockets to become ready, then runs the code waiting on that socket) and handles all of them with little memory and no locking. The code that waits is written as a coroutine: a function declared with async def that can pause at each await (for example await reader.readline()) and hand control back to the loop, which resumes it when the awaited thing is ready. Use threads when the code you must call is blocking and has no async version. Use multiprocessing only when the work is CPU-bound Python and you need several cores.
What the GIL changes
The GIL (Global Interpreter Lock) is a lock in the standard CPython interpreter that lets only one thread execute Python bytecode (the low-level instructions the interpreter runs) at a time. Consequences:
- Threads running pure Python do not run in parallel: two CPU-bound threads take turns.
- A thread waiting on I/O (
socket.recv,time.sleep, file reads) releases the GIL, so other threads run while it waits. Threads therefore do help I/O-bound work. - Some C extensions (NumPy, hashing, compression) release the GIL during heavy work, so threads can speed those up too.
- Processes each have their own interpreter and their own GIL, so they run in parallel across cores, at the cost of process start-up, memory per process, and pickling (serializing an object to bytes with Python's
picklemodule so it can cross a process boundary) every argument and result. - asyncio is single-threaded, so the GIL is not even a factor; the cost is that any blocking call inside a coroutine freezes the entire loop. A tiny demonstration: a ticker prints every 0.2 s while another coroutine does
time.sleep(1.0)(a blocking call) after 0.3 s:
import asyncio, time
async def ticker(t0):
for _ in range(5):
print(f"tick at {time.perf_counter() - t0:.1f} s")
await asyncio.sleep(0.2)
async def blocker(t0):
await asyncio.sleep(0.3)
time.sleep(1.0) # blocking call: freezes the whole loop
async def main():
t0 = time.perf_counter()
await asyncio.gather(ticker(t0), blocker(t0))
asyncio.run(main())
A sample run in a python:3.12-slim container printed:
tick at 0.0 s
tick at 0.2 s
tick at 1.3 s
tick at 1.5 s
tick at 1.7 s
The ticks at 0.4 s to 1.2 s are missing: the whole loop stood still for the second the blocking call held it. Replacing time.sleep(1.0) with await asyncio.sleep(1.0) in blocker gave ticks at 0.0, 0.2, 0.4, 0.6 and 0.8 s with no gap. (Exact tick times shift by a few hundredths of a second between runs; the missing stretch is what matters.)
Status of the GIL (this is background depth; a PEP is a Python Enhancement Proposal, the design document for a language change): PEP 703 added an optional free-threaded build (no GIL) of CPython, experimental in 3.13. PEP 779 (status Final) made that build officially supported but still optional from Python 3.14; the default build still has the GIL, and extension compatibility and single-thread speed are the trade-offs (the 3.14 release notes put the single-thread penalty at roughly 5 to 10 percent, depending on platform and compiler). So do not design around a GIL-free interpreter being present: on Python 3.13 and later, sys._is_gil_enabled() tells you whether the interpreter you will deploy to is running with the GIL.
Measured comparison
import asyncio, time, threading
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
def cpu_work(n):
s = 0
for i in range(n):
s += i * i
return s
def io_work(_):
time.sleep(0.5) # stands in for a blocking network call; releases the GIL
async def aio_work(_):
await asyncio.sleep(0.5)
def timed(label, fn):
t = time.perf_counter(); fn(); print(f"{label:<28}{time.perf_counter() - t:6.2f} s")
def main():
N, W = 4_000_000, 4
print("CPU-bound, 4 jobs of 4M iterations")
timed(" sequential", lambda: [cpu_work(N) for _ in range(W)])
with ThreadPoolExecutor(W) as ex:
timed(" 4 threads", lambda: list(ex.map(cpu_work, [N] * W)))
with ProcessPoolExecutor(W) as ex:
list(ex.map(cpu_work, [1] * W)) # warm up the worker processes
timed(" 4 processes", lambda: list(ex.map(cpu_work, [N] * W)))
print("I/O-bound, 40 waits of 0.5 s")
timed(" sequential", lambda: [io_work(0) for _ in range(40)])
with ThreadPoolExecutor(40) as ex:
timed(" 40 threads", lambda: list(ex.map(io_work, range(40))))
async def many(): await asyncio.gather(*(aio_work(i) for i in range(40)))
timed(" asyncio, 40 tasks", lambda: asyncio.run(many()))
if __name__ == "__main__":
main()
One sample run, Python 3.12.15 in a python:3.12-slim container on a 14-core aarch64 host:
CPU-bound, 4 jobs of 4M iterations
sequential 0.42 s
4 threads 0.41 s
4 processes 0.15 s
I/O-bound, 40 waits of 0.5 s
sequential 20.04 s
40 threads 0.51 s
asyncio, 40 tasks 0.51 s
Reading it: ex.map(f, items) runs f on every item in the pool and returns results in order. The warm-up line sends tiny jobs first so process start-up is not counted in the timing, and asyncio.gather(*coros) starts all the coroutines at once and waits for all of them. For CPU-bound work, threads gave no speed-up (0.41 s against 0.42 s sequential) because of the GIL, and processes were about 2.8 times faster (0.15 s; with 4 workers the ideal would be 4 times, and the gap is process overhead on a short job). That speed-up needs at least 4 free cores: the same script in a container limited to 1 CPU showed the 4 processes no faster than sequential. For waiting, threads and asyncio both collapse 40 waits of 0.5 s to the length of one wait, 0.51 s. Timings vary run to run; the pattern is the finding.
Polling thousands of sockets
import asyncio, resource, time
N = 3000 # concurrent connections; 3000 client + 3000 server sockets = 6000 descriptors
# Many shells start with a soft limit of 1,024: raise it as far as the hard limit allows.
soft, hard = resource.getrlimit(resource.RLIMIT_NOFILE)
if soft < 8192:
resource.setrlimit(resource.RLIMIT_NOFILE, (8192 if hard == resource.RLIM_INFINITY else min(8192, hard), hard))
async def handle(reader, writer):
data = await reader.readline()
writer.write(data.upper()); await writer.drain()
writer.close()
async def client(port, i, ok):
r, w = await asyncio.open_connection("127.0.0.1", port)
w.write(f"ping {i}\n".encode()); await w.drain()
reply = await r.readline()
ok.append(reply.startswith(b"PING"))
w.close()
async def main():
server = await asyncio.start_server(handle, "127.0.0.1", 0, backlog=N)
port = server.sockets[0].getsockname()[1]
ok = []
t = time.perf_counter()
await asyncio.gather(*(client(port, i, ok) for i in range(N)))
print(f"{len(ok)} of {N} sockets answered, all correct: {all(ok)}, {time.perf_counter() - t:.2f} s, one thread")
server.close()
print("fd limit:", resource.getrlimit(resource.RLIMIT_NOFILE)[0])
asyncio.run(main())
In the same container (soft limit already 20,480, so no raise was needed) it printed:
fd limit: 20480
3000 of 3000 sockets answered, all correct: True, 0.26 s, one thread
Walk through the program. asyncio.start_server(handle, host, 0, backlog=N) listens on a free port (0 lets the OS choose) and runs handle as a new coroutine for each accepted connection; backlog is how many not-yet-accepted connections the OS may queue, raised so a burst of 3,000 is not refused. reader.readline() waits for one line without blocking the loop. writer.write(...) queues bytes and await writer.drain() waits until the send buffer has room. asyncio.open_connection is the client side and returns a reader and writer pair. The 3,000 client coroutines run concurrently under gather. That is 3,000 client connections and 3,000 matching server connections (6,000 file descriptors) served from one thread on loopback. The file-descriptor limit printed first matters: each socket uses one file descriptor (the small integer the OS uses to name an open file or socket), and many systems default to 1,024, so a script that polls thousands of sockets must raise ulimit -n (the shell command that sets the per-process descriptor limit) or call resource.setrlimit, as the program above does, regardless of the concurrency model. Without that, this program fails with OSError: [Errno 24] Too many open files when the soft limit is 1,024. Loopback with an instant echo is the easy case; against real remote servers the elapsed time is dominated by network latency, which is exactly what concurrency hides.
Choosing, with the numbers
| Model | Memory and cost for 3,000 waiters | GIL effect | Pick it when |
|---|---|---|---|
| asyncio | 3,000 small coroutine objects on one thread | none | the work is waiting on sockets and you can use async libraries |
| threads | 3,000 OS threads, each with its own stack | released during I/O, serializes pure Python | you must call blocking libraries; keep the pool bounded (tens to low hundreds) |
| multiprocessing | one interpreter per process, pickling | each has its own GIL | CPU-bound pure-Python work on several cores |
For the socket-polling script, recommend asyncio with a cap on concurrent connections (an asyncio.Semaphore) and timeouts (asyncio.timeout, available from Python 3.11, or asyncio.wait_for, which also works on older versions) so one dead peer cannot hang the run. On Linux, asyncio uses epoll (a kernel facility that reports which of many sockets are ready without scanning them all) through the selectors module (Python's wrapper over such facilities), whereas a hand-written select() loop cannot watch descriptors numbered 1,024 or higher (the FD_SETSIZE limit, a fixed size of the descriptor set that select() accepts). You do not need these names to use asyncio; they explain why it scales. If part of the job is blocking or CPU-heavy (parsing large replies), hand that part to loop.run_in_executor or asyncio.to_thread, so the loop keeps answering sockets. What would flip the choice: an existing blocking client library with no async version (threads), or heavy per-reply computation (add a process pool behind the loop).
Implement a thread-safe bounded blocking queue in Java that supports multiple producers and consumers. Explain your synchronization strategy (wait/notify vs Lock/Condition), how you avoid deadlocks and spurious wakeups, and how you'd test correctness under high concurrency. Provide code and narrate the reasoning during the implementation.
Sample Answer
Direct answer
A bounded blocking queue is a circular array guarded by one lock, with two wait queues: producers wait while the queue is full, consumers wait while it is empty. Use ReentrantLock with two Condition objects (notFull, notEmpty) so a put wakes only consumers and a take wakes only producers. Every wait sits inside a while loop that re-tests the condition, because a woken thread is not guaranteed the condition still holds. The same design with synchronized and wait/notifyAll is correct but wakes everyone, since a monitor has a single wait set. (A monitor is the lock built into every Java object, the one synchronized takes. Its wait set is the list of threads parked inside wait() on that object.)
Synchronization strategy: wait/notify versus Lock/Condition
synchronized + wait/notify | ReentrantLock + Condition | |
|---|---|---|
| Wait sets per lock | One. Producers and consumers sleep in the same set. | As many as you create (here two). |
| Wake-up call | notifyAll() is the safe choice, because notify() picks an arbitrary thread (the Javadoc says the choice is arbitrary) and might wake another producer when a consumer was needed. | signal() on the right condition wakes a thread that can actually proceed. |
| Timed or interruptible waits | wait(timeout) and interruption work, but lock acquisition itself cannot be interrupted or timed out. | lockInterruptibly(), tryLock(timeout) and awaitNanos give interruptible and timed operations (used in poll below). |
| Cost | notifyAll wakes threads that go straight back to sleep (a "thundering herd": many threads woken at once to compete for one thing). | Fewer useless wake-ups. |
Recommendation: use Lock and Condition when you need two wait conditions or timeouts, which a real queue does. Use plain synchronized when one condition is enough and simplicity matters.
Spurious wakeups and why while
Two different things can wake a thread with the condition still false, and while covers both:
-
Spurious wakeup. The
Object.waitJavadoc says a thread can wake without being notified, interrupted or timing out, and recommends "to check the condition being awaited in a while loop around the call to wait". TheConditionJavadoc says a "spurious wakeup" is permitted to occur, as a concession to the underlying platform semantics, and that aCondition"should always be waited upon in a loop". It is rare and allowed by the specification, so the code must tolerate it. -
Stolen wakeup. A thread that was signalled does not run immediately; it has to reacquire the lock first, and another thread can get there before it. A timeline with real values, capacity 8, queue empty:
- Consumer A calls
take, seescount == 0and callsawait. The lock is released and A sleeps. - A producer calls
put(42):countbecomes 1 and it signals. A is now runnable but must queue for the lock. - Consumer B, which was never waiting, calls
takeand gets the lock first.countis 1, so B takes 42 andcountis back to 0. - A reacquires the lock and returns from
await. Withwhile, A re-testscount == 0, finds it true and waits again. Withif, A readsitems[head], an empty slot.
This case needs no spurious wakeup at all, which is why the
ifversion in the demo below fails on every run. - Consumer A calls
How deadlock is avoided
A deadlock needs threads holding one resource while waiting for another. Here there is exactly one lock, it is never held while calling out to user code, and it is released in finally. await() atomically releases the lock while the thread sleeps and reacquires it before returning, so a waiting producer never blocks the consumer that would wake it. The remaining hang risk is a missed signal: a waiter that goes to sleep after the one signal that would have woken it has already been sent. Every state change that can make a waiter's condition true (put makes the queue non-empty, take makes it non-full) is therefore followed by a signal on that condition, made while still holding the lock. Holding the lock is what rules out the missed signal: a consumer tests count == 0 and calls await as one step with respect to the lock (await releases it only once the thread is parked), so a put either runs entirely before the consumer's test, and the consumer sees count == 1, or runs after the consumer is parked, and its signal reaches it. There is no gap in between. The language also enforces the rule: Condition.signal() called without holding the lock throws IllegalMonitorStateException for a ReentrantLock condition, and so does notify() outside a synchronized block. Both throw IllegalMonitorStateException on Java 21, so signalling after unlock() is not a quiet bug but an immediate exception.
Complete runnable code (Java 21)
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicIntegerArray;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.ReentrantLock;
public class BoundedQueueDemo {
interface BQ {
void put(int x) throws InterruptedException;
int take() throws InterruptedException;
}
/** One lock, two conditions: producers wait on notFull, consumers wait on notEmpty. */
static class LockQueue implements BQ {
private final int[] items;
private int head, tail, count;
private final ReentrantLock lock = new ReentrantLock();
private final Condition notFull = lock.newCondition();
private final Condition notEmpty = lock.newCondition();
LockQueue(int capacity) { items = new int[capacity]; }
public void put(int x) throws InterruptedException {
lock.lockInterruptibly();
try {
while (count == items.length) notFull.await(); // while, never if
items[tail] = x;
tail = (tail + 1) % items.length;
count++;
notEmpty.signal(); // one consumer can proceed
} finally { lock.unlock(); }
}
public int take() throws InterruptedException {
lock.lockInterruptibly();
try {
while (count == 0) notEmpty.await();
int x = items[head];
head = (head + 1) % items.length;
count--;
notFull.signal(); // one producer can proceed
return x;
} finally { lock.unlock(); }
}
/** Timed take: returns -1 on timeout. awaitNanos returns the remaining time. */
public int poll(long timeout, TimeUnit unit) throws InterruptedException {
long nanos = unit.toNanos(timeout);
lock.lockInterruptibly();
try {
while (count == 0) {
if (nanos <= 0) return -1;
nanos = notEmpty.awaitNanos(nanos);
}
int x = items[head];
head = (head + 1) % items.length;
count--;
notFull.signal();
return x;
} finally { lock.unlock(); }
}
}
/** Same logic with synchronized + wait/notifyAll: one wait set, so notifyAll is required. */
static class MonitorQueue implements BQ {
private final int[] items;
private int head, tail, count;
MonitorQueue(int capacity) { items = new int[capacity]; }
public synchronized void put(int x) throws InterruptedException {
while (count == items.length) wait();
items[tail] = x;
tail = (tail + 1) % items.length;
count++;
notifyAll(); // wakes producers AND consumers; each re-checks its own condition
}
public synchronized int take() throws InterruptedException {
while (count == 0) wait();
int x = items[head];
head = (head + 1) % items.length;
count--;
notifyAll();
return x;
}
}
/** Deliberately wrong: `if` instead of `while`. Used only to prove the harness can fail. */
static class IfQueue implements BQ {
private final int[] items;
private int head, tail, count;
IfQueue(int capacity) { items = new int[capacity]; }
public synchronized void put(int x) throws InterruptedException {
if (count == items.length) wait();
items[tail] = x;
tail = (tail + 1) % items.length;
count++;
notifyAll();
}
public synchronized int take() throws InterruptedException {
if (count == 0) wait();
if (count == 0) throw new IllegalStateException("woke up to an empty queue");
int x = items[head];
head = (head + 1) % items.length;
count--;
notifyAll();
return x;
}
}
static final int PRODUCERS = 4, CONSUMERS = 4, PER_PRODUCER = 50_000, CAPACITY = 8;
/** Returns a short verdict string; "OK" only if nothing was lost, duplicated or reordered. */
static String run(BQ q) throws Exception {
int total = PRODUCERS * PER_PRODUCER;
AtomicIntegerArray seen = new AtomicIntegerArray(total);
AtomicInteger problems = new AtomicInteger();
AtomicInteger crashed = new AtomicInteger();
Thread[] ts = new Thread[PRODUCERS + CONSUMERS];
for (int p = 0; p < PRODUCERS; p++) {
final int base = p * PER_PRODUCER; // value = base + sequence number
ts[p] = new Thread(() -> {
try { for (int i = 0; i < PER_PRODUCER; i++) q.put(base + i); }
catch (Throwable t) { crashed.incrementAndGet(); }
});
}
for (int c = 0; c < CONSUMERS; c++) {
ts[PRODUCERS + c] = new Thread(() -> {
int[] lastFrom = new int[PRODUCERS];
java.util.Arrays.fill(lastFrom, -1);
try {
for (int i = 0; i < total / CONSUMERS; i++) {
int v = q.take();
if (seen.incrementAndGet(v) != 1) problems.incrementAndGet(); // duplicate
int p = v / PER_PRODUCER, seq = v % PER_PRODUCER;
if (seq <= lastFrom[p]) problems.incrementAndGet(); // FIFO per producer
lastFrom[p] = seq;
}
} catch (Throwable t) { crashed.incrementAndGet(); }
});
}
for (Thread t : ts) t.start();
for (Thread t : ts) t.join(20_000);
boolean stuck = false;
for (Thread t : ts) if (t.isAlive()) { stuck = true; t.interrupt(); }
int missing = 0;
for (int i = 0; i < total; i++) if (seen.get(i) == 0) missing++;
if (!stuck && crashed.get() == 0 && problems.get() == 0 && missing == 0) return "OK";
return "FAIL stuck=" + stuck + " crashed=" + crashed + " dupOrOrder=" + problems + " missing=" + missing;
}
public static void main(String[] args) throws Exception {
for (int round = 1; round <= 3; round++) {
System.out.println("round " + round
+ " LockQueue: " + run(new LockQueue(CAPACITY))
+ " | MonitorQueue: " + run(new MonitorQueue(CAPACITY)));
}
System.out.println("IfQueue (wrong on purpose): " + run(new IfQueue(CAPACITY)));
LockQueue empty = new LockQueue(2);
long t0 = System.nanoTime();
int r = empty.poll(100, TimeUnit.MILLISECONDS);
long ms = (System.nanoTime() - t0) / 1_000_000;
System.out.println("poll on empty queue returned " + r + ", waited at least 100 ms: " + (ms >= 100));
}
}
Run with java BoundedQueueDemo.java (Java 21). The IfQueue class is wrong on purpose: it exists to show that the test harness can fail. Output when run in an eclipse-temurin:21-jdk Linux container. The three rounds print the same every time; the IfQueue numbers are different on every run, so they are shown as placeholders:
round 1 LockQueue: OK | MonitorQueue: OK
round 2 LockQueue: OK | MonitorQueue: OK
round 3 LockQueue: OK | MonitorQueue: OK
IfQueue (wrong on purpose): FAIL stuck=<false or true> crashed=<varies> dupOrOrder=<large> missing=<large>
poll on empty queue returned -1, waited at least 100 ms: true
The two correct queues pass every round (4 producers and 4 consumers moving 200,000 items through 8 slots) with no lost, duplicated or reordered item, and IfQueue fails. The counts on its failure line, including how many threads crash (sometimes none) and whether stuck is true, differ from run to run; the failure does not.
Reading the harness. Each value is base + sequence number, where base = producer index * 50,000, so v / PER_PRODUCER recovers the producer and v % PER_PRODUCER its sequence number. It checks four things:
seenis anAtomicIntegerArraywith one counter per value.seen.incrementAndGet(v) != 1means this value was already taken once: a duplicate. After the join, any counter still 0 is a value that never arrived (missing).lastFrom[p]is each consumer's private memory of the last sequence number it saw from producerp. One producer's items enter the queue in order and one consumer takes them in queue order, so a sequence number less than or equal to the last one means an order violation. Both this and a duplicate add toproblems, which the output callsdupOrOrder.crashedcounts threads that died with an exception (forIfQueue, the deliberateIllegalStateException("woke up to an empty queue")).t.join(20_000)waits at most 20 seconds per thread; a thread still alive after that isstuck, which is how a lost signal shows up.
Reading the failure line. IfQueue breaks the queue's invariants in more than one way, so the counts on that line are large and change from run to run, and a run can even finish with crashed=0. When a consumer is woken by the if version and finds an empty queue, it throws (that is what crashed counts). A second effect does most of the damage: IfQueue.put tests count == items.length with equality, so once a producer wakes while the queue is still full and writes anyway, count can pass the capacity and the test may never be true again. Producers then stop waiting and overwrite slots no consumer has read, while consumers re-read stale slots. Re-read slots show up as duplicates or out-of-order values (dupOrOrder), and the overwritten values never reach anyone (missing). Because one value can be both a duplicate and out of order, dupOrOrder can even exceed the 200,000 items. Almost every one of the 200,000 values ends up in one of those two counts, so the numbers are of the order of the total, not of a few unlucky items. stuck is usually false because the failure here is lost and repeated data rather than a lost signal, but it can be true: after the overwrites, a consumer can be left waiting for an item that never arrives and is still alive when the 20-second join times out (each thread still stuck costs the harness up to 20 more seconds, because the joins wait one after another, so such a run is slow).
How I would test correctness under high concurrency
- Conservation: put 200,000 distinct values, take all of them, assert each arrives exactly once (the
AtomicIntegerArray). - Order: FIFO means each consumer sees any one producer's values in increasing order.
- Capacity: with capacity 8 and 4 producers, producers constantly block, so the
notFullpath is exercised. Add tiny capacities (1 and 2) as boundary cases. - Liveness: join with a timeout and fail on a live thread; a lost signal shows up as a hang, not a wrong value.
- Timed operations:
pollon an empty queue returns the timeout result after at least the requested time (checked in the demo). - Limits of testing: passing runs never prove absence of a race. Mutation (deliberately breaking the code, here
ifforwhile, and checking the test notices) shows the test has teeth. A tool such as OpenJDK's jcstress (a harness that runs tiny multi-thread scenarios millions of times and records every outcome seen) targets memory-ordering bugs, where one thread sees another's writes late or in an unexpected order, that a stress loop can miss.
Complexity and edge cases
put and take are O(1). Memory is O(capacity). Edge cases: capacity 1 works (a single slot, head == tail always), capacity 0 or negative must be rejected in the constructor, null is impossible here because the slots hold int, an interrupt during await throws InterruptedException with the lock reacquired first and then released by finally, and the lock is non-fair by default (a fair lock would serve the longest-waiting thread first), so a waiting thread can lose to a newly arrived one (acceptable here; new ReentrantLock(true) trades throughput for fairness). In production code use java.util.concurrent.ArrayBlockingQueue, which implements this same design; the exercise is about being able to build it.
Unlock Full Question Bank
Get access to all 29 Language-Level Concurrency and Multithreading interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.