предисловие
В течение долгого времени я обнаружил, что многие разработчикиJ.U.C(java.util.concurrent), то есть использование параллельных пакетов Java редко, не говоря уже о понимании этого, но это также необходимый уровень для нашего продвижения.
Я уже более или менее делился релевантным контентом, но он не систематичен, поэтому я хочу организовать серию статей, связанных с параллельными пакетами.
В основном он включает в себя следующие части:
- Реализуйте инструмент параллелизма самостоятельно по определению.
- Стандартная реализация JDK.
- Практический случай.
Исходя из этих трех пунктов, я считаю, что эта часть контента никого не запутает.
Поскольку я открыл новую яму, я не очень хочу это делать, поэтому я собираюсь охватить большинство классов в этом списке.
Поэтому данное обсуждение сосредоточено наArrayBlockingQueue.
реализовать это самостоятельно
Прежде чем реализовать это самостоятельно, сначала разберитесь с некоторыми особенностями очереди блокировки:
- Основные функции очереди: первый пришел, первый ушел.
- Блокирует, когда место в очереди записи недоступно.
- Блокирует, когда очередь пуста при получении данных очереди.
Существует много способов реализации очередей, вообще говоря, это массивы и связанные списки, на самом деле нам нужно разобраться только с одним из них, а разные характеристики в основном заключаются в отличии массивов от связанных списков.
здесьArrayBlockingQueueИз названия очевидно, что он реализован массивом.
Попробуем реализовать его самостоятельно, исходя из этих трех характеристик.
Инициализировать очередь
Я настроил класс здесь:ArrayQueue, его конструктор выглядит следующим образом:
public ArrayQueue(int size) {
items = new Object[size];
}
Очевидно здесьitemsЭто массив для хранения данных, при инициализации нужно создать массив по размеру.
очередь записи
Запись в очередь относительно проста, вам нужно только хранить данные в этом массиве по очереди, как показано ниже:
Но есть еще несколько моментов, на которые стоит обратить внимание:
- Когда очередь заполнена, поток записи необходимо заблокировать.
- Когда количество операций записи в очередь превышает размер очереди, необходимо начать запись с первого нижнего индекса.
посмотри на первый队列满的时候,写入的线程需要被阻塞, сначала рассмотрим, как сделать потокблокировать, похоже, что поток застрял и ничего не может сделать.
Есть несколько вариантов достижения такого эффекта:
-
Thread.sleep(timeout)Нить спит. -
object.wait()пусть нить войдетwaitingгосударство.
Конечно, есть некоторые
join、LockSupport.partи т. д. не входят в рамки этого обсуждения.
Еще одна очень важная особенность блокирующих очередей заключается в том, что когда пространство очереди доступно (выход из очереди), поток записи необходимо разбудить, чтобы в него можно было записать данные.
так очевидноThread.sleep(timeout)Не подходит, он будет продолжать работать после истечения времени ожидания; не достигнутокогда есть местоПросыпайтесь только для продолжения работы этой функции.
Фактически, такая функция легко заставляет нас думать о механизме уведомления об ожидании в Java для реализации межпотоковой связи; больше потоков видят схемы связи, пожалуйста, обратитесь сюда:Глубокое понимание взаимодействия потоков
Итак, что я здесь делаю, так это то, что как только очередь заполнена, поток записи вызываетobject.wait()Войтиwaitingдо тех пор, пока не произойдет пробуждение, когда освободится место.
/**
* 队列满时的阻塞锁
*/
private Object full = new Object();
/**
* 队列空时的阻塞锁
*/
private Object empty = new Object();
Поэтому здесь объявляются два объекта для взаимного оповещения, когда очередь заполнена и пуста.
Необходимо использовать после успешной записи данныхempty.notify(), цель этого состоит в том, чтобы разбудить поток, который использует очередь, после успешной записи данных, когда очередь получения пуста.
Операции ожидания и уведомления здесь должны использоваться для соответствующих объектов.
synchronizedблок метода, потому что и ожидание, и уведомление должны получить свои собственные блокировки.
очередь потребления
Также упоминалось выше: когда очередь пуста, поток, который получает очередь, должен быть заблокирован до тех пор, пока в очереди не появятся данные для пробуждения.
Код очень похож на то, что было написано, и его легко понять, просто ожидание и пробуждение здесь как раз наоборот, что хорошо можно понять по следующей картинке:
В целом это:
- Когда очередь записи заполнена, она будет заблокирована до тех пор, пока поток получения не проснется после использования данных очереди.написать тему.
- Когда очередь потребления пуста, она будет заблокирована до тех пор, пока поток записи не проснется после записи данных очереди.потребительская нить.
контрольная работа
Начнем с базового теста: однопоточная запись и потребление.
3
123
1234
12345
В результатах нет ничего плохого.
Когда записываемые данные превышают размер очереди, они могут быть записаны только после потребления.
2019-04-09 16:24:41.040 [Thread-0] INFO c.c.concurrent.ArrayQueueTest - [Thread-0]123
2019-04-09 16:24:41.040 [main] INFO c.c.concurrent.ArrayQueueTest - size=3
2019-04-09 16:24:41.047 [main] INFO c.c.concurrent.ArrayQueueTest - 1234
2019-04-09 16:24:41.048 [main] INFO c.c.concurrent.ArrayQueueTest - 12345
2019-04-09 16:24:41.048 [main] INFO c.c.concurrent.ArrayQueueTest - 123456
Из текущих результатов также видно, что данные могут быть записаны в очередь только после их использования.
При отсутствии потребления запись данных в очередь приведет к блокировке потока записи.
Параллельное тестирование
Три потока одновременно записывают 300 фрагментов данных, а один поток потребляет один.
=====0
299
Окончательный размер очереди равен 299, и видимый поток также безопасен.
Вся очередь потокобезопасна, потому что операции как в методах записи, так и в методах получения должны получать блокировки для работы.
ArrayBlockingQueue
Давайте посмотрим на стандарт JDK.ArrayBlockingQueueРеализация вышеизложенного будет лучше понята с вышеуказанным основанием.
Инициализировать очередь
Кажется, что это сложнее, но на самом деле это легко понять после постепенного разделения:
Первый шаг фактически такой же, как мы написали сами, инициализируя массив размера очереди.
Второй шаг инициализирует повторную блокировку, которая на самом деле такая же, как и раньше.synchronizedпоследовательный;
Просто при инициализации реентерабельной блокировки по умолчанию стоит非公平锁, и, конечно, также может быть указано какtrueИспользуйте справедливые блокировки; это будет записывать и потреблять в порядке очереди.
больше о
ReentrantLockПожалуйста, обратитесь к использованию и принципу здесь:Принцип реализации ReentrantLock
Создаются три или четыре шагаnotEmpty notFullЭти два условия, его влияние на использование и ранее использовавшеесяobject.wait/notifyаналогичный.
Это содержание всей инициализации, которая на самом деле очень похожа на то, что мы реализовали сами.
очередь записи
На самом деле вы обнаружите, что принципы блокировки записи аналогичны, за исключением того, что Lock используется здесь для явного получения и освобождения блокировок.
в то же времяnotFull.await();notEmpty.signal();и что мы использовали раньшеobject.wait/notifyИспользование и функции одинаковы.
Конечно, он по-прежнему реализует блокировку тайм-аута.API.
Это также относительно просто, используя метод ожидания с тайм-аутом.
очередь потребления
Посмотрите еще раз на очередь потребления:
Это почти то же самое, вы можете понять это с первого взгляда.
Также используется API тайм-аутаnotEmpty.awaitNanos(nanos)Чтобы реализовать возврат таймаута, мы не будем подробно говорить об этом.
фактический случай
Сказав это, давайте посмотрим на реальный случай очереди.
Фон такой:
Есть запланированная задача, которая через определенный интервал считывает пакет данных из базы данных, ей нужно проверить данные и вызвать удаленный интерфейс.
Самый простой способ — использовать поток этой временной задачи для завершения всего процесса чтения данных, проверки сообщения и вызова интерфейса; но это будет иметь проблему:
Если предположить, что в вызове внешнего интерфейса есть аномалия, а нестабильность сети приводит к увеличению затрат времени, эффективность всей задачи будет снижена, потому что все последовательно и будет влиять друг на друга.
Итак, мы улучшили схему:
По сути, это типичная модель производитель-потребитель:
- Рабочий поток считывает сообщение из базы данных и помещает его в очередь.
- Поток-потребитель получает данные из очереди для бизнес-логики.
Таким образом, два потока могут быть разделены через очередь, не влияя друг на друга, и очередь может также играть роль буфера.
Но есть и некоторые мелкие детали, на которые стоит обратить внимание в процессе использования.
Поскольку этот внешний интерфейс поддерживает пакетное выполнение, после того, как поток-потребитель извлечет данные, он произведет накопление в памяти.После достижения порога или накопления периода времени накопленные данные будут обработаны.
Однако из-за невнимательности разработчика используется при потребленииqueue.take()Это блокирующий API, он отлично работает.
Но как только исходный источник данных, то есть данные в БД отсутствуют, потребляющий поток будет заблокирован после того, как данные в очереди также будут использованы.
Таким образом, данные, накопленные в памяти в предыдущем цикле, не могут быть использованы до тех пор, пока в источнике данных снова не появятся данные.Если интервал будет длинным, это может привести к серьезным бизнес-исключениям.
Так что нам лучше использоватьqueue.poll(timeout)Такой апи с таймаутом, если только нет четкого бизнес-требования к блокировке.
Эта привычка также применима к другим сценариям, таким как вызов http, rpc-интерфейсов и т. д., для всех которых необходимо установить разумный тайм-аут.
Суммировать
оArrayBlockingQueueЭто конец связанного совместного использования, а затем мы продолжим обновлять другие параллельные контейнеры и параллельные инструменты.
Если у вас есть какие-либо вопросы по этой статье, вы можете оставить сообщение для обсуждения.
Весь исходный код, задействованный в этой статье:
Ваши лайки и репост - лучшая поддержка для меня