Skip to content
BytePatterns

Producer-Consumer With a Bounded Buffer: Conditions and Sentinels

8 min readBytePatterns

Why the queue between producers and consumers must be bounded, how to build one from a lock and two condition variables, and how to shut consumers down.

Producer-consumer is the concurrency question that looks easy until you write it. One side makes work, the other side does it, and a queue sits in between. The queue is the whole design: its capacity decides what happens when one side is faster, its locking decides whether items get lost, and its shutdown rule decides whether your program ever exits.

The problem it solves

A web scraper downloads pages faster than it can parse them. A log shipper reads lines faster than the network accepts them. An image service receives uploads faster than it can resize them. In each case two stages run at different speeds, and you want them to run at the same time without either one knowing about the other.

A queue decouples them. The producer puts items in and moves on; the consumer takes items out when it is ready. Adding consumers scales the slow stage without touching the fast one.

The intuition

The important word in the title is bounded. An unbounded queue hides a speed mismatch until memory runs out: if the producer is 10% faster, the backlog grows forever. A bounded queue turns the mismatch into a signal. When the queue is full, put blocks, and the producer slows to the consumer's pace. That is backpressure, and it is a feature, not an error.

The mirror case is just as important. When the queue is empty, get blocks, and a consumer waits without burning CPU in a polling loop.

Building that behaviour takes one lock and two condition variables that share it. A condition variable lets a thread sleep until another thread says something may have changed:

  • Producers wait on not full; consumers signal it after taking an item.
  • Consumers wait on not empty; producers signal it after adding an item.

Each waiter re-checks its condition in a while loop after waking up, because by the time it runs, another thread may already have taken the slot or the item it was woken for.

Watch it run

The animation is the lesson's conveyor-belt sushi bar with a belt of two plates. The chef plates tuna and egg, the belt is full, and his next put blocks. A diner lifts the tuna, which frees one cell and unblocks the chef at once. When the belt runs empty, the diner's get blocks instead, and the chef ends the evening by putting the sentinel None, which the diner reads as "no more work" and exits.

Producer and Consumer

Step 1 of 11

A queue of size 2 sits between the two sides. Neither knows the other exists.

The same interactive animation as the lesson — step through it with the controls.

The code

A bounded buffer from first principles, in the shape interviewers usually ask for. peak records the fullest the buffer ever got, so the test can prove the bound held:

import threading
from collections import deque

class BoundedBuffer:
    def __init__(self, capacity):
        self.items = deque()
        self.capacity = capacity
        self.lock = threading.Lock()
        self.not_full = threading.Condition(self.lock)    # producers wait here
        self.not_empty = threading.Condition(self.lock)   # consumers wait here
        self.peak = 0

    def put(self, item):
        with self.not_full:
            while len(self.items) == self.capacity:   # while, not if: re-check on wake
                self.not_full.wait()
            self.items.append(item)
            self.peak = max(self.peak, len(self.items))
            self.not_empty.notify()                   # one item: wake one consumer

    def get(self):
        with self.not_empty:
            while not self.items:
                self.not_empty.wait()
            item = self.items.popleft()
            self.not_full.notify()                    # one free slot: wake one producer
            return item

Shutting down is the part most answers forget. A consumer blocked in get only wakes when something arrives, so after all producers finish, put one sentinel per consumer. A unique object() is safer than None if None could ever be real data:

STOP = object()                                       # sentinel: no more work

def run(capacity, producers, consumers, per_producer):
    buf, eaten, eaten_lock = BoundedBuffer(capacity), [], threading.Lock()

    def produce(p):
        for i in range(per_producer):
            buf.put((p, i))

    def consume():
        while True:
            item = buf.get()
            if item is STOP:
                return                                # one sentinel stops one consumer
            with eaten_lock:
                eaten.append(item)

    ps = [threading.Thread(target=produce, args=(p,)) for p in range(producers)]
    cs = [threading.Thread(target=consume) for _ in range(consumers)]
    for t in ps + cs:
        t.start()
    for t in ps:
        t.join()                                      # every real item is in the buffer
    for _ in cs:
        buf.put(STOP)                                 # then one sentinel per consumer
    for t in cs:
        t.join()
    return eaten, buf.peak

eaten, peak = run(capacity=2, producers=3, consumers=2, per_producer=100)
print(len(eaten), len(set(eaten)), peak <= 2)         # 300 300 True

In real Python code, queue.Queue(maxsize=...) already is this buffer. Its put also accepts a timeout, which is how a producer can choose to give up instead of waiting forever:

import queue

q = queue.Queue(maxsize=2)
q.put("tuna"); q.put("egg")
try:
    q.put("eel", timeout=0.1)                         # full: wait at most 0.1 s
except queue.Full:
    print("full: the producer must wait, drop, or shed load")
# full: the producer must wait, drop, or shed load

Thread timing changes from run to run, so the check does not compare interleavings. It compares outcomes against what must be true whatever the timing: across 60 random mixes of capacity, producer count, consumer count and item count, every produced item is consumed exactly once, and the buffer never held more than its capacity:

import random

random.seed(3)
ok = True
for _ in range(60):
    cap, np_, nc, per = (random.randint(1, 4), random.randint(1, 4),
                         random.randint(1, 4), random.randint(0, 40))
    eaten, peak = run(cap, np_, nc, per)
    want = sorted((p, i) for p in range(np_) for i in range(per))   # what was produced
    ok &= sorted(eaten) == want and peak <= cap
print(ok)                                             # True

The complexity

Each put and get is O(1) work under the lock: one append or pop on a deque and one notify. Memory is O(capacity), which is the point — it is fixed no matter how far the producer runs ahead. Throughput is set by the slower side; the buffer only absorbs short bursts. If the consumer is slower on average, no capacity is large enough, and the fix is more consumers or less input.

Where it goes wrong

  • if instead of while around wait. A woken thread may find the condition false again because another thread got there first; with if, a consumer pops from an empty deque.
  • One condition for both sides. With a single condition, notify can wake a thread of the wrong kind — a producer waking another producer — and the thread that could make progress stays asleep.
  • No shutdown signal. Consumers blocked in get never exit, and join hangs. Send one sentinel per consumer, after every producer has finished. Python 3.13 also added Queue.shutdown for this.
  • Unbounded queues under load. They trade a visible slowdown for an invisible memory leak.

How to say it in an interview

"I put a bounded queue between the two sides. When it's full, put blocks, which gives backpressure; when it's empty, get blocks, so consumers don't spin. Internally it's a lock with two condition variables, not-full and not-empty, and every wait is in a while loop because the condition can change between the notify and the wake-up. For shutdown, after the producers finish I send one sentinel per consumer. In Python I'd use queue.Queue with a maxsize rather than write it myself."

The lock underneath is covered in locks and mutexes, and what to do when the producer must not simply block is the subject of backpressure.