Producer-consumer with BlockingQueue
MediumProducers put() items on a BlockingQueue and consumers take() them. put() waits while a bounded queue is full and take() waits while it's empty, so the queue handles all locking and signalling. Stop consumers with a poison pill or by interrupting them.
How it works
- Shared buffer. Producers and consumers only touch the queue. They never share locks or call
wait/notifythemselves. - Blocking calls.
put(e)blocks while the queue is full;take()blocks while it's empty. Both throwInterruptedException. - Back-pressure. A bounded queue makes fast producers slow down to the consumers' pace instead of piling up memory.
- Choosing the call. Each operation comes in four flavours:
- throw:
add/remove - return a value:
offer/poll - block:
put/take - wait with a timeout:
offer(e, t, unit)/poll(t, unit)
- throw:
- Implementations.
ArrayBlockingQueue: fixed capacity, one lock for both ends.LinkedBlockingQueue: separate locks for head and tail, so a put and a take can run at once. Unbounded (Integer.MAX_VALUE) unless you pass a capacity.SynchronousQueue: zero capacity; each put waits for a matching take.PriorityBlockingQueue(unbounded, ordered) andDelayQueue(items become available after a delay).
Example
record Order(int id) {}
static final Order DONE = new Order(-1); // poison pill
public static void main(String[] args) throws Exception {
BlockingQueue<Order> queue = new ArrayBlockingQueue<>(100);
int consumers = 2;
Thread producer = Thread.ofPlatform().start(() -> {
try {
for (int i = 1; i <= 1_000; i++) queue.put(new Order(i));
for (int i = 0; i < consumers; i++) queue.put(DONE); // one pill each
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
List<Thread> workers = new ArrayList<>();
for (int c = 0; c < consumers; c++) {
workers.add(Thread.ofPlatform().start(() -> {
try {
for (Order o = queue.take(); o != DONE; o = queue.take()) {
System.out.println(Thread.currentThread().getName() + " packed " + o.id());
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}));
}
producer.join();
for (Thread w : workers) w.join();
}Edge cases
BlockingQueuerejectsnullwithNullPointerException, partly becausepoll()usesnullto mean "nothing there".- With several consumers you need one poison pill per consumer, or the extras block forever.
ArrayBlockingQueue(cap, true)makes waiting threads proceed in FIFO order, at some throughput cost.drainTo(list, max)pulls a batch in one call, useful when consumers write to a database in chunks.
Common mistakes
- Using
new LinkedBlockingQueue<>()without a capacity, then running out of memory when consumers fall behind. - Catching
InterruptedExceptionand ignoring it. Restore the flag withThread.currentThread().interrupt()and exit. - Calling
add()on a full bounded queue and gettingIllegalStateExceptioninstead of waiting.
Likely follow-up
"How does ExecutorService relate to this?" A thread pool is producer-consumer: submit() puts tasks on the pool's BlockingQueue and worker threads take() them. That's why a fixed pool with an unbounded queue can hide an ever-growing backlog.
Get every deep dive in the app
Coming soon to the App StoreComing soon to Google Play