从源码全面解析LinkedBlockingQueue的来龙去脉

news2025/1/10 20:34:57

一、引言

并发编程在互联网技术使用如此广泛,几乎所有的后端技术面试官都要在并发编程的使用和原理方面对小伙伴们进行 360° 的刁难。

二、使用

对于阻塞队列,想必大家应该都不陌生,我们这里简单的介绍一下,对于 Java 里面的阻塞队列,其使用了 生产者和消费者 的模型

对于生产者来说,主要有以下几部分:

add(E)     	// 添加数据到队列,如果队列满了,无法存储,抛出异常
offer(E)    // 添加数据到队列,如果队列满了,返回false
offer(E,timeout,unit)   // 添加数据到队列,如果队列满了,阻塞timeout时间,如果阻塞一段时间,依然没添加进入,返回false
put(E)      // 添加数据到队列,如果队列满了,挂起线程,等到队列中有位置,再扔数据进去,死等!
复制代码

对于消费者来说,主要有以下几部分:

remove()    // 从队列中移除数据,如果队列为空,抛出异常
poll()      // 从队列中移除数据,如果队列为空,返回null,么的数据
poll(timeout,unit)   // 从队列中移除数据,如果队列为空,挂起线程timeout时间,等生产者扔数据,再获取
take()     // 从队列中移除数据,如果队列为空,线程挂起,一直等到生产者扔数据,再获取
复制代码

我们本篇来讲讲堵塞队列中的第二员猛将,LinkedBlockingQueue 的故事

我们先来看其基本使用

public class LinkedBlockingQueueTest {
    public static void main(String[] args) throws Exception {
        LinkedBlockingQueue queue = new LinkedBlockingQueue();

        // 生产者扔数据
        queue.add("1");
        queue.offer("2");
        queue.offer("3", 2, TimeUnit.SECONDS);
        queue.put("2");

        // 消费者取数据
        System.out.println(queue.remove());
        System.out.println(queue.poll());
        System.out.println(queue.poll(2, TimeUnit.SECONDS));
        System.out.println(queue.take());
    }
}
复制代码

三、源码

1、初始化

由于我们的 LinkedBlockingQueue 底层是链表实现的,所以我们初始化的时候不需要指定其大小

LinkedBlockingQueue queue = new LinkedBlockingQueue();

// 如果我们不指定容量大小的话,这里的容量默认为Integer.MAX_VALUE
public LinkedBlockingQueue() {
    this(Integer.MAX_VALUE);
}

public LinkedBlockingQueue(int capacity) {
    // 如果容量传进来是小于等于0的,直接抛异常
    if (capacity <= 0){
        throw new IllegalArgumentException();
    }
    // 当前的容量赋值
    this.capacity = capacity;
    // 这里其实和我们的AQS有点像
    // 搞一个虚拟的头结点,减少后面的判空
    last = head = new Node<E>(null);
}
复制代码

当然,除了我们初始化的这些成员变量,我们还有一部分:

class Node<E> {
    // 当前的数据
    E item;
    // 指向下一个数据的指针
    Node<E> next;
    Node(E x) {
        item = x;
    }
}

// 当前链表中存在的数据数量
private final AtomicInteger count = new AtomicInteger();

// 读锁
private final ReentrantLock takeLock = new ReentrantLock();

// 唤醒消费者线程
private final Condition notEmpty = takeLock.newCondition();

// 写锁
private final ReentrantLock putLock = new ReentrantLock();

// 唤醒生产者线程
private final Condition notFull = putLock.newCondition();
复制代码

这里可能有的小伙伴有点懵逼,为什么这哥们(LinkedBlockingQueue)用了两个锁呢?为什么我 ArrayBlockingQueue 只能用一把锁?

不要急,我们慢慢的往下看他源码

2、生产者的源码

2.1 add()源码实现

public boolean add(E e) {
    return super.add(e);
}

// 走到这里会发现,我们的add方法就是调用了offer方法
// offer: 添加数据到队列,如果队列满了,返回false
// 所以这里offer满了,就会抛出异常:"Queue full"
public boolean add(E e) {
    if (offer(e))
        return true;
    else
        throw new IllegalStateException("Queue full");
}
复制代码

2.2 offer()源码实现

public boolean offer(E e) {
    // 如果是空值,直接抛出异常
    if (e == null) throw new NullPointerException();
    // 引用,上篇我们分析过
    final AtomicInteger count = this.count;
    // 判断当前数据量是否和我们总容量一样
    if (count.get() == capacity){
        return false;
    }
    // 标记位
    int c = -1;
    // 创建节点
    Node<E> node = new Node<E>(e);
    // 引用写锁
    final ReentrantLock putLock = this.putLock;
    // 上锁
    putLock.lock();
    try {
        // 如果当前数据量小于总容量
        // 这里我们上面也检查过,相当于DCL的意思
        if (count.get() < capacity) {
            // 插入队列
            enqueue(node);
            // 得到当前数据量
            // 这里需要注意:getAndIncrement先返回数据,再加一
            c = count.getAndIncrement();
            // 如果我们发现当前数据量还小于总容量
            // 也就是我们可以继续放数据
            if (c + 1 < capacity)
                // 唤醒其他的生产者线程扔数据
                // 当然这里稍微多说一点,这里的唤醒指的是将生产者从Condition队列放到AQS队列中
                // 具体什么时候执行还需要看AQS的调度
                notFull.signal();
        }
    } finally {
        // 解锁
        putLock.unlock();
    }
    // 如果我们当前数据量为0,代表队列中原来无数据
    // 但上面现在扔进去了一个
    if (c == 0)
        // 需要唤醒所有的消费者消费数据
        signalNotEmpty();
    return c >= 0;
}

private void enqueue(Node<E> node) {
    // 将当面节点挂在last节点后
    // 将last节点指向当前节点
    last = last.next = node;
}


// 这里我们的Condition聊过
// 必须持有当前锁资源才可以使用Condition的方法
private void signalNotEmpty() {
    // 拿到读锁
    final ReentrantLock takeLock = this.takeLock;
    // 加锁
    takeLock.lock();
    try {
        // 唤醒消费者线程
        notEmpty.signal();
    } finally {
        // 解锁
        takeLock.unlock();
    }
}
复制代码

2.3 offer(time)源码实现

public boolean offer(E e, long timeout, TimeUnit unit) throws InterruptedException {
    // 如果是空值,直接抛出异常
    if (e == null) throw new NullPointerException();
    // 转成统一的单位
    long nanos = unit.toNanos(timeout);
    int c = -1;
    // 写锁
    final ReentrantLock putLock = this.putLock;
    // 当前容量
    final AtomicInteger count = this.count;
    // 加锁
    putLock.lockInterruptibly();
    try {
        // 如果当前数据量小于总容量
        // 这里我们上面也检查过,相当于DCL的意思
        while (count.get() == capacity) {
            // 如果我们剩余时间小于0,直接失败即可
            if (nanos <= 0)
                return false;
            // 反之生产者线程写入挂起nanos时间
            nanos = notFull.awaitNanos(nanos);
        }
        // 添加至队列
        enqueue(new Node<E>(e));
        // 得到当前数据量
        // 这里需要注意:getAndIncrement先返回数据,再加一
        c = count.getAndIncrement();
        // 如果我们发现当前数据量还小于总容量
        // 也就是我们可以继续放数据
        if (c + 1 < capacity)
            // 唤醒其他的生产者线程扔数据
            // 当然这里稍微多说一点,这里的唤醒指的是将生产者从Condition队列放到AQS队列中
            // 具体什么时候执行还需要看AQS的调度
            notFull.signal();
    } finally {
        // 解锁
        putLock.unlock();
    }
    // 如果我们当前数据量为0,代表队列中原来无数据
    // 但现在扔进去了一个,唤醒消费者线程
    if (c == 0)
        signalNotEmpty();
    return true;
}
复制代码

2.4 put()源码实现

  • 这里就不写了,其实和我们的 offer 一样,大家自己看看就好
public void put(E e) throws InterruptedException {
    if (e == null) throw new NullPointerException();
    int c = -1;
    Node<E> node = new Node<E>(e);
    final ReentrantLock putLock = this.putLock;
    final AtomicInteger count = this.count;
    putLock.lockInterruptibly();
    try {
        while (count.get() == capacity) {
            notFull.await();
        }
        enqueue(node);
        c = count.getAndIncrement();
        if (c + 1 < capacity)
            notFull.signal();
    } finally {
        putLock.unlock();
    }
    if (c == 0)
        signalNotEmpty();
}
复制代码

3、消费者的源码

3.1 remove()源码实现

public E remove() {
    E x = poll();
    if (x != null)
        return x;
    else
        throw new NoSuchElementException();
}
复制代码

3.2 poll()源码实现

public E poll() {
    // 获取当前链表的数据量
    final AtomicInteger count = this.count;
    // 如果数据量为0,说明无数据
    // 消费者无法消费,直接返回null即可
    if (count.get() == 0)
        return null;
    E x = null;
    int c = -1;
    // 拿到读锁
    final ReentrantLock takeLock = this.takeLock;
    // 加锁
    takeLock.lock();
    try {
        // 如果数据量大于0,说明有数据
        // 这里我们上面也检查过,相当于DCL的意思
        if (count.get() > 0) {
            // 取数
            x = dequeue();
            // 得到当前数据量
            // 这里需要注意:getAndIncrement先返回数据,再减一
            c = count.getAndDecrement();
            // 如果我们的数据量大于1,则唤醒消费者来消费
            if (c > 1)
                notEmpty.signal();
        }
    } finally {
        // 解锁
        takeLock.unlock();
    }
    // 如果数据量等于当前的总容量
    // 说明当前的链表已经有空余了,唤醒生产者生产
    if (c == capacity)
        signalNotFull();
    return x;
}

// 这个取数据和我们的AQS有点像
// 去除当前数据并且将当前节点作为头结点
private E dequeue() {
    Node<E> h = head;
    Node<E> first = h.next;
    h.next = h; // help GC
    head = first;
    E x = first.item;
    first.item = null;
    return x;
}

private void signalNotFull() {
    // 拿到写锁
    final ReentrantLock putLock = this.putLock;
    // 上锁
    putLock.lock();
    try {
        // 唤醒生产者
        notFull.signal();
    } finally {
        // 解锁
        putLock.unlock();
    }
}
复制代码

3.3 poll(time)源码实现

public E poll(long timeout, TimeUnit unit) throws InterruptedException {
    E x = null;
    int c = -1;
    // 统一时间单位
    long nanos = unit.toNanos(timeout);
    // 拿到当前数据量 + 读锁
    final AtomicInteger count = this.count;
    final ReentrantLock takeLock = this.takeLock;
    // 加可中断锁
    takeLock.lockInterruptibly();
    try {
        // 如果当前的数据量为0
        while (count.get() == 0) {
            // 如果时间没有剩余,直接返回null即可
            if (nanos <= 0)
                return null;
            // 让消费者线程等待nanos时间
            nanos = notEmpty.awaitNanos(nanos);
        }
        // 取数据
        x = dequeue();
        // 后面都是一样的
        c = count.getAndDecrement();
        if (c > 1)
            notEmpty.signal();
    } finally {
        takeLock.unlock();
    }
    if (c == capacity)
        signalNotFull();
    return x;
}
复制代码

3.4 take()源码实现

  • 这个大家可以自己看一下补充,也算一个小测试
public E take() throws InterruptedException {
    E x;
    int c = -1;
    final AtomicInteger count = this.count;
    final ReentrantLock takeLock = this.takeLock;
    takeLock.lockInterruptibly();
    try {
        while (count.get() == 0) {
            notEmpty.await();
        }
        x = dequeue();
        c = count.getAndDecrement();
        if (c > 1)
            notEmpty.signal();
    } finally {
        takeLock.unlock();
    }
    if (c == capacity)
        signalNotFull();
    return x;
}
复制代码

4、疑惑

看到这里,我想大家可能有和我一样的疑惑?

之前我们聊 ArrayBlockingQueue 的时候,他只用了一把锁(互斥锁),但 LinkedBlockingQueue 却使用了两把锁(读锁、写锁)

这时候你脑子会不会有一种疑问,我 ArrayBlockingQueue 能不能使用两把锁(读锁、写锁)来进行访问

如果你有这种想法,说明你确实思考了,哈哈哈

没错,博主我查阅了相关的资料,ArrayBlockingQueue 确实可以使用两把锁进行逻辑的更改

整体的逻辑基本上是仿造 LinkedBlockingQueue 的业务逻辑改造的,经测试这种性能要比原始的 ArrayBlockingQueue 要快 20%~30% 左右,感兴趣的也可以自己去测试一下。

import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.ReentrantLock;
 
public class ArrayBlockingQueueUsingTwoLockApproach {
    
     /** The queued items */
    final Object[] items;
 
    /** items index for next take, poll, peek or remove */
    int takeIndex;
 
    /** items index for next put, offer, or add */
    int putIndex;
 
    /** Current number of elements */
    private final AtomicInteger count = new AtomicInteger();
 
    /** Lock held by take, poll, etc */
    private final ReentrantLock takeLock = new ReentrantLock();
 
    /** Wait queue for waiting takes */
    private final Condition notEmpty = takeLock.newCondition();
 
    /** Lock held by put, offer, etc */
    private final ReentrantLock putLock = new ReentrantLock();
 
    /** Wait queue for waiting puts */
    private final Condition notFull = putLock.newCondition();
 
    public ArrayBlockingQueueUsingTwoLockApproach(int capacity) {
        if (capacity <= 0)
            throw new IllegalArgumentException();
        this.items = new Object[capacity];
    }
    
    public void put(Object e) throws InterruptedException {
        if (e == null) throw new NullPointerException();
        int c = -1;
        final ReentrantLock putLock = this.putLock;
        final AtomicInteger count = this.count;
        putLock.lockInterruptibly();
        try {
            while (count.get() == items.length) {
                notFull.await();
            }
            enqueue(e);
            c = count.getAndIncrement();
            if (c + 1 < items.length)
                notFull.signal();
        } finally {
            putLock.unlock();
        }
        if (c == 0)
            signalNotEmpty();
    }
    
    public Object take() throws InterruptedException {
        Object x;
        int c = -1;
        final AtomicInteger count = this.count;
        final ReentrantLock takeLock = this.takeLock;
        takeLock.lockInterruptibly();
        try {
            while (count.get() == 0) {
                notEmpty.await();
            }
            x = dequeue();
            c = count.getAndDecrement();
            if (c > 1)
                notEmpty.signal();
        } finally {
            takeLock.unlock();
        }
        if (c == items.length)
            signalNotFull();
        return x;
    }
    
    
    private void signalNotEmpty() {
        final ReentrantLock takeLock = this.takeLock;
        takeLock.lock();
        try {
            notEmpty.signal();
        } finally {
            takeLock.unlock();
        }
    }
 
    
    /**
     * Inserts element at current put position, advances, and signals.
     * Call only when holding lock.
     */
    private void enqueue(Object x) {
        // assert lock.getHoldCount() == 1;
        // assert items[putIndex] == null;
        final Object[] items = this.items;
        items[putIndex] = x;
        if (++putIndex == items.length)
            putIndex = 0;
        count.incrementAndGet();
    }
 
    
    /**
     * Extracts element at current take position, advances, and signals.
     * Call only when holding lock.
     */
    private Object dequeue() {
        // assert lock.getHoldCount() == 1;
        // assert items[takeIndex] != null;
        final Object[] items = this.items;
        @SuppressWarnings("unchecked")
        Object x = (Object) items[takeIndex];
        items[takeIndex] = null;
        if (++takeIndex == items.length)
            takeIndex = 0;
        count.decrementAndGet();
        return x;
    }
    
    /**
     * Signals a waiting put. Called only from take/poll.
     */
    private void signalNotFull() {
        final ReentrantLock putLock = this.putLock;
        putLock.lock();
        try {
            notFull.signal();
        } finally {
            putLock.unlock();
        }
    }
}
复制代码

四、流程图

其实,我们 LinkedBlockingQueue 整体的代码逻辑和 ArrayBlockingQueue 类似,只不过底层数据结构不同罢了

我们这里简单的画一下,有兴趣的同学也可以自己画吆

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处:http://www.coloradmin.cn/o/463530.html

如若内容造成侵权/违法违规/事实不符,请联系多彩编程网进行投诉反馈,一经查实,立即删除!

相关文章

【 SpringBoot 统⼀功能处理 】

文章目录 引言一、⽤户登录权限效验Spring 拦截器拦截器实现原理扩展&#xff1a;统⼀访问前缀添加 二、统⼀异常处理三、统⼀数据返回格式四、ControllerAdvice 源码分析 引言 接下来是 Spring Boot 统⼀功能处理模块&#xff0c;是 AOP 的实战环节&#xff0c;要实现的课程⽬…

轨道交通信号系统的可靠性与安全性

01.引言 城市轨道交通系统作为大容量公共交通工具&#xff0c;其安全性直接关系到广大乘客的生命安全&#xff0c;所以要求城市轨道交通系统在如此高的运行密度下&#xff0c;还要保证安全和高效率的运行。而信号系统作为保证列车安全、正点、便捷、舒适、高密度不间断运行的重…

Filter 过滤器基本内容及案例改进

举个例子 假设在Web资源中&#xff0c;A资源要写5行代码&#xff0c;而B资源也要写一模一样的5行代码&#xff0c;这时就把这些代码都提取出来&#xff0c; 在过滤器里写这些代码&#xff0c;因为访问任何资源都要经过过滤器&#xff0c;在过滤器走一遍就可以&#xff0c;而不用…

性能优化之20个 Linux 服务器性能调优技巧

Linux是一种开源操作系统&#xff0c;它支持各种硬件平台&#xff0c;Linux服务器全球知名&#xff0c;它和Windows之间最主要的差异在于&#xff0c;Linux服务器默认情况下一般不提供GUI(图形用户界面)&#xff0c;而是命令行界面&#xff0c;它的主要目的是高效处理非交互式进…

【22】核心易中期刊推荐——人工智能与识别图像处理与应用

🚀🚀🚀NEW!!!核心易中期刊推荐栏目来啦 ~ 📚🍀 核心期刊在国内的应用范围非常广,核心期刊发表论文是国内很多作者晋升的硬性要求,并且在国内属于顶尖论文发表,具有很高的学术价值。在中文核心目录体系中,权威代表有CSSCI、CSCD和北大核心。其中,中文期刊的数…

网络编程代码实例:多进程版

文章目录 前言代码仓库内容代码&#xff08;有详细注释&#xff09;server.cclient.cMakefile 结果总结参考资料作者的话 前言 网络编程代码实例&#xff1a;多进程版。 代码仓库 yezhening/Environment-and-network-programming-examples: 环境和网络编程实例 (github.com)E…

商品如果要在美国商超出售,那么如何申请美国条形码呢?

美国条码注册是指向美国条码协会提出条码申请&#xff0c;通过条码协会的审核批准后&#xff0c;条码可以印在产品上使用。条码是进入各大商场&#xff0c;超市的身份证&#xff0c;企业有条形码&#xff0c;一是可以提高企业产品的知名度&#xff1b;二是增加产品的防伪力度&a…

TypeScript与JavaScript

目录 一、什么是javascript 二、TypeScript&#xff1a;静态类型检查器 1、类型化的 JavaScript 超集 1.1、句法 1.2、类型 1.3、运行时行为 1.4、擦除类型 2、学习 JavaScript 还是 TypeScript 一、什么是javascript JavaScript&#xff08;也称为 ECMAScript&#xf…

为何电商这么难做…...你是否忽略了这个问题?

物流时效是影响买家体验的重要环节&#xff0c;物流服务优劣也是买家网上购物时的重要参考依据。但电商企业对于快递公司的时效承诺、服务质量基本处于被动接受的状况&#xff0c;直到买家投诉才知道快递公司服务缺失&#xff0c;若买家不投诉也没法主动知道大量的订单是否按约…

Notes/Domino 11.0.1FP7以及在NAS上安装Domino等

大家好&#xff0c;才是真的好。 目前HCL在还是支持更新的Notes/Domino主要是三个版本&#xff0c;V10、11和12&#xff0c;这不,上周HCL Notes/Domino 11.0.1居然推出了FP7补丁包程序。 从V10.0.1开始&#xff0c;Domino的FP补丁包程序主要是用来修复对应主要版本中的一些问…

TCP FACK 与 RACK

3 个 dupacks 触发 fast retransmit 是一个经典启发算法&#xff0c;但在引入 SACK 之后仍然计数 SACKed 数量 > 3 触发 fast retransmit 似乎就没理由了。即使把 reordering 算进去&#xff0c;一个距离 una 很远的 seg 被 SACKed&#xff0c;也足以判定丢包了&#xff0c;…

2022公考经验分享

一、写在前面 2017南京邮电毕业后&#xff0c;5年来一直就职事业单位。单位解决北京户口&#xff0c;也赶上了一两年的忙碌期&#xff0c;存款加上公积金大概40万。期间经历过几段感情&#xff0c;2020年通过相亲认识现在的老婆。2020年12月瑞泽家园签位排名7000多&#xff0c;…

本地spingboot配置Promethus+granfana监控

记录如何配置与启动 1.在搭建好的应用加上依赖 <!-- 实现对 Actuator 的自动化配置 --><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-actuator</artifactId></dependency><!-- Micr…

超火爆的ChatGPT课,送ChatGPT账号啦~~

HOT! HOT! HOT! &#x1f525; &#x1f525; &#x1f525; 上周&#xff0c;ChatGPT全栈开发课程一经推出&#xff0c;就在程序员圈子中引起了广泛关注。这两天 都被挤爆了&#xff0c;纷纷表示对课程内容很是期待呢。 明天就要开班直播啦&#xff0c;还未报名的同学&…

npm 打包发布一个公用的组件包

1&#xff0c;首先创建一个需要发包发布的组件 2.3使用Vue插件模式 这一步是封装组件中的重点&#xff0c;用到了Vue提供的一个公开方法&#xff1a;install。这个方法会在你使用Vue.use(plugin)时被调用&#xff0c;这样使得我们的插件注册到了全局&#xff0c;在子组件的任何…

AI在网络安全中的应用:机器学习如何帮助我们更好地保护网络

章节一&#xff1a;引言 随着信息技术的飞速发展&#xff0c;网络攻击的手段也在不断地演变。传统的网络安全技术已经难以应对日益复杂的网络安全威胁。AI技术&#xff0c;特别是机器学习技术&#xff0c;为网络安全提供了一种新的解决方案。本文将介绍AI在网络安全中的应用&am…

打造高性能的视频和弹幕系统(一): 对象存储服务

这里写自定义目录标题 欢迎使用Markdown编辑器新的改变功能快捷键合理的创建标题&#xff0c;有助于目录的生成如何改变文本的样式插入链接与图片如何插入一段漂亮的代码片生成一个适合你的列表创建一个表格设定内容居中、居左、居右SmartyPants 创建一个自定义列表如何创建一个…

一文综述:自然语言处理技术NLP

自然语言处理技术综述1-到2020年 写在最前面摘要NLP简介Preprocessing预处理Tokenization令牌化、标记化Stop Words 停用词Stemming and Lemmatization词干提取和词形还原&#xff08;英文单词&#xff09;Parts-of-Speech Tagging词性标记Bag of Words and N-Grams词袋模型、N…

Redis数据库的安装(Windows10)

Redis数据库的安装 前言安装启动命令简单的几条语句 前言 本节开始学习Redis数据库。 Redis数据库的优势如下&#xff1a; 性能极高 – Redis能读的速度是110000次/s,写的速度是81000次/s 。 丰富的数据类型 – Redis支持二进制案例的 Strings, Lists, Hashes, Sets 及 Ord…

ubuntu22.04安装显卡驱动+cuda+cudnn

ubuntu22.04安装显卡驱动cudacudnn 1. 下载驱动和卸载、禁用自带驱动程序1.1 查看系统显卡型号1.2 从NVIDIA官网下载相应驱动1.3 卸载Ubuntu自带的驱动程序1.4 禁用自带的nouveau nvidia驱动1.5 更新1.6 重启电脑1.7 查看是否将自带的驱动屏蔽 2. 安装显卡驱动2.1 停止lightdm桌…