前言
本文涉及两种方法实现生产者消费者模式。
理想情况下,如果生产者产出数据的速度大于消费者消费的速度,并且当生产出来的数据累积到一定程度的时候,那么生产者必须暂停等待一下(阻塞生产者线程),以便等待消费者线程把累积的数据处理完毕。我们可以通过 BlockingQueue 完成。BlockingQueue 位于java.util.concurrent 包中,创建 java.util.concurrent 的目的就是要实现 Collection 框架对数据结构所执行的并发操作。不了解的请戳我的blog-BlockingQueue.
通过 PriorityQueue 等 ,java.util 下的集合类实现。自己处理同步问题。
BlockingQueue 实现
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.LinkedBlockingQueue;
/**
* 作者:杨剑飞. 邮箱:847564732@qq.com 创建日期:2016-3-7 功能:
*/
public class BlockingQueueTest {
public static void main(String[] args) throws InterruptedException {
// 声明一个容量为10的缓存队列
BlockingQueue<String> queue = new LinkedBlockingQueue<String>(10);
Producer producer1 = new Producer(queue);
Producer producer2 = new Producer(queue);
Producer producer3 = new Producer(queue);
Consumer consumer = new Consumer(queue);
// 借助Executors
ExecutorService service = Executors.newCachedThreadPool();
// 启动线程
new Thread(producer1).start();
new Thread(producer2).start();
new Thread(producer3).start();
new Thread(consumer).start();
// service.submit(producer1);
// service.submit(producer2);
// service.submit(producer3); //需要返回值,用 submint
// service.execute(consumer); //不需要返回值,用 execute
// service.invokeAll(..);
// 执行10s
Thread.sleep(2 * 1000);
producer1.stop(); // 生产者退出
producer2.stop(); // 生产者退出
producer3.stop(); // 生产者退出
System.out.println("所有生产者退出");
Thread.sleep(2000);
// 退出Executor
service.shutdown(); // 不可以再submit execute新的task,已经submit的将继续执行。
System.out.println("关闭线程池");
}
}
import java.util.Random;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
/**
* 作者:杨剑飞.
* 邮箱:847564732@qq.com
* 创建日期:2016-3-7
* 功能:生产者线程
*/
public class Producer implements Runnable {
private boolean isRunning = true;
private volatile BlockingQueue queue;
private static AtomicInteger count = new AtomicInteger();
private static final int DEFAULT_RANGE_FOR_SLEEP = 1000;
public Producer(BlockingQueue queue) {
this.queue = queue;
}
public void run() {
String data = null;
Random r = new Random();
System.out.println(Thread.currentThread().getName() + "_启动生产者线程!");
try {
while (isRunning) {
System.out.println(Thread.currentThread().getName()
+ "_正在生产数据...");
Thread.sleep(r.nextInt(DEFAULT_RANGE_FOR_SLEEP));
data = "data:" + count.incrementAndGet();
//将数据加入到 BlockingQueue 中
if (queue.offer(data, 5, TimeUnit.SECONDS)) {
System.out.println(Thread.currentThread().getName()
+ "_将数据:" + data + "放入队列...");
} else {
System.out.println(Thread.currentThread().getName()
+ "_放入数据失败:" + data);
}
}
} catch (InterruptedException e) {
e.printStackTrace();
Thread.currentThread().interrupt();
} finally {
System.out.println(Thread.currentThread().getName() + "_退出生产者线程!");
}
}
public void stop() {
isRunning = false;
}
}
import java.util.Random;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.TimeUnit;
/**
* 作者:杨剑飞.
* 邮箱:847564732@qq.com
* 创建日期:2016-3-7
* 功能:消费者线程
*/
public class Consumer implements Runnable {
private volatile BlockingQueue<String> queue;
private static final int DEFAULT_RANGE_FOR_SLEEP = 1000;
public Consumer(BlockingQueue<String> queue) {
this.queue = queue;
}
public void run() {
System.out.println(Thread.currentThread().getName() + "_启动消费者线程!");
Random r = new Random();
boolean isRunning = true;
try {
while (isRunning) {
System.out.println(Thread.currentThread().getName()
+ "\t\t_正从队列获取数据...");
String data;
data = queue.poll(2, TimeUnit.SECONDS);
// data = queue.take();
if (null != data) {
System.out.println(Thread.currentThread().getName()
+ "\t\t_从队列中取出产品,正在消费数据:" + data);
Thread.sleep(r.nextInt(DEFAULT_RANGE_FOR_SLEEP));
System.out.println(Thread.currentThread().getName()
+ "\t\t_完成消费数据:" + data);
} else {
// 超过2s还没数据,认为所有生产线程都已经退出,自动退出消费线程。
isRunning = false;
}
}
} catch (InterruptedException e) {
e.printStackTrace();
Thread.currentThread().interrupt();
} finally {
System.out.println(Thread.currentThread().getName()
+ "\t\t_退出消费者线程!");
}
}
}
Log
Thread-1_启动生产者线程!
Thread-0_启动生产者线程!
Thread-2_启动生产者线程!
Thread-3_启动消费者线程!
Thread-2_正在生产数据...
Thread-0_正在生产数据...
Thread-1_正在生产数据...
Thread-3 _正从队列获取数据...
Thread-0_将数据:data:1放入队列...
Thread-3 _从队列中取出产品,正在消费数据:data:1
Thread-0_正在生产数据...
Thread-1_将数据:data:2放入队列...
Thread-1_正在生产数据...
Thread-0_将数据:data:3放入队列...
Thread-0_正在生产数据...
Thread-1_将数据:data:4放入队列...
Thread-1_正在生产数据...
Thread-3 _完成消费数据:data:1
Thread-3 _正从队列获取数据...
Thread-3 _从队列中取出产品,正在消费数据:data:2
Thread-2_将数据:data:5放入队列...
Thread-2_正在生产数据...
Thread-1_将数据:data:6放入队列...
Thread-1_正在生产数据...
Thread-0_将数据:data:7放入队列...
Thread-0_正在生产数据...
Thread-0_将数据:data:8放入队列...
Thread-0_正在生产数据...
Thread-1_将数据:data:9放入队列...
Thread-1_正在生产数据...
Thread-2_将数据:data:10放入队列...
Thread-2_正在生产数据...
Thread-2_将数据:data:11放入队列...
Thread-2_正在生产数据...
Thread-2_将数据:data:12放入队列...
Thread-2_正在生产数据...
Thread-3 _完成消费数据:data:2
Thread-3 _正从队列获取数据...
Thread-3 _从队列中取出产品,正在消费数据:data:3
Thread-2_将数据:data:13放入队列...
Thread-2_正在生产数据...
所有生产者退出
Thread-3 _完成消费数据:data:3
Thread-3 _正从队列获取数据...
Thread-3 _从队列中取出产品,正在消费数据:data:4
Thread-0_将数据:data:14放入队列...
Thread-0_退出生产者线程!
Thread-3 _完成消费数据:data:4
Thread-3 _正从队列获取数据...
Thread-3 _从队列中取出产品,正在消费数据:data:5
Thread-1_将数据:data:15放入队列...
Thread-1_退出生产者线程!
Thread-3 _完成消费数据:data:5
Thread-3 _正从队列获取数据...
Thread-3 _从队列中取出产品,正在消费数据:data:6
Thread-2_将数据:data:16放入队列...
Thread-2_退出生产者线程!
关闭线程池
Thread-3 _完成消费数据:data:6
Thread-3 _正从队列获取数据...
Thread-3 _从队列中取出产品,正在消费数据:data:7
Thread-3 _完成消费数据:data:7
Thread-3 _正从队列获取数据...
Thread-3 _从队列中取出产品,正在消费数据:data:8
Thread-3 _完成消费数据:data:8
Thread-3 _正从队列获取数据...
Thread-3 _从队列中取出产品,正在消费数据:data:9
Thread-3 _完成消费数据:data:9
Thread-3 _正从队列获取数据...
Thread-3 _从队列中取出产品,正在消费数据:data:10
Thread-3 _完成消费数据:data:10
Thread-3 _正从队列获取数据...
Thread-3 _从队列中取出产品,正在消费数据:data:11
Thread-3 _完成消费数据:data:11
Thread-3 _正从队列获取数据...
Thread-3 _从队列中取出产品,正在消费数据:data:12
Thread-3 _完成消费数据:data:12
Thread-3 _正从队列获取数据...
Thread-3 _从队列中取出产品,正在消费数据:data:13
Thread-3 _完成消费数据:data:13
Thread-3 _正从队列获取数据...
Thread-3 _从队列中取出产品,正在消费数据:data:14
Thread-3 _完成消费数据:data:14
Thread-3 _正从队列获取数据...
Thread-3 _从队列中取出产品,正在消费数据:data:15
Thread-3 _完成消费数据:data:15
Thread-3 _正从队列获取数据...
Thread-3 _从队列中取出产品,正在消费数据:data:16
Thread-3 _完成消费数据:data:16
Thread-3 _正从队列获取数据...
Thread-3 _退出消费者线程!
PriorityQueue 实现
public class ProducerAndConsumer {
private int queueSize = 10;
private PriorityQueue<Integer> queue = new PriorityQueue<Integer>(queueSize);
public static void main(String[] args) {
ProducerAndConsumer producerAndConsumer = new ProducerAndConsumer();
Consumer consumer = producerAndConsumer.new Consumer();
Producer producer = producerAndConsumer.new Producer();
producer.start();
consumer.start();
}
class Producer extends Thread {
@Override
public void run() {
while (true) {
synchronized (queue) {
if (queue.size() == queueSize) {
try {
System.out.println("产品队列满,等待消费");
queue.wait();
} catch (InterruptedException e) {
queue.notify();
e.printStackTrace();
}
}
queue.offer(1);
queue.notifyAll();
System.out.println("生产一个新产品,还有" + (queueSize - queue.size()) + "个空位,现有 "+queue.size()+" 个产品");
try {
//模拟生产的耗时操作
Thread.sleep(25);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}
}
class Consumer extends Thread {
@Override
public void run() {
while (true) {
synchronized (queue) {
if (queue.size() == 0) {
try {
System.out.println("产品队列空,等待生产");
queue.wait();
} catch (InterruptedException e) {
queue.notify();
e.printStackTrace();
}
}
queue.poll();
queue.notify();
System.out.println("消耗一个新产品,还有" + queue.size() + "个产品,现有 "+queue.size()+" 个产品");
try {
//模拟消费的耗时操作
Thread.sleep(50);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}
}
}
log
生产一个新产品,还有9个空位,现有 1 个产品
生产一个新产品,还有8个空位,现有 2 个产品
消耗一个新产品,还有1个产品,现有 1 个产品
消耗一个新产品,还有0个产品,现有 0 个产品
产品队列空,等待生产
生产一个新产品,还有9个空位,现有 1 个产品
消耗一个新产品,还有0个产品,现有 0 个产品
产品队列空,等待生产
生产一个新产品,还有9个空位,现有 1 个产品
生产一个新产品,还有8个空位,现有 2 个产品
生产一个新产品,还有7个空位,现有 3 个产品
生产一个新产品,还有6个空位,现有 4 个产品
消耗一个新产品,还有3个产品,现有 3 个产品
消耗一个新产品,还有2个产品,现有 2 个产品
生产一个新产品,还有7个空位,现有 3 个产品
生产一个新产品,还有6个空位,现有 4 个产品
生产一个新产品,还有5个空位,现有 5 个产品
消耗一个新产品,还有4个产品,现有 4 个产品
生产一个新产品,还有5个空位,现有 5 个产品
消耗一个新产品,还有4个产品,现有 4 个产品
消耗一个新产品,还有3个产品,现有 3 个产品
生产一个新产品,还有6个空位,现有 4 个产品
生产一个新产品,还有5个空位,现有 5 个产品
生产一个新产品,还有4个空位,现有 6 个产品
生产一个新产品,还有3个空位,现有 7 个产品
消耗一个新产品,还有6个产品,现有 6 个产品
消耗一个新产品,还有5个产品,现有 5 个产品
消耗一个新产品,还有4个产品,现有 4 个产品
生产一个新产品,还有5个空位,现有 5 个产品
生产一个新产品,还有4个空位,现有 6 个产品
生产一个新产品,还有3个空位,现有 7 个产品
生产一个新产品,还有2个空位,现有 8 个产品
生产一个新产品,还有1个空位,现有 9 个产品
生产一个新产品,还有0个空位,现有 10 个产品
消耗一个新产品,还有9个产品,现有 9 个产品
生产一个新产品,还有0个空位,现有 10 个产品
产品队列满,等待消费
消耗一个新产品,还有9个产品,现有 9 个产品
生产一个新产品,还有0个空位,现有 10 个产品
产品队列满,等待消费
消耗一个新产品,还有9个产品,现有 9 个产品
生产一个新产品,还有0个空位,现有 10 个产品
产品队列满,等待消费
消耗一个新产品,还有9个产品,现有 9 个产品
消耗一个新产品,还有8个产品,现有 8 个产品
消耗一个新产品,还有7个产品,现有 7 个产品
消耗一个新产品,还有6个产品,现有 6 个产品
生产一个新产品,还有3个空位,现有 7 个产品
生产一个新产品,还有2个空位,现有 8 个产品
消耗一个新产品,还有7个产品,现有 7 个产品
消耗一个新产品,还有6个产品,现有 6 个产品
消耗一个新产品,还有5个产品,现有 5 个产品
消耗一个新产品,还有4个产品,现有 4 个产品
生产一个新产品,还有5个空位,现有 5 个产品
生产一个新产品,还有4个空位,现有 6 个产品
生产一个新产品,还有3个空位,现有 7 个产品
消耗一个新产品,还有6个产品,现有 6 个产品
消耗一个新产品,还有5个产品,现有 5 个产品
消耗一个新产品,还有4个产品,现有 4 个产品
生产一个新产品,还有5个空位,现有 5 个产品
生产一个新产品,还有4个空位,现有 6 个产品
生产一个新产品,还有3个空位,现有 7 个产品
生产一个新产品,还有2个空位,现有 8 个产品
消耗一个新产品,还有7个产品,现有 7 个产品
消耗一个新产品,还有6个产品,现有 6 个产品
消耗一个新产品,还有5个产品,现有 5 个产品
消耗一个新产品,还有4个产品,现有 4 个产品
消耗一个新产品,还有3个产品,现有 3 个产品
生产一个新产品,还有6个空位,现有 4 个产品
生产一个新产品,还有5个空位,现有 5 个产品
生产一个新产品,还有4个空位,现有 6 个产品
生产一个新产品,还有3个空位,现有 7 个产品
生产一个新产品,还有2个空位,现有 8 个产品
生产一个新产品,还有1个空位,现有 9 个产品
生产一个新产品,还有0个空位,现有 10 个产品
产品队列满,等待消费
消耗一个新产品,还有9个产品,现有 9 个产品
生产一个新产品,还有0个空位,现有 10 个产品
产品队列满,等待消费
消耗一个新产品,还有9个产品,现有 9 个产品
生产一个新产品,还有0个空位,现有 10 个产品
产品队列满,等待消费
消耗一个新产品,还有9个产品,现有 9 个产品
消耗一个新产品,还有8个产品,现有 8 个产品
消耗一个新产品,还有7个产品,现有 7 个产品
消耗一个新产品,还有6个产品,现有 6 个产品
消耗一个新产品,还有5个产品,现有 5 个产品
消耗一个新产品,还有4个产品,现有 4 个产品
消耗一个新产品,还有3个产品,现有 3 个产品
消耗一个新产品,还有2个产品,现有 2 个产品
消耗一个新产品,还有1个产品,现有 1 个产品
消耗一个新产品,还有0个产品,现有 0 个产品
产品队列空,等待生产
生产一个新产品,还有9个空位,现有 1 个产品
生产一个新产品,还有8个空位,现有 2 个产品
生产一个新产品,还有7个空位,现有 3 个产品
消耗一个新产品,还有2个产品,现有 2 个产品
消耗一个新产品,还有1个产品,现有 1 个产品
消耗一个新产品,还有0个产品,现有 0 个产品
产品队列空,等待生产
生产一个新产品,还有9个空位,现有 1 个产品
生产一个新产品,还有8个空位,现有 2 个产品
消耗一个新产品,还有1个产品,现有 1 个产品
消耗一个新产品,还有0个产品,现有 0 个产品
产品队列空,等待生产
生产一个新产品,还有9个空位,现有 1 个产品
生产一个新产品,还有8个空位,现有 2 个产品
生产一个新产品,还有7个空位,现有 3 个产品
消耗一个新产品,还有2个产品,现有 2 个产品
消耗一个新产品,还有1个产品,现有 1 个产品
。。。