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.
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