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
ifinstead ofwhilearoundwait. A woken thread may find the condition false again because another thread got there first; withif, a consumer pops from an empty deque.- One condition for both sides. With a single condition,
notifycan 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
getnever exit, andjoinhangs. Send one sentinel per consumer, after every producer has finished. Python 3.13 also addedQueue.shutdownfor 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.