Skip to content

Producer-Consumer problem

Published: at 02:26 PM
Modified: (12 min read)

Overview

Producer-Consumer problem is a classical synchronization problem in the operating system. With the presence of more than one process and limited resources in the system the synchronization problem arises. If one resource is shared between more than one process at the same time then it can lead to data inconsistency. In the producer-consumer problem, the producer produces an item and the consumer consumes the item produced by the producer.

Problem

Assumptions

At later point of time, we would be doing limited buffer + more than 1 producer + more than 1 consumer

1 Producer 1 Consumer - No buffer restriction

Simple 1 producer, 1 consumer

class SimpleQueue<T> {
    private val items = LinkedList<T>()

    val size: Int get() = items.size

    /** Producer side: new item goes on the back. */
    fun enqueue(item: T) {
        items.addLast(item)
    }

    /** Consumer side: oldest item comes off the front; null if nothing buffered. */
    fun dequeue(): T? = items.pollFirst()
}

fun main() {
    val queue = SimpleQueue<Int>()
    val totalItems = 1_000

    val producer = thread(name = "producer") {
        for (i in 1..totalItems) {
            queue.enqueue(i)
            println("produced $i  (size ${queue.size})")
            Thread.sleep(1)          // production pace
        }
        println("producer done: $totalItems items")
    }

    val consumer = thread(name = "consumer") {
        var consumed = 0
        var emptyPolls = 0
        while (consumed < totalItems) {
            val item = queue.dequeue()
            if (item == null) {
                emptyPolls++         // producer hasn't caught up; spin
            } else {
                consumed++
                println("consumed $item  (size ${queue.size}) $consumed")
            }
            Thread.sleep(2)          // consumption pace, deliberately slower
        }
        println("consumer done: $consumed items, $emptyPolls empty polls")
    }

    producer.join()
    consumer.join()
    println("main done, leftover in queue: ${queue.size}")
}

This above function runs perfectly fine because we have one producer, one consumer, and there is no shared variable which is being manipulated by these two of the threads. So these threads are just competing, and one is creating thousand elements and another one is consuming thousand elements and so consumer is also keeping a record of how many empty polls it had

Slow producer

In this above where we have only one producer and one consumer, even though the producer is slow, consumer would be waiting for him and the empty poles would be increasing for the consumer, so that is handled

Slow consumer

So if in the above situation the consumer is slow, then the producer will produce all the items, thousand items, and the consumer will consume all those thousand items with some time, and that is also handled in this case.

Why concurrent modification crash did not happen?? 1 thread adding, 1 thread removing??

Because nothing ever iterates the list, and ConcurrentModificationException is thrown only by iterators.

Image

Multiple Producer Single Consumer - No buffer restriction

class SimpleQueue<T> {
    private val items = LinkedList<T>()

    val size: Int get() = items.size

    /** Producer side: new item goes on the back. */
    fun enqueue(item: T) {
        items.addLast(item)
    }

    /** Consumer side: oldest item comes off the front; null if nothing buffered. */
    fun dequeue(): T? = items.pollFirst()
}

fun main() {...}

Output

// 1st run
received by consumer : 976

// 2nd run
received by consumer : 32

// see huge gap

Why we see duplicates, what is leading to duplicates?? Why we see unique items being so low?

Using ConcurrentLinkedQueue

private val items = ConcurrentLinkedQueue<T>()

output

offered by producers : 1000  (ranges: 1..334, 335..667, 668..1000)
received by consumer : 1000
distinct ids         : 1000
delivered twice      : 0
never delivered      : 0
stranded in queue    : 0
empty polls          : 0
producer exceptions  : 0, 0, 0
consumer exceptions  : 0
RESULT: looks correct this run — run it again

And everytime you run it, it gives u same result, hence its not a volatile issue, pure concurrency issue at producer end only and solved via thread safe data structure

1 Producer 1 consumer but Fixed sized buffer

class FixedBuffer<T>(private val capacity: Int = 10) {
    private val items = LinkedList<T>()

    private val lock = ReentrantLock()
    private val notFull: Condition = lock.newCondition()   // producers wait here
    private val notEmpty: Condition = lock.newCondition()  // consumers wait here

    /** Also takes the lock — reading size outside it is the same data race as before. */
    val size: Int get() = lock.withLock { items.size }

    /** Producer side: parks while the buffer is full. */
    fun put(item: T) = lock.withLock {
        // while, not if: await() can return spuriously, and another producer may
        // have taken the free slot before this thread reacquired the lock.        while (items.size == capacity) {
            notFull.await()
        }
        items.addLast(item)
        notEmpty.signal()
    }

    /** Consumer side: parks while the buffer is empty. */
    fun take(): T = lock.withLock {
        while (items.isEmpty()) {
            notEmpty.await()
        }
        val item = items.pollFirst()
        notFull.signal()
        item
    }
}

fun main() {...}

This above function does not works

producer done: offered 1000 items
consumer done: took 1000 items

capacity             : 10
largest size seen    : 10
offered by producer  : 1000
received by consumer : 1000
distinct ids         : 1000
delivered twice      : 0
never delivered      : 0
left in buffer       : 0
in order             : true

But there is a dataStructure with use this.

There is one ReentrantLock guarding the buffer, and two conditions on it: notFull and notEmpty. Each is named for the state a thread is waiting for. When the producer calls put on a full buffer, it already holds the lock. It sees size == capacity and calls notFull.await(), which atomically releases the lock and parks the thread — no CPU used, and the lock is now free for the consumer. The consumer acquires the lock, removes an item, and calls notFull.signal(). That doesn’t release anything; it moves the parked producer onto the lock’s queue. The lock is actually released when the consumer exits its withLock block. The producer then reacquires the lock, await() returns, and the while loop re-checks the condition — because between the signal and the reacquisition the state could have changed again. Seeing room, it adds its item and calls notEmpty.signal() for the consumer’s benefit. The empty case is the mirror image: the consumer waits on notEmpty, and put signals it. The two methods never run at the same time — one lock means strict alternation when both are active. That serialization is the correctness guarantee, and also the cost.

A diagram to show lock and unlock

ArrayBlockingQueue - Circular Queue

ArrayBlockingQueue answers it properly. It is a fixed array plus one lock and two conditions: put() parks the producer while full, take() parks the consumer while empty, and each wakes the other. A parked thread costs no CPU, which is the whole difference from the spin.

ArrayBlockingQueue does not solve Serialisation. It uses a single ReentrantLock for both put and take — the same design you just wrote by hand. What it solved was the hang, not the serialization. Those are two different problems.

class FixedBuffer<T>(capacity: Int = 10) {
    private val items = ArrayBlockingQueue<T>(capacity)

    val size: Int get() = items.size

    /** Producer side: blocks while the buffer is full. */
    fun put(item: T) = items.put(item)

    /** Consumer side: blocks while the buffer is empty. */
    fun take(): T = items.take()
}

fun main() {  ... }

ConcurrentLinkedQueue

class UnboundedBuffer<T> {
    private val items = ConcurrentLinkedQueue<T>()

    /** Careful: ConcurrentLinkedQueue.size() walks the list. O(n), not a field read. */
    val size: Int get() = items.size

    /** Producer side: never blocks, never full. */
    fun put(item: T) {
        items.offer(item)
    }

    /** Consumer side: never blocks; null means empty right now. */
    fun poll(): T? = items.poll()
}

fun main() {
    val buffer = UnboundedBuffer<Int>()
    val totalItems = 1_000

    var maxSeen = 0
    var emptyPolls = 0L
    val consumed = ArrayList<Int>(totalItems)

    val producer = thread(name = "producer") {
        for (item in 1..totalItems) {
            buffer.put(item)                        // never waits
            maxSeen = maxOf(maxSeen, buffer.size)
        }
        println("producer done: offered $totalItems items")
    }

    val consumer = thread(name = "consumer") {
        // No deadline needed: CAS guarantees nothing is lost, so this always
        // reaches totalItems. The cost is spinning while the queue is empty.        while (consumed.size < totalItems) {
            val item = buffer.poll()
            if (item == null) emptyPolls++ else consumed.add(item)
        }
        println("consumer done: took ${consumed.size} items")
    }

    producer.join()
    consumer.join()

    val distinct = consumed.toSet()
    val missing = (1..totalItems).toSet() - distinct

    println()
    println("capacity             : unbounded")
    println("largest size seen    : $maxSeen")
    println("offered by producer  : $totalItems")
    println("received by consumer : ${consumed.size}")
    println("distinct ids         : ${distinct.size}")
    println("delivered twice      : ${consumed.size - distinct.size}")
    println("never delivered      : ${missing.size}")
    println("left in buffer       : ${buffer.size}")
    println("consumer spun on empty: $emptyPolls")
    println("in order             : ${consumed == (1..totalItems).toList()}")
}

Works like charm

How Serialisation is the cost?

What each option actually gives you

capacityblocks / backpressureproducer & consumer in parallel
hand-written FixedBuffer✓✓✗ — one lock
ArrayBlockingQueue✓✓✗ — one lock, same design
LinkedBlockingQueue✓ optional✓✓ — two locks
ConcurrentLinkedQueue✗✗✓ — lock-free
Channel (coroutines)✓✓ suspendsno threads parked at all

LinkedBlockingQueue is the one that actually attacks serialization. It’s the two-lock queue algorithm: a separate putLock and takeLock, with an AtomicInteger count as the only shared state. A producer at the tail and a consumer at the head genuinely run at the same time. The price is a node allocation per item and a size() that’s a snapshot rather than exact.

ConcurrentLinkedQueue removes serialization entirely — CAS, no locks, threads never block each other. But it gives up exactly the things a fixed buffer is about: no capacity, no backpressure, poll() returns null instead of waiting, and size() is O(n). That’s why it fixes multi-producer losses but is the wrong tool for a fixed buffer.

CAS

Compare-And-Swap is a single atomic CPU instruction that writes a new value to a memory location only if that location still holds the value you expected, and reports whether it succeeded.

The one-line version for the article: it turns “read, then write” into a single step that fails instead of clobbering. Everything else — AtomicInteger, lock-free queues, even ReentrantLock itself — is a retry loop wrapped around that one instruction.

CAS(address, expected, new) =
    if (*address == expected) { *address = new; true }
    else                      { false }
StructureMechanismCAS in its own code?
ArrayDeque, LinkedList, ArrayList, HashMapnothingNo — not thread-safe at all
ArrayBlockingQueue1 ReentrantLock + 2 conditionsNo (the lock uses CAS internally)
PriorityBlockingQueue, DelayQueue1 ReentrantLockNo
CopyOnWriteArrayListlock for writes, volatile array swapNo
Hashtable, Collections.synchronizedListsynchronizedNo
LinkedBlockingQueue2 locks + AtomicInteger countPartly — CAS for the counter only
ConcurrentHashMapCAS to fill an empty bin, synchronized per bin on collisionHybrid
ConcurrentLinkedQueue / Dequelock-free Michael–ScottYes — primary
ConcurrentSkipListMap / Setlock-freeYes — primary
LinkedTransferQueue, SynchronousQueuelock-free dual queueYes — primary
AtomicInteger / Long / ReferenceCAS + retry loopYes — it is CAS
LongAdderstriped CAS across cellsYes
kotlinx Channellock-free segment queueYes

Rule of thumb

Two exceptions worth flagging so the rule doesn’t mislead: SynchronousQueue and LinkedTransferQueue are lock-free despite the naming, and LinkedBlockingQueue is a hybrid — locks for the links, CAS for the count.

JAVA Data structures

What is circular queue and whats not, ArrayBlockingQueue, LinkedBlockingQueue, ConcurrentLinkedQueue?

StorageCircular?CapacityAllocation per item
ArrayBlockingQueueone Object[], allocated onceYESfixed, physical — the array sizenone
LinkedBlockingQueuesingly-linked nodes + sentinel headnooptional, enforced by an AtomicInteger countone node
ConcurrentLinkedQueuesingly-linked nodes, CAS’d head/tailnonone — unboundedone node

Channels to the rescue

class SuspendingBuffer<T>(capacity: Int = 10) {
    private val channel = Channel<T>(capacity)

    /** Producer side: suspends while the buffer is full. */
    suspend fun put(item: T) = channel.send(item)

    /** Consumer side: suspends while the buffer is empty. */
    suspend fun take(): T = channel.receive()

    /** No more items will be sent. Anything already buffered can still be taken. */
    fun close() {
        channel.close()
    }
}

fun main() = runBlocking {
    val buffer = SuspendingBuffer<Int>(capacity = 10)
    val totalItems = 1_000

    val consumed = ArrayList<Int>(totalItems)

    // Dispatchers.Default so these really are on different threads, like the
    // earlier stages. On a single thread the result would be identical — which    // is itself the point: correctness here does not come from thread layout.    val producer = launch(Dispatchers.Default) {
        for (item in 1..totalItems) {
            buffer.put(item)                        // suspends here when full
        }
        buffer.close()
        println("producer done: offered $totalItems items")
    }

    val consumer = launch(Dispatchers.Default) {
        repeat(totalItems) {
            consumed.add(buffer.take())             // suspends here when empty
        }
        println("consumer done: took ${consumed.size} items")
    }

    producer.join()
    consumer.join()

    val distinct = consumed.toSet()
    val missing = (1..totalItems).toSet() - distinct

    println()
    println("capacity             : 10")
    println("offered by producer  : $totalItems")
    println("received by consumer : ${consumed.size}")
    println("distinct ids         : ${distinct.size}")
    println("delivered twice      : ${consumed.size - distinct.size}")
    println("never delivered      : ${missing.size}")
    println("in order             : ${consumed == (1..totalItems).toList()}")
}

BlockingQueue -> parks an OS THREAD. The kernel is involved, the thread’s stack sits idle in memory, and waking it is a context switch.

Channel -> suspends a COROUTINE. The continuation becomes a heap object, the thread is handed back to the dispatcher and goes off to run something else. No kernel, no context switch.