生产者-消费者模型是并发编程中的经典案例,它展示了如何通过多线程协调来处理共享资源。本文将实现一个基于阻塞队列的生产者-消费者模型,展示Java多线程编程的核心概念。

完整代码实现
java
复制代码
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.atomic.AtomicInteger;

/**

  • 生产者-消费者模型演示

  • 使用BlockingQueue作为线程安全的缓冲区
    */
    public class ProducerConsumerDemo {

    public static void main(String[] args) {
    // 创建容量为5的阻塞队列
    BlockingQueue queue = new ArrayBlockingQueue<>(5);

     // 创建生产者和消费者
     Producer producer = new Producer(queue);
     Consumer consumer = new Consumer(queue);
     
     // 启动线程
     new Thread(producer, "生产者线程").start();
     new Thread(consumer, "消费者线程").start();
    

    }
    }

/**

  • 生产者类
    */
    class Producer implements Runnable {
    private final BlockingQueue queue;
    private final AtomicInteger count = new AtomicInteger(0);
    private volatile boolean isRunning = true;

    public Producer(BlockingQueue queue) {
    this.queue = queue;
    }

    @Override
    public void run() {
    try {
    while (isRunning) {
    String data = “产品-” + count.incrementAndGet();
    // 将产品放入队列,如果队列满则阻塞
    queue.put(data);
    System.out.println(Thread.currentThread().getName() + " 生产: " + data);
    Thread.sleep(100); // 模拟生产时间
    }
    } catch (InterruptedException e) {
    Thread.currentThread().interrupt();
    System.out.println(“生产者被中断”);
    }
    }

    public void stop() {
    isRunning = false;
    }
    }

/**

  • 消费者类
    */
    class Consumer implements Runnable {
    private final BlockingQueue queue;
    private volatile boolean isRunning = true;

    public Consumer(BlockingQueue queue) {
    this.queue = queue;
    }

    @Override
    public void run() {
    try {
    while (isRunning) {
    // 从队列中取出产品,如果队列空则阻塞
    String data = queue.take();
    System.out.println(Thread.currentThread().getName() + " 消费: " + data);
    Thread.sleep(150); // 模拟消费时间
    }
    } catch (InterruptedException e) {
    Thread.currentThread().interrupt();
    System.out.println(“消费者被中断”);
    }
    }

    public void stop() {
    isRunning = false;
    }
    }
    关键特性解析

  1. 线程安全的数据结构
    使用ArrayBlockingQueue作为缓冲区,它内部实现了线程安全的put和take操作:

put(): 当队列满时自动阻塞
take(): 当队列空时自动阻塞
2. 优雅的线程终止
通过volatile boolean isRunning标志位实现可控的线程终止,避免使用已废弃的stop()方法。

  1. 原子操作
    使用AtomicInteger保证产品编号的原子性递增,避免使用同步锁。

  2. 异常处理
    正确处理InterruptedException,保持线程的中断状态。

运行结果示例
复制代码
生产者线程 生产: 产品-1
消费者线程 消费: 产品-1
生产者线程 生产: 产品-2
消费者线程 消费: 产品-2
生产者线程 生产: 产品-3

扩展建议
多生产者和多消费者:可以创建多个生产者和消费者线程来测试更复杂的场景
线程池管理:使用ExecutorService来管理线程生命周期
性能监控:添加统计信息来监控生产消费速率
优雅关闭:实现ShutdownHook来确保程序正常退出

更多推荐