ByteScrollGet the app
☰ Topics
producer-consumer-blockingqueue5 / 200‹›
JAVA / CONCURRENCY3 minute read

Producer-consumer with BlockingQueue

Medium

Producers 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

  1. Shared buffer. Producers and consumers only touch the queue. They never share locks or call wait/notify themselves.
  2. Blocking calls. put(e) blocks while the queue is full; take() blocks while it's empty. Both throw InterruptedException.
  3. Back-pressure. A bounded queue makes fast producers slow down to the consumers' pace instead of piling up memory.
  4. 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)
  5. 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) and DelayQueue (items become available after a delay).
put()put()take()take()full: put() waitsempty: take() waitsProducer 1Bounded queue,capacity 100Producer 2Consumer 1Consumer 2
put()put()take()take()full: put() waitsempty: take() waitsProducer 1Bounded queue,capacity 100Producer 2Consumer 1Consumer 2

Example

Example.javaJava
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

  • BlockingQueue rejects null with NullPointerException, partly because poll() uses null to 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 InterruptedException and ignoring it. Restore the flag with Thread.currentThread().interrupt() and exit.
  • Calling add() on a full bounded queue and getting IllegalStateException instead 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