The producer/consumer pattern splits work into two sides that run at different speeds. Producers create work items: requests, log lines, file chunks, messages from a socket. Consumers process them: write to a database, call an API, compress, index. Between them sits a queue. The queue decouples the two sides in time, absorbs short bursts, and, if it is bounded, pushes back on producers when consumers cannot keep up.
In Java, that queue is almost always a java.util.concurrent.BlockingQueue. This article explains the interface from first principles, how the main implementations differ inside, how to shut a pipeline down cleanly, how to size the queue and the consumer pool with a worked example, and how the same ideas appear inside ThreadPoolExecutor. If you want to see how a bounded buffer is built from a mutex and condition variables, read the condition variable article first; the queues here are production-grade versions of that design.
Why put a queue between them
Without a queue, a producer has to call the consumer directly and wait. The slowest step sets the speed of the whole system, and a 200 millisecond pause in the database stalls the socket reader, which then drops connections. A queue lets each side run at its own pace, with three properties that matter:
- Decoupling. Producers and consumers do not know about each other; each only knows the queue. You can change the number of consumers without touching producer code.
- Burst absorption. If producers briefly outrun consumers, items accumulate in the queue and drain later. The queue's capacity is the size of the burst you can absorb.
- Backpressure. If producers outrun consumers for long, a bounded queue fills and
put()blocks, slowing producers to the consumers' speed. An unbounded queue instead grows until the heap is exhausted.
The last point is the one that causes outages. A bounded queue turns sustained overload into slower producers, which is visible and recoverable. An unbounded queue turns it into an OutOfMemoryError minutes later, which is neither.
The BlockingQueue contract
The interface offers four ways to insert and four ways to remove, and the right one depends on what should happen when the queue is full or empty:
| Behaviour | Insert | Remove | Inspect |
|---|---|---|---|
| Throw an exception | add(e) | remove() | element() |
| Return a special value | offer(e) returns false | poll() returns null | peek() |
| Block until possible | put(e) | take() | none |
| Block with a timeout | offer(e, time, unit) | poll(time, unit) | none |
Consumers in a dedicated thread normally use take(), or poll(timeout) if they also need to wake up periodically to flush a partial batch. Producers on a thread you can afford to block use put(). Producers on a thread you must never block, such as a network event loop, use offer() and decide what to do on failure: drop, count, return an error to the caller. Blocking queues do not accept null elements, because null is the sentinel poll() uses for "nothing there".
Two more methods matter in practice. drainTo(collection, max) moves up to max available items in one call, taking the lock once instead of once per item. remainingCapacity() reports free space, which is useful for metrics but not for decisions, because it is stale the moment it returns.
Implementations and how they lock
The implementations differ in capacity, ordering and, most importantly for performance, locking:
| Implementation | Bounded | Inside | Use when |
|---|---|---|---|
ArrayBlockingQueue | Yes, fixed at construction | Circular array, one lock, two conditions (not empty, not full); optional fairness | Default choice for a bounded pipeline; no allocation per item |
LinkedBlockingQueue | Optional; default capacity Integer.MAX_VALUE | Linked nodes, separate put lock and take lock, atomic count | Producers and consumers both busy; always pass a capacity |
LinkedBlockingDeque | Optional | Linked nodes, one lock | You need both ends, for example work you may push back |
SynchronousQueue | No storage at all | Direct handoff; each put waits for a take | Handing work straight to an idle thread |
LinkedTransferQueue | No (unbounded) | Non-blocking linked structure; transfer() waits for a consumer | High-throughput handoff where producers may wait for receipt |
PriorityBlockingQueue | No (unbounded) | Binary heap, one lock | Consumers must take the most urgent item first |
DelayQueue | No (unbounded) | Priority heap ordered by delay | Items become available only after a delay, such as retries |
The lock design explains the performance difference. ArrayBlockingQueue uses a single ReentrantLock for both ends, so a producer and a consumer contend with each other. LinkedBlockingQueue uses two locks, so a put and a take can proceed at the same time; the price is allocating a node for every item, which adds garbage collection work. Neither is universally faster. Measure with your item sizes and thread counts. When the queue's own lock is genuinely the bottleneck, at millions of items per second, the answer is usually a different design, such as the ring buffer described in the Disruptor article.
Note the trap in the table: three of the implementations are unbounded and LinkedBlockingQueue is effectively unbounded unless you pass a capacity. new LinkedBlockingQueue<>() is the most common source of the out-of-memory failure described above.
A complete pipeline with clean shutdown
Here is a complete pipeline: one producer reading lines, a bounded queue, and a pool of consumers that write in batches. It shows the three details that most examples omit: batching with drainTo, interrupt handling, and clean shutdown with one poison pill per consumer.
import java.util.*;
import java.util.concurrent.*;
public final class IngestPipeline {
private static final String POISON = new String("POISON"); // unique identity
private final BlockingQueue<String> queue = new ArrayBlockingQueue<>(10_000);
private final int consumers;
private final ExecutorService pool;
IngestPipeline(int consumers) {
this.consumers = consumers;
this.pool = Executors.newFixedThreadPool(consumers);
}
void start(Sink sink) {
for (int i = 0; i < consumers; i++) {
pool.submit(() -> consume(sink));
}
}
/** Called by the producer thread; blocks when consumers fall behind. */
void publish(String line) throws InterruptedException {
queue.put(line);
}
/** Called once by the producer when input is exhausted. */
void finish() throws InterruptedException {
for (int i = 0; i < consumers; i++) queue.put(POISON);
pool.shutdown();
if (!pool.awaitTermination(60, TimeUnit.SECONDS)) pool.shutdownNow();
}
private void consume(Sink sink) {
List<String> batch = new ArrayList<>(500);
try {
while (true) {
String first = queue.take(); // wait for work
if (first == POISON) break;
batch.add(first);
queue.drainTo(batch, 499); // grab what is ready
int pills = 0;
for (String s : batch) if (s == POISON) pills++;
if (pills > 0) { // drained pill(s): ours plus peers'
batch.removeIf(s -> s == POISON);
sink.write(batch);
for (int i = 1; i < pills; i++) queue.put(POISON); // return peers' pills
return;
}
sink.write(batch);
batch.clear();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt(); // preserve the signal
}
}
interface Sink { void write(List<String> batch); }
}A few choices are deliberate. The poison pill compares by identity (==) with a freshly allocated string, so a real input line that happens to read "POISON" cannot stop a consumer. Because drainTo can scoop up pills meant for other consumers, the code keeps one and puts the rest back, so every consumer still receives exactly one. A pill is better than shutdownNow() for normal shutdown because it lets consumers finish everything queued before it; interrupts are for the case where you want to abandon work. Catching InterruptedException and re-setting the interrupt flag lets the executor see that the thread was asked to stop; swallowing it silently is a classic bug that makes shutdown hang.
Worked example: sizing consumers and capacity
Sizing comes from two numbers: how fast items arrive and how fast one consumer can process them. Little's law says the average number of items in a stable system equals the arrival rate times the time each item spends in it. Applied to the consumer side, the number of busy consumers you need is the arrival rate times the service time per item.
Suppose a producer receives log events at a steady 20,000 per second, with bursts to 60,000 per second lasting up to 2 seconds. Each consumer writes batches of 500 events, and one batch insert takes 40 milliseconds, so one consumer handles 500 / 0.040 = 12,500 events per second. Two consumers give 25,000 per second, which covers the steady rate with 25 percent headroom. Use three: if the database slows down by half during compaction, three consumers still give 18,750 per second, close to the arrival rate, and the queue absorbs the difference for a while.
During a burst, items arrive at 60,000 per second and leave at up to 37,500 per second with three consumers, so the queue grows by 22,500 per second, or 45,000 items over two seconds. A capacity of 50,000 absorbs the burst without blocking producers. At an average of 300 bytes per event plus object overhead, roughly 30 to 40 MB of heap, which is affordable. If the burst lasts longer, the queue fills and the producer blocks, which is exactly the behaviour you want: the overload is pushed back towards the source instead of into the heap.
Do not make the queue much larger than the burst you need to absorb. A bigger queue does not increase throughput; it only increases latency, because each item waits behind everything ahead of it, and it increases the amount of work lost if the process crashes. Monitor queue depth: a depth that keeps hitting capacity means you need more consumers or a faster sink, not a bigger queue.
The queue inside ThreadPoolExecutor
ThreadPoolExecutor is a producer/consumer system with the queue built in, and its behaviour surprises people because the queue choice changes how threads are created. On execute(), the executor starts a new thread if fewer than corePoolSize are running; otherwise it offers the task to the queue; only if the queue refuses does it start threads up to maximumPoolSize; and if that also fails it calls the rejection handler. Pool sizing more generally is covered in the thread pools article.
- With an unbounded
LinkedBlockingQueuethe queue never refuses, somaximumPoolSizeis never used and tasks pile up without limit.Executors.newFixedThreadPoolis built this way. - With a
SynchronousQueueevery task needs a free thread immediately, so the pool grows tomaximumPoolSize.Executors.newCachedThreadPooluses this with an effectively unlimited maximum, which can create thousands of threads under load. - With a bounded
ArrayBlockingQueueandThreadPoolExecutor.CallerRunsPolicy, a full pool makes the submitting thread run the task itself. That slows the submitter down to the pool's pace, which is a simple and effective form of backpressure.
ThreadPoolExecutor exec = new ThreadPoolExecutor(
8, 8, // core == max: a fixed pool
0L, TimeUnit.MILLISECONDS,
new ArrayBlockingQueue<>(1_000), // bounded: overload becomes visible
new ThreadPoolExecutor.CallerRunsPolicy());
Failure modes
The failures are well known, which is why they are worth listing:
- Unbounded queue under sustained overload. Memory grows until the JVM spends all its time in garbage collection and then fails. Always set a capacity.
- Missing poison pills. One pill for three consumers stops one of them; the other two block forever in
take()and the process never exits. - Swallowed interrupts. An empty
catch (InterruptedException e) {}loses the shutdown request and the loop carries on. - Consumer exceptions. An exception thrown from a task submitted with
submit()is captured in theFutureand never logged. The consumer loop dies quietly and the queue fills. Catch and log per item inside the loop, or check the futures. - Blocking an event loop. Calling
put()from a Netty or other event-loop thread stalls every connection on that loop. Useoffer()there. - Ordering assumptions. With more than one consumer, items finish out of order even though the queue is FIFO. If order matters per key, route each key to a fixed consumer with its own queue.
- Lost work on crash. An in-memory queue is lost when the process dies. If items must survive, use a durable log or broker and treat the in-memory queue as a buffer only.
Trade-offs and alternatives
Bounded blocking queues trade some latency and a little throughput for predictable memory and natural backpressure, which is almost always the right trade inside one process. For cross-process pipelines, a broker such as Kafka replaces the queue, and flow control becomes consumer lag and producer quotas. For asynchronous code that must never block a thread, reactive streams make the consumer request a number of items explicitly; see reactive streams backpressure.
Virtual threads, final in Java 21, change the cost model but not the pattern. Blocking in put() or take() on a virtual thread is cheap because the carrier thread is released while it waits, and the queues use java.util.concurrent locks rather than synchronized, so they do not pin the carrier. You can afford one virtual thread per producer, but you still need a bounded queue: cheap threads make it easier, not harder, to create more work than the sink can absorb. More detail is in virtual threads in production.
What to do next
- Search your codebase for
new LinkedBlockingQueue<>(),newFixedThreadPoolandnewCachedThreadPooland give each one an explicit bound or a justification. - For each pipeline, write down the arrival rate, the per-item service time and the largest burst, and compute consumer count and capacity with Little's law.
- Export queue depth and rejected or dropped item counts as metrics, and alert when depth stays near capacity.
- Implement shutdown with one poison pill per consumer, and test it: the process must exit within a bounded time with no items lost.
- Replace per-item writes with
drainTobatching where the sink supports batch operations, and measure the difference. - Choose a rejection policy deliberately; use
CallerRunsPolicywhere slowing the submitter is acceptable and an explicit error where it is not.