Linux System Programming · advanced · ~14 min

The producer/consumer pattern

- By the end you can build a bounded, thread-safe queue from a circular buffer, one mutex, and two condition variables. - By the end you can explain why a condition variable must always be waited on inside a `while` loop that re-checks a predicate. - By the end you can pair each signal with the right condition variable so producers and consumers never wake the wrong side. - By the end you can shut the system down cleanly, releasing blocked consumers without losing or duplicating items. - By the end you can choose correctly between `pthread_cond_signal` and `pthread_cond_broadcast`.

Overview

You already know from mutex-basics how to make a mutex guard shared state so only one thread mutates it at a time, and from deadlocks how threads can freeze forever when they wait on each other in the wrong order. The producer/consumer pattern combines both ideas into the single most common multi-threaded design in systems programming: some threads make work, other threads do work, and a shared queue in the middle lets them run at their own pace.

The new ingredient here is the condition variable — a primitive that lets a thread sleep until another thread tells it the state it cares about has changed, instead of burning CPU checking in a loop. A mutex answers "is it safe to touch the queue right now?"; a condition variable answers "is there anything worth touching yet?". This lesson shows how the two work together, and how the same wait-ordering discipline you learned for avoiding deadlocks keeps this pattern correct.

Why it matters

Almost every server, pipeline, and background-job system is a producer/consumer at heart: a network thread accepts requests and hands them to a pool of workers; a logging thread drains a queue of messages other threads produced; a video decoder fills a frame buffer that the display thread empties. Get the synchronization wrong and you get the nastiest class of bugs there is — races that corrupt the queue, lost wakeups that hang the whole program, and consumers that read from an empty buffer and hand attacker-influenced garbage to the rest of the system. Because these bugs are timing-dependent they often pass every test on your laptop and only surface under load in production, which is exactly when they are most expensive and most dangerous.

Core concepts

The four pieces

A correct bounded producer/consumer needs exactly four things working together:

  1. A buffer with a fixed capacity to hold items. A circular buffer (a fixed array plus head, tail, and count) is the usual choice: it never allocates, never grows, and gives you natural back-pressure when it fills.
  2. A mutex that guards every field of the queue. head, tail, count, and the array cells are all shared state; touching any of them without the lock is a data race and therefore undefined behaviour.
  3. A condition variable not_empty. A consumer that finds the queue empty waits here. A producer signals it after adding an item.
  4. A condition variable not_full. A producer that finds the queue full waits here. A consumer signals it after removing an item.

Two conditions, two condition variables. Mixing them up is the classic bug — you signal the side that isn't waiting, and the side that is waiting sleeps forever.

How a condition variable actually works

The magic of pthread_cond_wait(&cv, &m) is that it does three things atomically: it unlocks the mutex, puts the thread to sleep, and — when woken — re-acquires the mutex before returning. This atomic "unlock-and-sleep" is why you must hold the mutex when you call it. If waiting simply meant "unlock, then sleep" as two separate steps, another thread could sneak in between them, change the state, signal the condition, and the signal would land on nobody — a lost wakeup that hangs your program.

Here is the timeline for a consumer arriving at an empty queue while a producer is about to add one item:

  consumer thread                 producer thread
  --------------                  --------------
  lock(m)
  count == 0?  yes
  cond_wait(not_empty, m) ......  (m is now free)
     |  sleeps, m released              lock(m)
     |                                  buf[tail]=x; count=1
     |                                  cond_signal(not_empty)
     |  <---- woken -------------       unlock(m)
  re-acquires m automatically
  count == 0?  no  -> take item
  cond_signal(not_full)
  unlock(m)

Always wait in a while loop, never an if

Wait on a condition variable inside a while loop that re-checks the predicate:

while (q->count == 0)
    pthread_cond_wait(&q->not_empty, &q->m);

There are two independent reasons a plain if is wrong:

  • Spurious wakeups. POSIX explicitly permits pthread_cond_wait to return without any matching signal. The standard allows this because it makes condition variables far cheaper to implement on real hardware. Your code must treat every wakeup as merely a hint to re-check, not a guarantee.
  • Stolen wakeups. With multiple consumers, a signal wakes one thread, but before it re-acquires the mutex a different consumer could run and take the item. The woken thread must re-check and find the queue empty again.

Either way the fix is identical: loop and re-test. The while makes the code correct against both.

Knowledge check: You replace the while around pthread_cond_wait(&not_empty, ...) with an if. There is exactly one producer and one consumer, so wakeups are never "stolen". Is the code now safe?

No. Even with a single producer and single consumer, POSIX still permits spurious wakeups: pthread_cond_wait may return with no signal ever sent. With an if, the consumer falls straight through to buf[head] while count == 0, reading a stale or never-written slot. The while is mandatory regardless of thread counts.

signal vs broadcast

pthread_cond_signal wakes at least one waiter; pthread_cond_broadcast wakes all of them.

Situation Use Why
Added exactly one item signal(not_empty) Only one consumer can take it; waking more just makes them re-check and sleep again.
Removed exactly one item signal(not_full) Symmetric: one slot freed, one producer can proceed.
Shutting the queue down broadcast(not_empty) Every blocked consumer must wake, see "closed", and exit. A single signal would leave the others asleep forever.
Changed state that satisfies several different predicates broadcast When you can't be sure which/how many waiters can now proceed, waking all and letting each re-check is always safe.

When in doubt, broadcast is never incorrect — only potentially slower (the "thundering herd" of woken threads all contending for the mutex). signal is the optimization you use when you can prove exactly one waiter can make progress.

Clean shutdown

A subtle real-world need: how do consumers know to stop? If producers simply finish, consumers blocked in pthread_cond_wait(&not_empty, ...) wait forever because no more items will ever arrive. The idiom is a closed flag: after every producer has joined, set closed = true and broadcast(not_empty). Each consumer's wait predicate becomes while (count == 0 && !closed), and after the loop it checks: if the queue is empty and closed, return "done". This drains every remaining item first, then lets each consumer exit exactly once.

Syntax notes

int pthread_cond_wait(pthread_cond_t *cv, pthread_mutex_t *m);

Must be called with m already locked. Atomically unlocks m and blocks on cv; on return m is locked again. Returns 0 on success. Always call inside while (!predicate).

int pthread_cond_signal(pthread_cond_t *cv);     // wake >= 1 waiter
int pthread_cond_broadcast(pthread_cond_t *cv);  // wake ALL waiters

Safe to call with or without the mutex held; holding it while signalling gives predictable scheduling and avoids a class of races, so prefer signalling before you unlock. If no thread is waiting, the signal is simply lost (harmless — this is why you re-check the predicate rather than rely on the signal).

int pthread_cond_init(pthread_cond_t *cv, const pthread_condattr_t *attr); // attr NULL = defaults
int pthread_cond_destroy(pthread_cond_t *cv);   // only when no thread is waiting on it

Static alternative: pthread_cond_t cv = PTHREAD_COND_INITIALIZER; (no destroy strictly required, but pairing init/destroy is good hygiene).

int pthread_mutex_lock(pthread_mutex_t *m);
int pthread_mutex_unlock(pthread_mutex_t *m);

Guard every access to shared queue fields. The same mutex passed to pthread_cond_wait must be the one protecting the predicate it tests.

All of these return an errno-style int (0 = success). In production you check it; examples often omit the check for brevity. Compile and link pthreads with cc -std=c11 file.c -o file -lpthread.

Lesson

The pattern

Producer/consumer is the classic multi-threaded design.

  • One or more producer threads add items to a queue.
  • One or more consumer threads remove items from that queue.

The queue is the meeting point where the threads synchronise.

Building blocks

You need four pieces:

  • A buffer to hold the items. This can be a circular buffer (a fixed-size array reused in a loop) or a linked list.
  • A mutex that protects the queue, so only one thread touches it at a time.
  • A condition variable cv_not_empty. Consumers wait on it when the queue is empty. A producer signals it with pthread_cond_signal after adding an item.
  • A condition variable cv_not_full. Producers wait on it when the queue is full. A consumer signals it after removing an item.

A condition variable lets a thread sleep until another thread tells it that something changed.

Always wait in a while loop

Wait on a condition variable inside a while loop that checks the predicate (the condition you are waiting for). Do not use a plain if.

The reason is spurious wakeups: POSIX explicitly allows pthread_cond_wait to return even when no signal was sent. The while loop re-checks the condition and goes back to sleep if it is not yet true.

while (count == 0)
    pthread_cond_wait(&cv_not_empty, &m);

Code examples

#include <pthread.h>
#include <stdio.h>
#include <stdlib.h>
#include <stdbool.h>

#define CAP 4          /* small on purpose: forces the "full" path to trigger */
#define ITEMS 12       /* total items each producer will make */
#define NPROD 2
#define NCONS 3

typedef struct {
    int      buf[CAP];
    int      head, tail, count;   /* circular-buffer indices + fill level */
    bool     closed;              /* producers set this when done */
    pthread_mutex_t m;
    pthread_cond_t  not_empty;    /* consumers wait here */
    pthread_cond_t  not_full;     /* producers wait here */
} queue_t;

static void q_init(queue_t *q) {
    q->head = q->tail = q->count = 0;
    q->closed = false;
    pthread_mutex_init(&q->m, NULL);
    pthread_cond_init(&q->not_empty, NULL);
    pthread_cond_init(&q->not_full, NULL);
}

static void q_destroy(queue_t *q) {
    pthread_mutex_destroy(&q->m);
    pthread_cond_destroy(&q->not_empty);
    pthread_cond_destroy(&q->not_full);
}

/* Push one item; blocks while the queue is full. */
static void q_put(queue_t *q, int x) {
    pthread_mutex_lock(&q->m);
    while (q->count == CAP)                      /* WHILE, not if */
        pthread_cond_wait(&q->not_full, &q->m);
    q->buf[q->tail] = x;
    q->tail = (q->tail + 1) % CAP;
    q->count++;
    pthread_cond_signal(&q->not_empty);          /* wake one consumer */
    pthread_mutex_unlock(&q->m);
}

/* Pop one item into *out. Returns false when the queue is closed AND drained. */
static bool q_get(queue_t *q, int *out) {
    pthread_mutex_lock(&q->m);
    while (q->count == 0 && !q->closed)          /* WHILE, not if */
        pthread_cond_wait(&q->not_empty, &q->m);
    if (q->count == 0 && q->closed) {            /* nothing left, ever */
        pthread_mutex_unlock(&q->m);
        return false;
    }
    *out = q->buf[q->head];
    q->head = (q->head + 1) % CAP;
    q->count--;
    pthread_cond_signal(&q->not_full);           /* wake one producer */
    pthread_mutex_unlock(&q->m);
    return true;
}

/* Called once, after every producer has finished, to release blocked consumers. */
static void q_close(queue_t *q) {
    pthread_mutex_lock(&q->m);
    q->closed = true;
    pthread_cond_broadcast(&q->not_empty);       /* wake ALL waiters */
    pthread_mutex_unlock(&q->m);
}

static queue_t q;
static long consumed_sum = 0;   /* checked at the end */
static int  consumed_cnt = 0;
static pthread_mutex_t tally = PTHREAD_MUTEX_INITIALIZER;

static void *producer(void *arg) {
    long id = (long)arg;
    for (int i = 0; i < ITEMS; i++) {
        int value = (int)(id * 1000 + i);
        q_put(&q, value);
    }
    return NULL;
}

static void *consumer(void *arg) {
    (void)arg;
    int x;
    while (q_get(&q, &x)) {
        pthread_mutex_lock(&tally);
        consumed_sum += x;
        consumed_cnt++;
        pthread_mutex_unlock(&tally);
    }
    return NULL;
}

int main(void) {
    q_init(&q);
    pthread_t prod[NPROD], cons[NCONS];

    for (long i = 0; i < NCONS; i++)
        pthread_create(&cons[i], NULL, consumer, NULL);
    for (long i = 0; i < NPROD; i++)
        pthread_create(&prod[i], NULL, producer, (void *)i);

    for (int i = 0; i < NPROD; i++)
        pthread_join(prod[i], NULL);
    q_close(&q);                       /* only after all producers joined */
    for (int i = 0; i < NCONS; i++)
        pthread_join(cons[i], NULL);

    long expected = 0;
    for (long id = 0; id < NPROD; id++)
        for (int i = 0; i < ITEMS; i++)
            expected += id * 1000 + i;

    printf("produced %d items across %d producers\n", NPROD * ITEMS, NPROD);
    printf("consumed %d items, sum=%ld (expected %ld) -> %s\n",
           consumed_cnt, consumed_sum, expected,
           (consumed_cnt == NPROD * ITEMS && consumed_sum == expected)
               ? "OK: nothing lost or duplicated" : "MISMATCH");

    q_destroy(&q);
    return 0;
}

Line by line

  • typedef struct { ... } queue_t; — every field the two sides share lives in one struct: the array, the three indices, the closed flag, the mutex, and both condition variables. Keeping the lock next to the data it protects is a good habit; it documents what the mutex guards.
  • CAP 4 with 24 total items — the capacity is deliberately tiny so producers actually hit the full queue and block on not_full, exercising the back-pressure path rather than a queue that never fills.
  • q_put → while (q->count == CAP) pthread_cond_wait(&q->not_full, &q->m); — the producer holds the mutex, and if the queue is full it atomically releases the lock and sleeps on not_full. The while re-checks after every wakeup (spurious or real).
  • q->tail = (q->tail + 1) % CAP; — the modulo is what makes the buffer circular: tail wraps back to 0 after the last slot, reusing the array forever.
  • pthread_cond_signal(&q->not_empty) inside q_put — after adding, we wake a consumer. Note we signal the opposite condition from the one we waited on: producers wait on not_full, signal not_empty.
  • q_get predicate while (q->count == 0 && !q->closed) — a consumer sleeps only while the queue is empty and still open. Once closed, it stops waiting so it can drain or exit.
  • The if (q->count == 0 && q->closed) block — this is the exit condition. Empty and closed means no item will ever come; the consumer unlocks and returns false, ending its loop in consumer(). Crucially we drain remaining items first (the predicate above lets it fall through while count > 0).
  • q_close sets closed = true and broadcasts — a single signal would wake only one of the three consumers; the other two would block forever. broadcast wakes all so each can observe the flag and exit.
  • The tally mutex — consumed_sum/consumed_cnt are shared across all three consumers, so they need their own lock. It is separate from the queue mutex to keep critical sections short.
  • Ordering in main — consumers are started first (harmless: they just block on the empty queue), producers run, and only after joining every producer do we q_close. Closing earlier could strand items still being produced. Then we join consumers.
  • The final check — because every produced value is distinct and we sum them, matching the expected sum and count proves nothing was lost, duplicated, or corrupted by a race.

Common mistakes

1. Using if instead of while around the wait.

if (q->count == 0)                       // WRONG
    pthread_cond_wait(&q->not_empty, &q->m);
int x = q->buf[q->head];                 // may run with count == 0

A spurious or stolen wakeup lets the consumer fall through and read an empty/stale slot — undefined behaviour, corrupted output, or a crash.

while (q->count == 0 && !q->closed)      // FIXED: re-check every wakeup
    pthread_cond_wait(&q->not_empty, &q->m);

2. Signalling the wrong condition variable.

q->count--;
pthread_cond_signal(&q->not_empty);      // WRONG: we freed a slot, wake a PRODUCER

After consuming, a slot opened up, so a producer waiting on not_full should wake — but this wakes consumers instead, and blocked producers sleep forever (a hang).

q->count--;
pthread_cond_signal(&q->not_full);       // FIXED

3. Touching shared state without the mutex.

void q_put(queue_t *q, int x) {
    while (q->count == CAP) { }           // WRONG: unlocked read + busy-wait
    q->buf[q->tail] = x; q->count++;      // data race
}

Reading/writing count, tail, and the buffer without the lock is a data race (undefined behaviour) and the spin wastes a CPU. Hold the mutex and wait, don't spin.

pthread_mutex_lock(&q->m);
while (q->count == CAP) pthread_cond_wait(&q->not_full, &q->m);
q->buf[q->tail] = x; q->count++;
pthread_mutex_unlock(&q->m);              // FIXED

4. Never signalling shutdown, so consumers hang.

// producers finish, main joins them, then just returns...  // WRONG

Consumers blocked in pthread_cond_wait(&not_empty, ...) never wake; pthread_join on them hangs forever.

q_close(&q);                              // FIXED: sets closed + broadcast(not_empty)
for (int i = 0; i < NCONS; i++) pthread_join(cons[i], NULL);

Debugging tips

  • Hangs (the program never exits): almost always a lost wakeup or a missing/mismatched signal. Attach with gdb -p <pid>, then thread apply all bt. Threads parked in pthread_cond_wait show the frame clearly; check which condition variable each is on and whether anyone ever signals it. A consumer stuck on not_empty with all producers gone means you forgot to close/broadcast.
  • Wrong or corrupted output / occasional crashes: run under a race detector. cc -fsanitize=thread (ThreadSanitizer) instruments the build and prints the exact two stack traces of a data race — invaluable for spotting an unlocked access to count or the buffer. Note TSan and Valgrind's --tool=helgrind are alternatives; helgrind also flags lock-ordering problems.
  • valgrind --tool=helgrind ./m reports mutex/condition-variable misuse: destroying a locked mutex, signalling patterns it deems suspicious, and potential deadlocks from inconsistent lock order (the deadlock skills apply directly).
  • printf debugging done right: print count, head, tail, and the thread id at the start of each critical section — but do the printing inside the lock, or the log lines interleave and lie to you.
  • Reproduce timing bugs: shrink CAP to 1 and add many threads, or insert a tiny usleep right after a cond_signal to widen the window where a lost wakeup or stolen item would manifest. If the bug appears only with the delay, you have a synchronization flaw, not bad luck.

Memory safety

  • Data races are undefined behaviour. Every read or write of head, tail, count, closed, or any buffer slot must happen with the queue mutex held. "It printed the right thing once" is not correctness; the compiler and CPU may reorder or cache unlocked accesses.
  • The predicate and its wait must share one mutex. The condition variable you wait on must be paired with the exact mutex that protects the predicate you test. Using a different mutex breaks the atomic unlock-and-sleep guarantee and reintroduces lost wakeups.
  • Reading an empty buffer. The if/while bug (mistake #1) is a use-of-uninitialised / stale-data hazard: a consumer dequeues from a slot that was never written or was already consumed. In security terms this can leak previous contents or feed uninitialised bytes downstream.
  • Bounded buffer = bounded memory. A fixed-capacity circular buffer gives natural back-pressure: producers block instead of allocating without limit. An unbounded queue under a fast producer and slow consumer is a memory-exhaustion (DoS) risk — a real concern for network-facing services. Prefer a bound and decide explicitly what to do when it's hit (block, or drop with a counter).
  • Lifetime and destroy order. Never pthread_cond_destroy/pthread_mutex_destroy while a thread might still be waiting or locking — join all users first (as main does). Destroying a live sync object is undefined behaviour.
  • Don't leak the mutex on early return. Every path out of a critical section — including the return false in q_get — must unlock first. A forgotten unlock on an error path is an instant deadlock.

Real-world uses

  • Thread-pool servers: an acceptor thread produces connections/requests into a bounded queue; a fixed pool of worker threads consumes them. This is the backbone of web servers, RPC frameworks, and databases.
  • Asynchronous logging: application threads produce log records into a queue; one background thread consumes and writes them, so slow disk I/O never blocks request handling.
  • Media and data pipelines: a decoder produces frames into a ring buffer that the renderer consumes; audio callbacks consume samples a producer fills. Ring buffers dominate here for their zero-allocation, cache-friendly behaviour.
  • Best practice: keep critical sections tiny — do the queue surgery under the lock, do the actual work (parsing, I/O, computation) outside it. Always bound the queue and define an overflow policy. Prefer signal for single-item hand-offs and reserve broadcast for state changes (like shutdown) that many waiters must observe. For very high throughput, look at lock-free ring buffers or per-worker queues, but reach for those only after a correct mutex+condvar version is measured and proven to be the bottleneck.

Practice tasks

  1. Instrument the queue. Add a long total_put, total_get pair (guarded by the queue mutex) and print them at shutdown. Confirm they match NPROD * ITEMS, then deliberately change one signal to the wrong condition variable and observe the hang under gdb.

  2. Prove the while matters. Temporarily replace the while in q_get with an if, then add a usleep(1000) right after the consumer's wait returns. Add many consumers and a CAP of 1, and show the program can read past count == 0. Then restore the while and show it's fixed.

  3. Add a blocking q_put timeout. Implement bool q_put_timed(queue_t*, int x, int ms) using pthread_cond_timedwait and clock_gettime(CLOCK_REALTIME, ...), returning false if the queue stays full past the deadline instead of blocking forever.

  4. Generalise to any element type. Change the buffer to hold void * items and make the queue store pointers, so it can carry heap-allocated jobs. Decide and document who frees each item, and verify with valgrind --leak-check=full that nothing leaks.

  5. Graceful drain-then-stop under multiple producers. Extend shutdown so a separate signal (e.g. Ctrl-C via a flag) tells producers to stop and guarantees every already-queued item is still consumed before consumers exit. Verify the consumed sum equals the produced sum for whatever count was reached.

Summary

  • Producer/consumer = bounded (circular) buffer + one mutex guarding all queue state + two condition variables, not_empty (consumers wait) and not_full (producers wait).
  • Producers wait while the queue is full and signal not_empty after adding; consumers wait while it is empty and signal not_full after removing. Signal the opposite condition from the one you waited on.
  • Always wait inside while (predicate-not-yet-true), never a plain if — spurious and stolen wakeups mean every wakeup is only a hint to re-check.
  • pthread_cond_wait atomically unlocks the mutex, sleeps, and re-locks on return; that atomicity is why you must hold the mutex and why the predicate and the wait share one mutex.
  • Use signal for a single-item hand-off; use broadcast for shutdown so every blocked consumer wakes, sees closed, drains, and exits exactly once.
  • Touch shared state only under the lock, unlock on every exit path, bound the queue to avoid memory-exhaustion, and join all threads before destroying the sync objects.

Practice with these exercises