惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

W
WeLiveSecurity
Jina AI
Jina AI
博客园 - 司徒正美
雷峰网
雷峰网
宝玉的分享
宝玉的分享
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
博客园_首页
WordPress大学
WordPress大学
Google DeepMind News
Google DeepMind News
GbyAI
GbyAI
MyScale Blog
MyScale Blog
Apple Machine Learning Research
Apple Machine Learning Research
美团技术团队
I
InfoQ
博客园 - Franky
CTFtime.org: upcoming CTF events
CTFtime.org: upcoming CTF events
博客园 - 叶小钗
阮一峰的网络日志
阮一峰的网络日志
Cyberwarzone
Cyberwarzone
C
CXSECURITY Database RSS Feed - CXSecurity.com
S
Schneier on Security
P
Privacy & Cybersecurity Law Blog
T
Threatpost
Cloudbric
Cloudbric
D
Docker
M
MIT News - Artificial intelligence
Recent Commits to openclaw:main
Recent Commits to openclaw:main
Vercel News
Vercel News
Martin Fowler
Martin Fowler
J
Java Code Geeks
AWS News Blog
AWS News Blog
The Cloudflare Blog
酷 壳 – CoolShell
酷 壳 – CoolShell
L
Lohrmann on Cybersecurity
Hacker News: Ask HN
Hacker News: Ask HN
Last Week in AI
Last Week in AI
S
Security @ Cisco Blogs
Help Net Security
Help Net Security
C
Cisco Blogs
V
V2EX
博客园 - 【当耐特】
I
Intezer
爱范儿
爱范儿
F
Fortinet All Blogs
cs.AI updates on arXiv.org
cs.AI updates on arXiv.org
P
Privacy International News Feed
IT之家
IT之家
L
LINUX DO - 最新话题
B
Blog RSS Feed
K
KPMG report finds enterprise disconnect between AI and its ROI | CIO

小松鼠的博客

记录一次线上k8s工作节点无法创建容器的问题排查思路与解决办法 记一次线上GoLang项目OOM排查过程 从LastPass转向拥抱开源KeePass的心路历程 故障定位与 AI 结合前后端编码实践 FileBeat收集nginx-ingress-controller日志 K8s云原生环境下文件描述符占用过高查询思路 2024年最新关闭火绒安全工具的开机自启方法 Kubernetes任务调度实践-Go语言实现Job和CronJob对比分析 离线更新k8s环境下的trivy漏洞库方法 使用Go语言接入Choerodon实现基于OAuth2的统一身份认证登录 在Vue2中自定义Switch组件并实现父子组件双向数据绑定 关于docker jdk1.8镜像中的GB18030-2022标准支持及验证 Go框架gin中的session存储gin-contrib-sessions和go-session 关于修改node_module中的源码问题记录 docker-compose网络和内网服务IP冲突问题 慎用存储过程:一条语句引发的数据库存储100%占用 Spring Boot中4种文件下载方法的实现 避坑-不能将specific类型的gitlab-runner改变为share类型 Docker compose中的MySQL主从复制模式和percona-toolkit工具使用 在minio中开启https访问以及使用rclone备份minio桶 在多机Docker环境下部署Choerodon的解决方案 Prometheus中Monitor添加对SpringBoot Actuator的Basic认证 在Nginx的容器镜像中隐藏Nginx的Server响应头 K8s中的两种nginx-ingress-controller及其区别 两个docker工具:runlike和whaler Grafana中的邮件报警和截图插件grafana-image-enderer K8s中externalName-service和services-without-selectors maven配置文件settings.xml中的一些概念总结 K8s中flexvolume插件驱动的安装 K8s中的coredns无法解析svc问题排查 K8s中使用Ingress访问请求体过大问题解决 关于k8s中对于SpringBoot应用TCP类型的就绪探针不准确的问题发现 K8s中的环境变量与应用程序的对应关系与操作 SpringMVC4升级为SpringBoot2实战 在Vmware中Ubuntu22.04的vm-tools和网络问题 修改k8s节点主机名并重新加入集群 离线安装Grafana插件 Spring Data Jpa 中使用CriteriaBuilder动态拼接SQL 在SpringBoot项目配置Liquibase数据库版本管理 记录Vue中父子组件传值的实战应用 实现单例模式的8种方法 使用两个线程交替打印0-100的奇偶数 关于部署于JBoss5中的Spring应用获取项目真实部署路径的问题 获取下一个完全对称日 通过短信验证码验证修改密码的解决方案 在Win10中使用Win+R快速启动软件 使用RSA加解密时注意Cipher.getInstance(String var0,Provider var1)提供的Provider是否正确 在RestEasy2.x中解决接口重复提交问题 几道简单的CTF题目思路 重温Spring---Spring事务控制与基于XML和注解的配置方法 重温Spring---Spring AOP基于XML和注解的配置 重温Spring---AOP动态代理和Spring AOP及其基本原理 重温Spring---Spring IOC基于XML和注解的配置和比较 在Windows10中安装MySQL5.7 Zip版本及常用配置 重温Spring---使用Spring IOC解决程序耦合 策略模式与责任链模式实战应用 Linux上直接打开war包修改文件 在Windows上运行两个微信的简单脚本 ThreadPoolExecutor的使用方法与分页查询数据实例 IDEA中Shelve Changes 和 Git Stash 通过resteasy发布RESTful接口 解决前端请求后台接口,后台报错Can not deserialize instance of java.util.ArrayList out of START_OBJECT token 使用VBA脚本汇总Excel文档 使用Jenkins+GitLab实现自动部署vue项目 Kubernetes:使用hostPath挂载nginx集群的配置文件和html 彻底搞定VirtualBox虚拟机的网络设定 在Docker中安装MySQL5.7并开启远程访问(附授权和修改密码方式) 利用git命令和java文件流 获取自己改动过的文件 浅谈Spring定时任务的使用(Scheduled注解) 在Spring项目简单配置Flyway(V4.2版本)数据库版本管理 解决Spring单元测试中因外键关联导致的失败integrity constraint violation:foreign key no action Redis安装与哨兵模式配置入门 关于Vue中使用Element-UI样式row-class-name失效的问题 Element-UI中实现可动态增加行列和可编辑单元格的表格 Windows系统查看端口占用、结束进程方法和命令 层次分析法(AHP)分析步骤与计算方法 源码分析之解决layui框架重载表格时额外参数不清空的问题 Spring Data Jpa 返回自定义对象(实体部分属性、多表联查) 如何将一个jar放到本地maven仓库中 关于SSM项目停止Tomcat时Log4j出现java.lang.NoClassDefFoundError: 获取el-table单元格值并根据该值对元素自定义样式渲染 解决Git每次push都要重新输入账号密码和HttpRequestException encountered的问题 解决前后端分离项目中Vue不带cookies的问题 SSM集成Shiro自定义权限过滤器不执行解决方案 SSM集成Shiro不进入自定义Realm的doGetAuthorizationInfo的解决方案 Vue+SSM中使用Token验证登录 Git拉代码推送代码提示密码错误如何修改 Git配置SSH Key(Git配置多个账户) 安装Tomcat服务器以及错误汇总(tomcat8.0、jdk8) 关于我
三种常用的生产者消费者模式实现
ycyin · 2022-01-12 · via 小松鼠的博客

三种常用的生产者消费者模式实现

2022年1月12日大约 8 分钟设计模式生产者与消费者模式软件设计


本篇是在学习5.Thread和Object中线程相关的重要方法 (ycyin.eu.org)时对notify()wait()的相关用法记录。本篇除代码外多处引用网上文字,具体出处见文末参考。

生产者消费者模式并不是GOF提出的23种设计模式之一,23种设计模式都是建立在面向对象的基础之上的,但其实面向过程的编程中也有很多高效的编程模式,生产者消费者模式便是其中之一,它是我们编程过程中最常用的一种设计模式。

在实际的软件开发过程中,经常会碰到如下场景:某个模块负责产生数据,这些数据由另一个模块来负责处理(此处的模块是广义的,可以是类、函数、线程、进程等)。产生数据的模块,就形象地称为生产者;而处理数据的模块,就称为消费者。

单单抽象出生产者和消费者,还不算是生产者/消费者模式。该模式还需要有一个缓冲区处于生产者和消费者之间,作为一个中介。

实现消费者生产者模式有很多种方式,可以在文末参考中找到,我认为常用的有BlockingQueueconditionwait/notify三种实现方式,虽说是三种实现方式但是本次都是差不多的。

下面的例子,使用三种方法实现生产者消费者模式,完成生产者生产100条数据,消费者消费100条数据。

方式一:用BlockingQueue 实现

使用这种方式,只需要使用ArrayBlockingQueue类型的BlockingQueue命名为queue,并指定固定容量。然后生产者使用queue.put()负责往队列添加数据,消费者使用queue.take()负责消费数据。

public class ConsumerProducerWithBlockingQueue {

    public static void main(String[] args) {
        ConsumerProducerWithBlockingQueue cpwb = new ConsumerProducerWithBlockingQueue();
        ArrayBlockingQueue<Object> blockingQueue = new ArrayBlockingQueue<>(10);
        Producer producer = cpwb.new Producer(blockingQueue);

        Consumer consumer = cpwb.new Consumer(blockingQueue);

        producer.start();
        consumer.start();
    }

    class Consumer extends Thread {
        private BlockingQueue<Object> queue;
        public Consumer(BlockingQueue queue) {
            this.queue = queue;
        }

        @Override
        public void run() {
            try {
                for (int i = 0; i <= 100; i++) {
                    System.out.println("消费者取出:" + queue.take()+"仓库还有"+queue.size());
                }
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }
    }

    class Producer extends Thread {
        private BlockingQueue<Object> queue;
        public Producer(BlockingQueue queue) {
            this.queue = queue;
        }

        @Override
        public void run() {
            try {
                for (int i = 0; i <= 100; i++) {
                    System.out.println("生产者放入:"+i);
                    queue.put(i);
                }
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }
    }
}

虽然代码非常简单,但实际上ArrayBlockingQueue已经在背后完成了很多工作,比如队列满了就去阻塞生产者线程,队列有空就去唤醒生产者线程等。比如从ArrayBlockingQueuetake()源码中可以看出,它为我们使用Condition 实现的方式实现了。

public E take() throws InterruptedException {
    final ReentrantLock lock = this.lock;
    lock.lockInterruptibly();
    try {
        while (count == 0)
            notEmpty.await();
        return dequeue();
    } finally {
        lock.unlock();
    }
}

方式二:用 Condition 实现

利用 Condition 实现生产者消费者模式,与BlockingQueue背后的实现原理非常相似,相当于我们自己实现一个简易版的 BlockingQueue:

定义了一个队列变量 queue 并设置最大容量为 10;定义了一个 ReentrantLock 类型的 Lock 锁,并在 Lock 锁的基础上创建两个 Condition,一个是 notEmpty,另一个是 notFull,分别代表队列没有空和没有满的条件;声明put 和 take 这两个核心方法。

public class ConsumerProducerWithCondition {

    private Queue queue;

    private int maxSize = 10;

    private ReentrantLock lock = new ReentrantLock();

    private Condition notFull = lock.newCondition();

    private Condition notEmpty = lock.newCondition();

    public ConsumerProducerWithCondition() {
        this.queue = new LinkedList();
    }

    public void put(Object obj) throws InterruptedException {
        lock.lock();
        try {
            while (queue.size() == maxSize) {
                notFull.await(); //不要写成了 .wait()
            }
            queue.add(obj);
            notEmpty.signalAll();
        }finally {
            lock.unlock();
        }
    }

    public Object take() throws InterruptedException {
        lock.lock();
        try {
            while (queue.isEmpty()) {
                notEmpty.await();
            }
            Object item = queue.remove();
            notFull.signalAll();
            return item;
        }finally {
            lock.unlock();
        }
    }


    public static void main(String[] args) {
        ConsumerProducerWithCondition consumerProducer = new ConsumerProducerWithCondition();

        ConsumerProducerWithCondition.Consumer consumer = consumerProducer.new Consumer(consumerProducer);

        ConsumerProducerWithCondition.Producer producer = consumerProducer.new Producer(consumerProducer);

        producer.start();
        consumer.start();
    }

    class Consumer extends Thread {

        private ConsumerProducerWithCondition consumerProducer;

        public Consumer(ConsumerProducerWithCondition consumerProducer) {
            this.consumerProducer = consumerProducer;
        }

        @Override
        public void run() {
            try {
                for (int i = 0; i < 100; i++) {
                    Object take = consumerProducer.take();
                    System.out.println("消费者取出:" + take+"仓库还有"+queue.size());
                }
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }
    }

    class Producer extends Thread {

        private ConsumerProducerWithCondition consumerProducer;


        public Producer(ConsumerProducerWithCondition consumerProducer) {
            this.consumerProducer = consumerProducer;
        }

        @Override
        public void run() {
            try {
                for (int i = 0; i < 100; i++) {
                    System.out.println("生产者放入:"+i);
                    consumerProducer.put(i);
                }
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }
    }
}

因为生产者消费者模式通常是面对多线程的场景,需要一定的同步措施保障线程安全,所以在 put 方法中先将 Lock 锁上,然后,在 while 的条件里检测 queue 是不是已经满了,如果已经满了,则调用 notFull 的 await() 阻塞生产者线程并释放 Lock,如果没有满,则往队列放入数据并利用 notEmpty.signalAll() 通知正在等待的所有消费者并唤醒它们。最后在 finally 中利用 lock.unlock() 方法解锁,把 unlock 方法放在 finally 中是一个基本原则,否则可能会产生无法释放锁的情况。

take 方法实际上是与 put 方法相互对应的,同样是通过 while 检查队列是否为空,如果为空,消费者开始等待,如果不为空则从队列中获取数据并通知生产者队列有空余位置,最后在 finally 中解锁。

**注意:**这里需要注意判断队列是否为空、是否为满的状态需要使用while来判断,和下面的方式三一样。这在oracle官方Object对象的wait()方法有说明:As in the one argument version, interrupts and spurious wakeups are possible, and this method should always be used in a loop。

为什么在take()方法中使用while( queue.size() == 0 ) 判断而不是用if来判断?

思考这样一种情况,因为生产者消费者往往是多线程的,我们假设有两个消费者,第一个消费者线程获取数据时,发现队列为空,便进入等待状态;因为第一个线程在等待时会释放 Lock 锁,所以第二个消费者可以进入并执行 if( queue.size() == 0 ),也发现队列为空,于是第二个线程也进入等待;而此时,如果生产者生产了一个数据,便会唤醒两个消费者线程,而两个线程中只有一个线程可以拿到锁,并执行 queue.remove 操作,另外一个线程因为没有拿到锁而卡在被唤醒的地方,而第一个线程执行完操作后会在 finally 中通过 unlock 解锁,而此时第二个线程便可以拿到被第一个线程释放的锁,继续执行操作,也会去调用 queue.remove 操作,然而这个时候队列已经为空了,所以会抛出 NoSuchElementException 异常,这不符合我们的逻辑。而如果用 while 做检查,当第一个消费者被唤醒得到锁并移除数据之后,第二个线程在执行 remove 前仍会进行 while 检查,发现此时依然满足 queue.size() == 0 的条件,就会继续执行 await 方法,避免了获取的数据为 null 或抛出异常的情况。

方式三:用 wait/notify 实现

使用 wait/notify 实现生产者消费者模式的方法,实际上实现原理和前面两种是非常类似的。这为我们理解wait和notify提供了很大的帮助。

最主要的部分仍是 take 与 put 方法,put 方法被 synchronized 保护,while检查队列是否为满,如果不满就往里放入数据并通过notify()notifyAll()唤醒其他线程。同样,take方法也被synchronized修饰,while检查队列是否为空,如果不为空就获取数据并唤醒其他线程。

public class ConsumerProducerWithWaitNotify {
    public static void main(String[] args) {
        ConsumerProducerWithWaitNotify consumerProducer = new ConsumerProducerWithWaitNotify();

        EventStorge eventStorge = consumerProducer.new EventStorge();

        Consumer consumer = consumerProducer.new Consumer(eventStorge);

        Producer producer = consumerProducer.new Producer(eventStorge);

        producer.start();
        consumer.start();
    }

    class Consumer extends Thread {
        private EventStorge storge;
        public Consumer(EventStorge storge) {
            this.storge = storge;
        }

        @Override
        public void run() {
            try {
                for (int i = 0; i < 100; i++) {
                    storge.take();
                }
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }
    }

    class Producer extends Thread {
        private EventStorge storge;
        public Producer(EventStorge storge) {
            this.storge = storge;
        }

        @Override
        public void run() {
            try {
                for (int i = 0; i < 100; i++) {
                    storge.put(i);
                }
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }
    }

    class EventStorge {
        private LinkedList<Integer> list=null;
        private int maxSize = 0;

        public EventStorge() {
            this.list = new LinkedList<>();
            this.maxSize = 10;
        }

        public synchronized void put(int num) throws InterruptedException {
            while (list.size() == maxSize) { // not if
                wait();
            }
            System.out.println("生产者放入"+num+"---");
            list.add(num);
            notify();
        }

        public synchronized void take() throws InterruptedException {
            while (list.isEmpty()) {
                wait();
            }
            notify();
            System.out.println("消费者取出 " + list.poll() +"还有" + list.size()+"个");
        }

    }

}

总结

这三种方式实现生产者消费者模式,从编码难度上来说第一种使用BlockingQueue实现是最简单的,但实际上其底层的逻辑在方式二、方式三中得以体现。这三种 方式均需要一个队列,然后要有锁来控制生产者和消费者。

除此三种方式外,还有使用信号量(数据库连接池)、管道流(只适合两个线程间通信)等方式实现生产者消费者模式,我觉得都是差不多的原理。

参考

  1. 生产者/消费者模式的理解及实现(整理) - Luego - 博客园 (cnblogs.com)
  2. Java多种方式解决生产者消费者问题(十分详细)_爱你の大表哥的博客-CSDN博客_java生产者消费者
  3. Java并发编程(六)——三种方式实现生产者消费者模式_易水寒的博客-CSDN博客
  4. java多线程wait时为什么要用while而不是if_worldchinalee的博客-CSDN博客_java wait while
  5. Object (Java Platform SE 8 ) (oracle.com)