BlockingQueue<E> is the hand-off point between threads that produce work and threads that consume it. It is a thread-safe queue with one extra promise: a consumer that asks for an element when the queue is empty can wait, and a producer that offers an element to a full bounded queue can wait too. That second wait is the important one. It is back-pressure: when consumers fall behind, producers slow down automatically instead of filling the heap.
This article explains the interface contract, how the main implementations work inside, how to size a bounded queue with arithmetic rather than guesswork, and how to shut a pipeline down without losing or duplicating work. It finishes with a complete batching pipeline, the failure modes that show up in production and a checklist. Everything here applies to Java 8 onwards; the notes on virtual threads need JDK 21 or later.
The contract: four ways to insert and remove
Every operation comes in four forms that differ only in what happens when the operation cannot proceed immediately. Picking the wrong form is the most common BlockingQueue bug, so learn this table first:
| Operation | Throws | Returns a value | Blocks | Times out |
|---|---|---|---|---|
| Insert | add(e) | offer(e) returns false | put(e) | offer(e, time, unit) |
| Remove | remove() | poll() returns null | take() | poll(time, unit) |
| Examine | element() | peek() | n/a | n/a |
Three more rules are part of the contract. First, null elements are rejected with NullPointerException, because poll() uses null to mean empty. Second, the blocking and timed methods throw InterruptedException, which is how another thread asks a waiting thread to stop. Third, the queue gives you a memory-visibility guarantee: everything a producer did before inserting an element happens-before everything a consumer does after removing it. You can build a request object with plain fields, put it, and the consumer will see every field fully written, with no volatile and no extra lock.
Use the blocking pair put/take for steady pipelines, the timed forms when a thread must also notice shutdown or deadlines, and offer/poll when the caller has something better to do than wait, such as returning HTTP 503. The throwing forms rarely fit concurrent code.
Inside the implementations: locks and conditions
Knowing the internals explains the performance and the failure modes. ArrayBlockingQueue is a circular array with a fixed capacity chosen at construction, guarded by one ReentrantLock with two conditions, notEmpty and notFull. A put takes the lock, waits on notFull while the array is full, stores the element, and signals notEmpty. Because there is a single lock, a producer and a consumer can never be inside the queue at the same moment. An optional fairness flag, new ArrayBlockingQueue<>(cap, true), grants the lock to waiting threads in roughly FIFO order, which removes starvation at a real throughput cost.
LinkedBlockingQueue is a linked list with two locks, one for the tail (putLock) and one for the head (takeLock), plus an atomic element count. Producers and consumers work on different ends and only meet when one side has to wake the other. That is why it usually sustains higher throughput under mixed load. The price is one node allocation per element, which adds garbage-collection work, and a default capacity of Integer.MAX_VALUE, which is effectively unbounded. Always pass a capacity.
Here is a teaching version of the single-lock design. It shows the essential shape, including the while loop that guards against spurious wake-ups:
final class TinyBoundedQueue<E> {
private final Object[] items;
private int putIndex, takeIndex, count;
private final ReentrantLock lock = new ReentrantLock();
private final Condition notEmpty = lock.newCondition();
private final Condition notFull = lock.newCondition();
TinyBoundedQueue(int capacity) { items = new Object[capacity]; }
void put(E e) throws InterruptedException {
Objects.requireNonNull(e);
lock.lockInterruptibly();
try {
while (count == items.length) notFull.await(); // re-check after every wake-up
items[putIndex] = e;
putIndex = (putIndex + 1) % items.length;
count++;
notEmpty.signal();
} finally { lock.unlock(); }
}
@SuppressWarnings("unchecked")
E take() throws InterruptedException {
lock.lockInterruptibly();
try {
while (count == 0) notEmpty.await();
E e = (E) items[takeIndex];
items[takeIndex] = null; // let the GC reclaim it
takeIndex = (takeIndex + 1) % items.length;
count--;
notFull.signal();
return e;
} finally { lock.unlock(); }
}
}
Choosing an implementation
The JDK ships seven implementations. They differ in whether they are bounded, how they order elements and whether a producer can ever block:
| Implementation | Bounded? | Order | Use it for |
|---|---|---|---|
ArrayBlockingQueue | always | FIFO | a fixed buffer with predictable memory |
LinkedBlockingQueue | optional (bound it) | FIFO | high-throughput producer/consumer pipelines |
LinkedBlockingDeque | optional | both ends | work that is sometimes re-queued at the front |
PriorityBlockingQueue | no, put never blocks | comparator | priority scheduling, with a separate admission limit |
DelayQueue | no | earliest expiry first | retry and timeout schedules; elements implement Delayed |
SynchronousQueue | zero capacity | direct hand-off | passing work only when a receiver is ready |
LinkedTransferQueue | no | FIFO | transfer(), which waits until a consumer has received the element |
Two cautions. The unbounded variants never push back on producers, so they need a limit somewhere else, for example a Semaphore around submission. And the iterator of a PriorityBlockingQueue is not in priority order; only removal is.
Sizing a bounded queue with Little&amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;amp;#x27;s law
A bounded queue needs a number, and Little's law gives you one. In a stable system, the average number of items in the queue equals the arrival rate multiplied by the average time each item waits. Turn it around: decide how long an item may wait before the result is useless, and the capacity follows.
Worked example. An ingestion service receives a peak of 4,000 events per second. Four consumer threads each process 1,200 events per second, so total capacity is 4,800 per second and the system keeps up on average. Bursts are the problem: traffic arrives in clumps, and the downstream contract says an event must be handled within 500 ms. The queue drains at the consumers' rate of 4,800 per second, so the last element in a full queue waits capacity / 4,800 seconds. Capacity is therefore about 4,800 × 0.5 = 2,400 elements. Larger than that, and items sit in memory past their deadline. Smaller, and harmless bursts trigger back-pressure. Check the memory too: at about 1 KB per event, 2,400 events is about 2.4 MB, which is nothing. If they hold 1 MB payloads, store references to off-heap buffers or reduce the bound.
Worked example: a batching log shipper
Here is a complete pipeline that ships log records to a remote store in batches. Producers are request threads that must never block for long. Consumers batch records to amortise network calls and flush at least every 200 ms. Shutdown drains everything already accepted.
public final class LogShipper implements AutoCloseable {
private static final LogRecord POISON = new LogRecord("__poison__");
private final BlockingQueue<LogRecord> queue = new ArrayBlockingQueue<>(2_000);
private final List<Thread> workers = new ArrayList<>();
private final AtomicLong dropped = new AtomicLong();
private volatile boolean closed;
private final Sink sink;
public LogShipper(Sink sink, int consumers) {
this.sink = sink;
for (int i = 0; i < consumers; i++) {
Thread t = new Thread(this::consume, "log-shipper-" + i);
t.start();
workers.add(t);
}
}
/** Called on request threads: wait at most 5 ms, then shed load. */
public boolean submit(LogRecord r) {
if (closed) return false;
try {
if (queue.offer(r, 5, TimeUnit.MILLISECONDS)) return true;
} catch (InterruptedException e) {
Thread.currentThread().interrupt(); // never swallow the interrupt
}
dropped.incrementAndGet(); // export this as a metric
return false;
}
private void consume() {
List<LogRecord> batch = new ArrayList<>(500);
try {
while (true) {
LogRecord first = queue.poll(200, TimeUnit.MILLISECONDS);
if (first != null) {
batch.add(first);
queue.drainTo(batch, 499); // grab what is ready, no waiting
}
long pills = batch.stream().filter(x -> x == POISON).count();
batch.removeIf(x -> x == POISON); // identity, never equals()
if (!batch.isEmpty()) { sink.write(batch); batch.clear(); }
if (pills > 0) { // drainTo may grab other workers' pills
for (long i = 1; i < pills; i++) queue.put(POISON);
return;
}
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt(); // forced stop: exit promptly
}
}
@Override public void close() throws InterruptedException {
closed = true;
for (int i = 0; i < workers.size(); i++) queue.put(POISON); // one pill per consumer
for (Thread t : workers) t.join(10_000);
}
}Walk through the data flow. A request thread calls submit. If the queue has room, the record is stored and the call returns in microseconds. If the queue is full for 5 ms, the record is counted as dropped and the request carries on, because losing a log line is better than stalling a user request. A consumer waits up to 200 ms for the first record, then uses drainTo to take up to 499 more that are already waiting, so a busy queue produces full batches and a quiet one still flushes on time. Note that drainTo is not an atomic snapshot: it moves what it finds, and other consumers may take elements concurrently. On close, one poison pill per consumer goes in behind the accepted records, so every record before the pill is written before its consumer exits. Because drainTo can pick up several pills in one sweep, a consumer that finds more than one puts the extras back for its siblings.
Shutdown: poison pills and interruption
There are two ways to stop consumers, and they mean different things. A poison pill is a graceful stop: it travels through the queue like any other element, so all earlier work completes. Interruption is a forced stop: a thread blocked in take() or put() throws InterruptedException immediately. Use the pill for normal shutdown and interruption as the fallback when the pill does not finish within a deadline, which is exactly what ExecutorService.shutdownNow() does to pooled consumers.
When you catch InterruptedException, either rethrow it or restore the flag with Thread.currentThread().interrupt(). Swallowing it is the classic bug: the thread forgets it was asked to stop, loops back into take() and keeps the JVM alive at shutdown. Decide too what happens to elements still in the queue after a forced stop. Either persist them, or accept and document that they are lost. A queue is memory, not a durable log; if work must survive a crash, it belongs in a broker or a database before you acknowledge it.
Failure modes
These are the failures that actually reach production:
- The unbounded backlog. A
LinkedBlockingQueuebuilt without a capacity, often hidden insideExecutors.newFixedThreadPool, grows until the heap is exhausted. Long GC pauses come first, thenOutOfMemoryError. The ExecutorService article shows how to build pools with bounded queues. - Wrong pill count. One pill and four consumers leaves three threads blocked forever. Send one per consumer, or have each consumer put the pill back before exiting.
- Deciding with size().
if (queue.size() < cap) queue.put(x)is a race: the answer is stale before you act on it. Use the boolean result ofofferor the value ofpoll, which reflect the queue at the moment of the operation.remainingCapacity()has the same limitation, and on an unbounded queue it always returnsInteger.MAX_VALUE. - Head-of-line blocking. One slow element holds a consumer while fast ones wait behind it. Put deadlines on processing, or split slow and fast work into separate queues with separate consumers.
- Blocking a thread that must not block. Calling
putfrom an event-loop thread, such as Netty's, stalls every connection on that loop. Useofferthere.
Operating queues in production
Export three numbers per queue: depth (size() is fine for a gauge), the rejection or drop count, and the time each element waited. Waiting time is the most useful and the one queues do not give you: stamp each element with System.nanoTime() when it is enqueued and record the difference when it is taken. A queue that is always near full means consumers are under-provisioned; one that is always empty with busy producers means the queue is the wrong place to look.
Virtual threads, on JDK 21 or later, combine well with blocking queues. The JDK queues are built on java.util.concurrent locks, so a virtual thread that blocks in take() unmounts and frees its carrier thread. That makes one virtual thread per connection, each feeding a shared bounded queue, a cheap design. See virtual threads for how mounting works.
Trade-offs and alternatives
A BlockingQueue is the right default for moving work between thread pools within one JVM. Its costs are lock hand-offs and, when a thread has to wait, a park and an unpark. For most services that is far below the cost of the work itself. The alternatives make sense in narrower cases:
| Option | Choose it when | Give up |
|---|---|---|
ConcurrentLinkedQueue | consumers poll and never need to wait | back-pressure and blocking |
| LMAX Disruptor, JCTools queues | millions of events per second, latency-critical | simplicity; often spinning CPU |
SubmissionPublisher (Flow API) | reactive-style demand signalling | the plain put/take model |
| Kafka, SQS or another broker | work must survive a crash or cross processes | in-memory speed |
Measure before switching. Most queue bottlenecks are slow consumers, not slow queues. For the wider picture of Java's concurrency toolbox, see Java concurrency in depth and ReentrantLock, the lock these queues are built on.
What to do next
Use this list to review the queues in your own code:
- Search the codebase for
new LinkedBlockingQueue<>()andnewFixedThreadPool, and give every queue an explicit capacity. - Size each bound from your latency budget with Little's law, then check its memory footprint.
- Decide per producer whether it may block. Use
putfor internal pipelines and a timedofferplus a drop or 503 path for request threads. - Audit every
catch (InterruptedException e): rethrow it or restore the interrupt flag. - Implement shutdown as one poison pill per consumer, with an interrupt fallback after a deadline.
- Export depth, drops and per-element wait time, and alert on sustained high wait time.
- Load-test with consumers deliberately slowed, and confirm memory stays flat and drops appear.