Producer-Consumer Pattern
The producer-consumer patern is a fundamental concurrency design pattern in multithreaded programming that addresses data exchange and synchronization between producers and consumers. In this model, producers generate data and place it into a shared queue, while consumers retrieve and process data from the same queue. This approach decouples producers and consumers, anabling efficient collaboration.
Important Consideration: Thread safety must be carefully managed to prevent race conditions and data inconsistency in multithreaded environments.
Core Objective
The primary goal is to ensure one-to-one production and consumption of data items.
Common Implementation Approaches
- Blocking Queue Approach: Java's
BlockingQueueinterface provides thread-safe queue implementations suitable for producer-consumer scenarios. It automatically handles thread synchronization and blocking operations. - Wait/Notify Mechanism: Using
Object'swait(),notify(), andnotifyAll()methods combined withsynchronizedblocks to coordinate threads. Producers wait when the queue is full, and consumers wait when empty. - Lock and Condition Approach: Utilizing
ReentrantLockandConditionobjects for more flexible thread control compared to traditional wait/notify mechanisms.
Implementation Example
Problem Analysis:
- Product Class: Phone (with brand and price attributes)
- Producer Thread: Producer
- Consumer Thread: Consumer
Key Challenges:
- Multiple producers and consumers operating on the same product object
- Data corruption issues like:
- null -- 0.0
- Huawei -- 0.0
- Race conditions where brand is set but price update is interrupted by consumer threads
- Solution: Implement locking to ensure complete product creation before consumption
- Maintain one-produce-one-consume cycle using inventory tracking
public class MainApplication {
public static void main(String[] args) {
Phone product = new Phone();
Producer prod1 = new Producer(product);
Producer prod2 = new Producer(product);
Consumer cons1 = new Consumer(product);
Consumer cons2 = new Consumer(product);
prod1.start();
prod2.start();
cons1.start();
cons2.start();
}
}
class Consumer extends Thread {
private Phone phone;
public Consumer(Phone phone) {
this.phone = phone;
}
@Override
public void run() {
while(true) {
synchronized(phone) {
while(!phone.hasInventory()) {
try {
phone.wait();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
System.out.println(phone.getName() + " -- " + phone.getValue());
phone.setInventoryStatus(false);
phone.notifyAll();
}
}
}
}
class Phone {
private String name;
private double value;
private boolean available;
// Constructors and getters/setters
public Phone() {}
public String getName() { return name; }
public void setName(String name) { this.name = name; }
public double getValue() { return value; }
public void setValue(double value) { this.value = value; }
public boolean hasInventory() { return available; }
public void setInventoryStatus(boolean status) { this.available = status; }
}
class Producer extends Thread {
private Phone phone;
private static boolean toggle = true;
public Producer(Phone phone) {
this.phone = phone;
}
@Override
public void run() {
while(true) {
synchronized(phone) {
while(phone.hasInventory()) {
try {
phone.wait();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
if(toggle) {
phone.setName("Huawei");
phone.setValue(3999);
} else {
phone.setName("Xiaomi");
phone.setValue(1999);
}
toggle = !toggle;
phone.setInventoryStatus(true);
phone.notifyAll();
}
}
}
}
Storage Model
The storage model represents shared resources in producer-consumer problems, functioning as a warehouse or buffer zone. Producers deposit products into storage while consumers retrieve them for processing. Thread safety is crucial for maintaining data integrity in concurrent environments.
Essential Components
- Shared Resource: Typically implemented as a queue or buffer storing produced items for consumer access
- Producers: Responsible for creating and depositing products, waiting when storage is full
- Consumers: Retrieve and process products from storage, waiting when empty
- Thread Safety: Proper synchronization mechanisms prevent race conditions and ensure data consistency
Implementation Strategies
Similar approaches apply including blocking queues, wait/notify patterns, and lock/condition mechanisms. Effective synchronization design ensures reliable data exchange between producers and consumers.
Warehouse Implementation Example
Requirements:
- Producer threads continuously create cakes and add them to warehouse
- Consumer threads continuously remove and consume cakes
- First-in-first-out (FIFO) processing order
Design Elements:
- Cake class and Warehouse class
- Warehouse contains LinkedList collection for cake storage
- Support for multiple producer and consumer threads
import java.time.LocalDateTime;
import java.util.LinkedList;
public class WarehouseDemo {
public static void main(String[] args) {
Warehouse storage = new Warehouse(20);
BakeryProducer bp1 = new BakeryProducer(storage);
BakeryProducer bp2 = new BakeryProducer(storage);
CustomerConsumer cc1 = new CustomerConsumer(storage);
CustomerConsumer cc2 = new CustomerConsumer(storage);
bp1.start();
bp2.start();
cc1.start();
cc2.start();
}
}
class Pastry {
private String manufacturer;
private String productionTime;
public Pastry() {}
public Pastry(String maker, String time) {
this.manufacturer = maker;
this.productionTime = time;
}
// Getters and setters
public String getManufacturer() { return manufacturer; }
public void setManufacturer(String manufacturer) { this.manufacturer = manufacturer; }
public String getProductionTime() { return productionTime; }
public void setProductionTime(String time) { this.productionTime = time; }
@Override
public String toString() {
return "Pastry [maker=" + manufacturer + ", time=" + productionTime + "]";
}
}
class CustomerConsumer extends Thread {
private Warehouse warehouse;
public CustomerConsumer(Warehouse warehouse) {
this.warehouse = warehouse;
}
@Override
public void run() {
while(true) {
warehouse.retrieve();
}
}
}
class BakeryProducer extends Thread {
private Warehouse warehouse;
public BakeryProducer(Warehouse warehouse) {
this.warehouse = warehouse;
}
@Override
public void run() {
while(true) {
Pastry item = new Pastry("Taoli", LocalDateTime.now().toString());
warehouse.deposit(item);
}
}
}
class Warehouse {
private int currentItems;
private int maximumCapacity;
private LinkedList<Pastry> container;
public Warehouse(int maxCapacity) {
this.maximumCapacity = maxCapacity;
this.container = new LinkedList<>();
}
public synchronized void deposit(Pastry pastry) {
while(currentItems >= maximumCapacity) {
try {
this.wait();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
container.add(pastry);
currentItems++;
System.out.println("Deposit - Current stock: " + currentItems);
this.notifyAll();
}
public synchronized Pastry retrieve() {
while(currentItems <= 0) {
try {
this.wait();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
Pastry item = container.removeFirst();
currentItems--;
System.out.println("Retrieve - Current stock: " + currentItems + ", Sold item: " + item);
this.notifyAll();
return item;
}
}