Головокружительный параллелизм Java (4) Очередь блокировки — дело о расследовании роста ЦП

Java
Головокружительный параллелизм Java (4) Очередь блокировки — дело о расследовании роста ЦП

В предыдущей статье я представил AQS от Niubi и примерно объяснил идею синхронизации в JUC. Не придумал, что написать в этой статье, Буквально на прошлой неделе в коде коллеги возникла проблема, после расследования выяснилось, что она вызвана неправильным использованием блокирующих очередей, поэтому в этой статье было решено ввести блокирующие очереди.

реальный кейс

Случай ошибки:

Также довольно совпадение, что в тот день iMac коллеги сменился на Macbook Pro. Потом запускал как обычно разные службы, а через какое-то время компьютерный вентилятор бешено работал и шумел.Потому что проект IDEA на iMac обычно занимает много памяти и выполняется долго, он зависнет, так что он было все равно. Но после этого мы подумали, что это очень странно, а потом открыли его монитор активности и обнаружили, что некий Java-процесс на самом деле занимает 90% процессора, а потом подтвердили, что это за проект, и, наконец, проверили проект через jstack. ситуация с потоком , находит пользовательский поток, а затем просматривает код и находит следующее:

MyThreadPool.exportEnclosurePool.execute(() -> {
    while (true) {
        BlockingQueue<EnclosureRequest> blockingQueue = requestQueue.getBlockingQueue();
           while (!blockingQueue.isEmpty()) {
               System.out.println("开始消费");
               EnclosureRequest one = null;
               try {
                   one = blockingQueue.take();
                   ossService.exportEnclosureToLocalServer(one.getEnclosureList(), one.getSobId(), one.getUserUuid(), one.getUserName(), one.getTmpFileName(), one.getZipUuidList());
               } catch (Exception e) {
                   e.printStackTrace();
               }
            }
    }
}

Требование коллеги состоит в том, чтобы сделать функцию экспорта вложений в очереди, поэтому он выбрал режим производитель-потребитель и реализовал его, заблокировав очереди; но, поскольку он не был знаком с этим, он написал этот проблемный код, который привел к бесконечному циклу; к счастью, этот код был найден в тестовой ветке и официально не публиковался. Правильная реализация потребительского кода выглядит следующим образом:

Правильная реализация:

MyThreadPool.exportEnclosurePool.execute(() -> {
    BlockingQueue<EnclosureRequest> blockingQueue = requestQueue.getBlockingQueue();
    while (true) {
        try {
            EnclosureRequest one = blockingQueue.take();
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        System.out.println("开始消费");
        ossService.exportEnclosureToLocalServer(one.getEnclosureList(), one.getSobId(), one.getUserUuid(), one.getUserName(), one.getTmpFileName(), one.getZipUuidList());
    }
}

Введение в блокирующие очереди

Блокирующая очередь — это очередь, методы вставки и удаления которой поддерживают дополнительные операции;
  • Методы вставки, поддерживающие блокировку: когда очередь блокировки заполнена, очередь блокирует поток, который вставляет элементы, до тех пор, пока очередь не будет заполнена.
  • Методы удаления блокировки: когда очередь блокировки пуста, поток, который получает элемент очереди, будет заблокирован до тех пор, пока очередь не станет пустой.

Четыре метода обработки:

Способ\метод обработки Метод вставки метод удаления
Выбросить исключение add(e) remove()
вернуть логическое значение offer(e) poll()
блокировать put(e) take()
тайм-аут offer(e,time,unit) poll(time,unit)

💡Советы:Если это неограниченная блокирующая очередь, очередь не может быть заполнена, поэтому метод put() никогда не будет заблокирован, а метод offer() всегда будет возвращать значение true.

Введение в блокирующие очереди в Java

  • ArrayBlockingQueue: ограниченная очередь блокировки на основе массива, которая поддерживает настройку политик справедливости.
  • LinkedBlockingQueue: неограниченная (по умолчанию Integer.MAX_VALUE) очередь блокировки на основе связанного списка, рабочая очередь, используемая newFixedThreadPool() и newSingleThreadExecutor() в Executors, поэтому Executors не рекомендуются.
  • LinkedBlockingDeque: неограниченная (по умолчанию Integer.MAX_VALUE) двунаправленная очередь блокировки на основе связанного списка
  • LinkedTransferQueue: неограниченная очередь блокировки на основе связанного списка. Очередь предоставляет метод передачи (e). Если потребитель ожидает, элемент отправляется непосредственно потребителю. В противном случае элемент помещается в хвостовой узел очереди. и блокируется до тех пор, пока элемент не будет израсходован.
  • PriorityBlockingQueue: неограниченная очередь блокировки, которая поддерживает сортировку по приоритету. По умолчанию она сортируется в естественном порядке в порядке возрастания. Вы также можете переопределить метод compareTo(), чтобы указать правила сортировки элементов, или указать параметр построения Comparator для сортировки, когда очередь инициализируется.
  • DelayQueue: неограниченная очередь блокировки задержки, реализованная с использованием PriorityQueue.
  • SynchronousQueue: очередь блокировки, в которой не хранятся элементы. Каждая операция размещения должна блокироваться до тех пор, пока не произойдет операция взятия, иначе элементы не могут быть добавлены. Поддерживает настройку политик справедливости.

Анализ принципа реализации очереди блокировки (LinkedBlockingQueue)

LinkedBlockingQueue представляет собой односвязную структуру списка, состоящую из переменных-членов Node. Емкость по умолчанию — максимальное значение Integer. Существуют также две блокировки ReentrantLock, putLock и takeLock, которые используются для обеспечения потокобезопасности вставки и удаления (блокировка ReentrantLock используется в других очередях блокировки). ), две очереди ожидания состояния notEmpty и notFull используются для хранения потоков, заблокированных с помощью take() и put(). Здесь я кратко анализирую два более важных метода put() и take().

Анализ исходного кода

/**
 * 由Node节点组成单链表结构
 */
static class Node<E> {
    E item;
    Node<E> next;
    Node(E x) { item = x; }
}
/** 用于移除操作的锁 */
private final ReentrantLock takeLock = new ReentrantLock();

/** 阻塞于take的等待队列 */
private final Condition notEmpty = takeLock.newCondition();

/** 用于插入操作的锁 */
private final ReentrantLock putLock = new ReentrantLock();

/** 阻塞于put的等待队列 */
private final Condition notFull = putLock.newCondition();

/**
 * 不指定容量默认是Integer的最大值
 */
public LinkedBlockingQueue() {
    this(Integer.MAX_VALUE);
}

/**
 * 阻塞式插入元素(队列为满则阻塞)
 */
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();
}
/**
 * 阻塞式移除元素(队列为空则阻塞)
 */
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;
    }

Графический анализ

Следует отметить, что операция put() добавляет элемент в очередь и освобождает блокировку после определения того, меньше ли емкость, чем верхний предел, и уведомляет очередь ожидания notFull.Прежде чем уведомлять очередь notEmpty, вам необходимо сначала получить takeLock То же верно и для операции take().

💡Советы:Существует большая разница между методами put() и take() LinkedBlockingQueue и другими блокирующими очередями. Другие очереди блокировки будут уведомлять соответствующую ожидающую очередь каждый раз, когда put() и take(), но LinkedBlockingQueue будет уведомлять notEmpty только в том случае, если она пуста перед размещением, и уведомлять очередь ожидания notFull, если она заполнена до выполнения, и это не будет заполнен после пут.Очередь ожидания notFull уведомляется, а очередь ожидания notEmpty уведомляется, когда дубль не пуст. Мое личное понимание этого вопроса состоит в том, что в LinkedBlockingQueue есть блокировки чтения-записи. пробуждение потока в очереди ожидания notFull уменьшает количество блокировок. В другой очереди блокировки есть только одна блокировка, поэтому уведомлениям не нужно бороться за другую блокировку. Конечно, это только мое личное мнение, надеюсь, кто-то из понимающих сможет дать мне совет.

Суммировать

Очереди с блокировкой очень важны для параллелизма. Очереди с блокировкой используются в пулах потоков, представленных ранее. Модель производитель-потребитель также может быть реализована с помощью очередей с блокировкой. До сих пор были введены AQS, очереди с блокировкой и пулы потоков. Надеюсь, вы может связать их.Понимание углубляет впечатление.

Прошлые статьи:

Добро пожаловать друзья, которые также заинтересованы в обсуждении