NOTE

6.30 Producer-Consumer

1. Using BlockingQueue. 2. Using wait/notify. 3. Using Lock/Condition. Compared with wait/notify, two Conditions are used so producers and consumers are not woken together; each side wakes only the other side.

JavaCreated Updated 1 min readhistorical

This is a historical learning note and may contain outdated or incomplete understanding.

1. Using BlockingQueue

public class ProducerConsumer
{

    public static void main(String[] args) throws InterruptedException
    {
        SynchronousQueue<Integer> queue = new SynchronousQueue<>();
        Producer producer = new Producer(queue);
        Consumer consumer = new Consumer(queue);

        Thread thread1 = new Thread(producer);
        Thread thread2 = new Thread(consumer);

        thread1.start();
        thread2.start();
        thread1.join();
        thread2.join();
    }
}


class Producer implements Runnable
{

    private SynchronousQueue<Integer> queue;

    public Producer(SynchronousQueue<Integer> queue)
    {
        this.queue = queue;
    }

    @Override
    public void run()
    {
        for (int i = 0; i < 10000; i++)
        {
            try
            {
                TimeUnit.SECONDS.sleep(1);
                System.out.println("Producer: " + i);
                queue.put(i);
            }
            catch (InterruptedException e)
            {
                e.printStackTrace();
            }

        }
    }
}


class Consumer implements Runnable
{
    private SynchronousQueue<Integer> queue;

    public Consumer(SynchronousQueue<Integer> queue)
    {
        this.queue = queue;
    }
    @Override
    public void run()
    {
        while (true)
        {
            try
            {
                Integer val = queue.take();
                System.out.println("Consumer: " + val);

            }
            catch (InterruptedException e)
            {
                e.printStackTrace();
            }
        }
    }
}

2. Using wait / notify

public class ProducerConsumer2
{

    static class Producer implements Runnable
    {
        private final List<Integer> queue;
        private final int fullSize;
        private final Object lock;

        public Producer(List<Integer> queue, int fullSize, Object lock)
        {
            this.queue = queue;
            this.fullSize = fullSize;
            this.lock = lock;
        }

        @Override
        public void run()
        {

            for (int i = 0; i < 10000000; i++)
            {
                synchronized (lock)
                {
                    while (queue.size() == fullSize)
                    {
                        try
                        {
                            lock.wait();
                        }
                        catch (InterruptedException e)
                        {
                            e.printStackTrace();
                        }
                    }
                    System.out.println(Thread.currentThread().getName() + " put val: " + i);
                    queue.add(i);
                    lock.notifyAll();
                }

            }
        }
    }

    static class Consumer implements Runnable
    {
        private final List<Integer> queue;
        private final int fullSize;
        private final Object lock;

        public Consumer(List<Integer> queue, int fullSize, Object lock)
        {
            this.queue = queue;
            this.fullSize = fullSize;
            this.lock = lock;
        }

        @Override
        public void run()
        {

            while (true)
            {
                synchronized (lock)
                {
                    while (queue.isEmpty())
                    {
                        try
                        {
                            lock.wait();
                        }
                        catch (InterruptedException e)
                        {
                            e.printStackTrace();
                        }
                    }

                    Integer val = queue.remove(0);
                    System.out.println(Thread.currentThread().getName() + " get val: " + val);

                    lock.notifyAll();
                }
            }
        }
    }


    public static void main(String[] args) throws InterruptedException
    {
        final List<Integer> queue = new ArrayList<>();
        final int fullSize = 5;
        final Object lock = ProducerConsumer2.class;
        Producer producer = new Producer(queue, fullSize, lock);
        Consumer consumer = new Consumer(queue, fullSize, lock);

        Thread thread1 = new Thread(producer, "producer");
        Thread thread2 = new Thread(consumer, "consumer");


        thread1.start();
        thread2.start();
        thread1.join();
        thread2.join();

    }
}

3. Using Lock / Condition

Compared with the wait / notify implementation above, this version uses two Condition objects. This means a wake-up does not wake producers and consumers together; only the relevant side needs to be woken (the producer wakes the consumer, and the consumer wakes the producer).

public class ConditionTest
{
    private Lock lock;// One lock means reads and writes are mutually exclusive
    private int capacity;
    private List<Object> items;
    private Condition notFull;// Used to wake writer threads
    private Condition notEmpty;// Used to wake reader threads

    public ConditionTest(int capacity)
    {
        this.capacity = capacity;
        this.items = new ArrayList<>();
        this.lock = new ReentrantLock();
        this.notFull = lock.newCondition();
        this.notEmpty = lock.newCondition();
    }

    public void add(Object data) throws InterruptedException
    {
        try
        {
            lock.lock();
            // When adding, if it is already full, wait for the not-full condition to wake this thread
            while (this.items.size() == capacity)
            {
                this.notFull.await();
            }

            // An element was added, so signal not-empty
            this.items.add(data);
            this.notEmpty.signalAll();

        }
        finally
        {
            lock.unlock();
        }
    }

    public Object remove() throws InterruptedException
    {
        try
        {
            lock.lock();

            // When removing, if it is already empty, wait for the not-empty condition to wake this thread
            while (this.items.size() == 0)
            {
                this.notEmpty.await();
            }

            // An element was removed, so signal not-full
            Object data = this.items.remove(0);
            this.notFull.signalAll();
            return data;
        }
        finally
        {
            lock.unlock();
        }
    }

    public static void main(String[] args)
    {
        ConditionTest conditionTest = new ConditionTest(5);
        new Thread(() -> {
            for (int i = 0; i < 1000; i++)
            {
                try
                {
                    conditionTest.add(i);
                    System.out.println(String.format("Producer puts %d", i));
                }
                catch (InterruptedException e)
                {
                    e.printStackTrace();
                }
            }
        }).start();

        new Thread(() -> {
            try
            {
                while (true)
                {
                    Object data = conditionTest.remove();
                    System.out.println(String.format("Consumer consumes %d", data));
                }
            }
            catch (InterruptedException e)
            {
                e.printStackTrace();
            }
        }).start();



    }
}

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub