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.

Producer 1socket readerProducer 2file tailerProducer N...BlockingQueue (bounded)capacity = burst bufferConsumer 1drainTo batchConsumer 2drainTo batchConsumer M...put()take()full queue: put() blocks, producers slow down (backpressure)empty queue: take() blocks, consumers sleep without spinningDatabase / APIthe real bottleneck
Producers and consumers share only the queue. The bounded capacity both absorbs bursts and limits memory; when it is full, producers wait.

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:

BehaviourInsertRemoveInspect
Throw an exceptionadd(e)remove()element()
Return a special valueoffer(e) returns falsepoll() returns nullpeek()
Block until possibleput(e)take()none
Block with a timeoutoffer(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:

ImplementationBoundedInsideUse when
ArrayBlockingQueueYes, fixed at constructionCircular array, one lock, two conditions (not empty, not full); optional fairnessDefault choice for a bounded pipeline; no allocation per item
LinkedBlockingQueueOptional; default capacity Integer.MAX_VALUELinked nodes, separate put lock and take lock, atomic countProducers and consumers both busy; always pass a capacity
LinkedBlockingDequeOptionalLinked nodes, one lockYou need both ends, for example work you may push back
SynchronousQueueNo storage at allDirect handoff; each put waits for a takeHanding work straight to an idle thread
LinkedTransferQueueNo (unbounded)Non-blocking linked structure; transfer() waits for a consumerHigh-throughput handoff where producers may wait for receipt
PriorityBlockingQueueNo (unbounded)Binary heap, one lockConsumers must take the most urgent item first
DelayQueueNo (unbounded)Priority heap ordered by delayItems 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 LinkedBlockingQueue the queue never refuses, so maximumPoolSize is never used and tasks pile up without limit. Executors.newFixedThreadPool is built this way.
  • With a SynchronousQueue every task needs a free thread immediately, so the pool grows to maximumPoolSize. Executors.newCachedThreadPool uses this with an effectively unlimited maximum, which can create thousands of threads under load.
  • With a bounded ArrayBlockingQueue and ThreadPoolExecutor.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 the Future and 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. Use offer() 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

  1. Search your codebase for new LinkedBlockingQueue<>(), newFixedThreadPool and newCachedThreadPool and give each one an explicit bound or a justification.
  2. 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.
  3. Export queue depth and rejected or dropped item counts as metrics, and alert when depth stays near capacity.
  4. Implement shutdown with one poison pill per consumer, and test it: the process must exit within a bounded time with no items lost.
  5. Replace per-item writes with drainTo batching where the sink supports batch operations, and measure the difference.
  6. Choose a rejection policy deliberately; use CallerRunsPolicy where slowing the submitter is acceptable and an explicit error where it is not.
Key takeaway: A BlockingQueue decouples producers from consumers, absorbs bursts and, when bounded, slows producers instead of exhausting memory. Always set a capacity sized from arrival rate, service time and burst length, batch with drainTo, shut down with one poison pill per consumer, preserve interrupts, and choose the ThreadPoolExecutor queue and rejection policy on purpose.