Producer-Consumer Pattern and Storage Model Implementation

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

  1. Blocking Queue Approach: Java's BlockingQueue interface provides thread-safe queue implementations suitable for producer-consumer scenarios. It automatically handles thread synchronization and blocking operations.
  2. Wait/Notify Mechanism: Using Object's wait(), notify(), and notifyAll() methods combined with synchronized blocks to coordinate threads. Producers wait when the queue is full, and consumers wait when empty.
  3. Lock and Condition Approach: Utilizing ReentrantLock and Condition objects 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:

  1. Multiple producers and consumers operating on the same product object
  2. Data corruption issues like:
    • null -- 0.0
    • Huawei -- 0.0
  3. Race conditions where brand is set but price update is interrupted by consumer threads
  4. Solution: Implement locking to ensure complete product creation before consumption
  5. 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

  1. Shared Resource: Typically implemented as a queue or buffer storing produced items for consumer access
  2. Producers: Responsible for creating and depositing products, waiting when storage is full
  3. Consumers: Retrieve and process products from storage, waiting when empty
  4. 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;
    }
}

Tags: java multithreading Concurrency producer-consumer pattern-design

Posted on Sat, 03 Oct 2026 16:54:32 +0000 by eRott