Concurrency
Blocking Queue — Producer-Consumer
Implement a thread-safe bounded blocking queue from scratch using ReentrantLock and Condition variables — the canonical producer-consumer concurrency problem.
Problem#
Implement a thread-safe bounded queue with a fixed capacity, supporting:
| Method | Behaviour |
|---|---|
| enqueue(x) | Adds an element. If the queue is full, the calling thread blocks until space is available. |
| dequeue() | Removes and returns the front element. If the queue is empty, the calling thread blocks until an element is available. |
| size() | Returns the current number of elements. |
Multiple producer threads and multiple consumer threads call these concurrently. The implementation must be correct under contention — no lost items, no reading from an empty queue, no writing to a full one.
Think Before Coding#
Work through these before looking at the solution:
-
What shared state needs protection? The underlying queue and its size are accessed and modified by multiple threads — they must always be accessed under a lock.
-
What are the two distinct conditions threads wait on?
- Producers wait when the queue is full (size == capacity)
- Consumers wait when the queue is empty (size == 0)
-
Why two separate condition variables instead of one? Using a single condition (or notifyAll) wakes all waiting threads — both producers and consumers — whenever anything changes. That's wasted work and introduces subtle bugs. Two conditions let you wake only the right party:
- After a successful enqueue → signal consumers (notEmpty)
- After a successful dequeue → signal producers (notFull)
Implementation#
package BlockingQueue;
import java.util.LinkedList;
import java.util.Queue;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.ReentrantLock;
public class BlockingQueue {
private final Queue<Integer> q = new LinkedList<>();
private final int capacity;
private final ReentrantLock lock = new ReentrantLock();
private final Condition notFull = lock.newCondition();
private final Condition notEmpty = lock.newCondition();
public BlockingQueue(int cap) {
this.capacity = cap;
}
public void enqueue(int item) throws InterruptedException {
lock.lock();
try {
while (q.size() == capacity) // loop — not if — guards against spurious wakeups
notFull.await();
q.offer(item);
notEmpty.signal(); // wake one waiting consumer
} finally {
lock.unlock();
}
}
public int dequeue() throws InterruptedException {
lock.lock();
try {
while (q.isEmpty()) // loop — same reason
notEmpty.await();
int val = q.poll();
notFull.signal(); // wake one waiting producer
return val;
} finally {
lock.unlock();
}
}
public int size() {
lock.lock();
try {
return q.size();
} finally {
lock.unlock();
}
}
}
Key Design Decisions#
ReentrantLock over synchronized
synchronized gives you a single implicit condition (via wait/notify). ReentrantLock lets you create multiple named conditions from one lock — exactly what we need here to wake producers and consumers separately.
while loop, not if, around await()
await() can return spuriously (OS-level wakeup unrelated to a signal). If you use if, the thread proceeds even though the condition may still not hold. The while loop re-checks and goes back to sleep if needed. Always loop around await().
signal() not signalAll()
We only need to wake one thread — the next producer or consumer in line. signalAll() is safe but wastes CPU by waking all waiting threads, which then contend for the lock and mostly go right back to sleep.
Two separate Conditions
notFull — producers park here when the queue is full. notEmpty — consumers park here when the queue is empty.
After enqueue, we signal notEmpty (a consumer might now be able to proceed). After dequeue, we signal notFull (a producer might now be able to proceed). Signalling the wrong condition would leave threads parked indefinitely.
Demo#
package BlockingQueue;
public class Demo {
public static void main(String[] args) throws InterruptedException {
// Part 1: basic single-threaded sanity check
BlockingQueue queue = new BlockingQueue(5);
queue.enqueue(1);
queue.enqueue(2);
queue.enqueue(3);
System.out.println("Queue size: " + queue.size()); // 3
queue.dequeue();
System.out.println("Queue size after dequeue: " + queue.size()); // 2
System.out.println("---- concurrent demo ----");
// Part 2: real producer/consumer — capacity-3 queue forces blocking
BlockingQueue bq = new BlockingQueue(3);
// Producer: tries to add 6 items; blocks when queue fills up
Thread producer = new Thread(() -> {
try {
for (int i = 1; i <= 6; i++) {
bq.enqueue(i);
System.out.println("Produced: " + i + " (size=" + bq.size() + ")");
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, "producer");
// Consumer: drains slowly — forces the producer to wait for space
Thread consumer = new Thread(() -> {
try {
for (int i = 1; i <= 6; i++) {
Thread.sleep(200);
System.out.println(" Consumed: " + bq.dequeue());
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, "consumer");
producer.start();
consumer.start();
producer.join();
consumer.join();
System.out.println("Done. Final size: " + bq.size()); // 0
}
}
What to observe in the output#
The producer fills the queue to capacity 3 and then blocks. You'll see 3 "Produced" lines appear instantly, then silence until the consumer dequeues one (after 200ms sleep), which wakes the producer to add the next item. The interleaving makes the blocking behaviour visible.
Queue size: 3
Queue size after dequeue: 2
---- concurrent demo ----
Produced: 1 (size=1)
Produced: 2 (size=2)
Produced: 3 (size=3)
Consumed: 1
Produced: 4 (size=3)
Consumed: 2
Produced: 5 (size=3)
Consumed: 3
Produced: 6 (size=3)
Consumed: 4
Consumed: 5
Consumed: 6
Done. Final size: 0
Common Mistakes#
| Mistake | Why it breaks |
|---|---|
| Using if instead of while around await() | Spurious wakeups cause a thread to proceed when the condition still doesn't hold |
| One condition variable for both states | Can't wake only producers or only consumers — leads to deadlock or wasted wake-ups |
| signal() before releasing the lock | Not an issue with Condition.signal() (it's legal to signal before unlock), but the signalled thread won't actually run until the lock is released |
| Forgetting lock.unlock() in finally | Lock is never released on exception — every other thread waits forever |
| Using notify() instead of signal() | notify() is tied to synchronized; mixing with ReentrantLock throws IllegalMonitorStateException |