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
- producer produces things
- consumer consumes things
- there are producers, and there are consumers
- 1 producer 1 consumer is a simple problem, producer adds item to queue, consumer dequeue the queue, easy peasy
- What happens on multiple producers and muliple consumers
- consider a factory
- where
producerput things - and from this factory consumer put things
- where
Assumptions
- Buffer is too large, no constraint on its size
- Only 1 producer and only 1 consumer
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.

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?
- Because, of multi-threading producers, We cannot be sure that all thousands were created because of the concurrency between the producer itself because we have five stages in the creation of the new node producers interleave among themselves and hence we cannot see thousand unique elements were pushed into the queue itself. First thing is that.
- ConcurrentLinkedQueue
- One thread is pushing something, another thread is unable to read it (Volatile issue?). Let’s see
- Is that the problem??
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?
- The producer touches the tail, the consumer touches the head. Logically they’re independent operations on opposite ends of the buffer
- But one lock forces them to take turns anyway. On a multicore machine, one core idles while the other works — you’re paying for parallelism you can’t use
- The concrete costs:
- Lost parallelism — the two threads could have run simultaneously and don’t
- Context switches — a contended
lock()parks the thread; park/unpark are OS-level and cost microseconds, vastly more than the work being protected - Cache-line ping-pong — the lock word bounces between cores on every acquisition
- A throughput ceiling — every single item funnels through one lock, so max throughput is
1 / (critical section + lock overhead), no matter how many cores you own
What each option actually gives you
| capacity | blocks / backpressure | producer & consumer in parallel | |
|---|---|---|---|
hand-written FixedBuffer | ✓ | ✓ | ✗ — one lock |
ArrayBlockingQueue | ✓ | ✓ | ✗ — one lock, same design |
LinkedBlockingQueue | ✓ optional | ✓ | ✓ — two locks |
ConcurrentLinkedQueue | ✗ | ✗ | ✓ — lock-free |
Channel (coroutines) | ✓ | ✓ suspends | no 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 }
| Structure | Mechanism | CAS in its own code? |
|---|---|---|
ArrayDeque, LinkedList, ArrayList, HashMap | nothing | No — not thread-safe at all |
ArrayBlockingQueue | 1 ReentrantLock + 2 conditions | No (the lock uses CAS internally) |
PriorityBlockingQueue, DelayQueue | 1 ReentrantLock | No |
CopyOnWriteArrayList | lock for writes, volatile array swap | No |
Hashtable, Collections.synchronizedList | synchronized | No |
LinkedBlockingQueue | 2 locks + AtomicInteger count | Partly — CAS for the counter only |
ConcurrentHashMap | CAS to fill an empty bin, synchronized per bin on collision | Hybrid |
ConcurrentLinkedQueue / Deque | lock-free Michael–Scott | Yes — primary |
ConcurrentSkipListMap / Set | lock-free | Yes — primary |
LinkedTransferQueue, SynchronousQueue | lock-free dual queue | Yes — primary |
AtomicInteger / Long / Reference | CAS + retry loop | Yes — it is CAS |
LongAdder | striped CAS across cells | Yes |
kotlinx Channel | lock-free segment queue | Yes |
Rule of thumb
Concurrent*in the name → CAS-primary, lock-free. No thread ever blocks another; losers retry*BlockingQueuein the name → lock-primary. Threads park and are signalled
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?
| Storage | Circular? | Capacity | Allocation per item | |
|---|---|---|---|---|
ArrayBlockingQueue | one Object[], allocated once | YES | fixed, physical — the array size | none |
LinkedBlockingQueue | singly-linked nodes + sentinel head | no | optional, enforced by an AtomicInteger count | one node |
ConcurrentLinkedQueue | singly-linked nodes, CAS’d head/tail | no | none — unbounded | one 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.