предисловие
Будут ли пулы потоков с неограниченными очередями вызывать всплески памяти? Интервьюеры часто задают этот вопрос.В этой статье будет проанализирована проблема увеличения памяти, вызванная пулом потоков newFixedThreadPool на основе исходного кода, в надежде углубить ваше понимание.
Проблема увеличения памяти повторяется
пример кода
ExecutorService executor = Executors.newFixedThreadPool(10);
for (int i = 0; i < Integer.MAX_VALUE; i++) {
executor.execute(() -> {
try {
Thread.sleep(10000);
} catch (InterruptedException e) {
//do nothing
}
});
}
Настроить параметры JVM
IDE указывает параметры JVM: -Xmx8m -Xms8m :
Результаты
Запуск приведенного выше кода вызовет OOM:
Анализ исходного кода пула потоков
Вышепример кода,НаnewFixedThreadPoolиexecuteметод. Во-первых, давайте взглянем на исходный код метода newFixedThreadPool.
Исходный код newFixedThreadPool
public static ExecutorService newFixedThreadPool(int nThreads) {
return new ThreadPoolExecutor(nThreads, nThreads,
0L, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<Runnable>());
}
Исходный код этого раздела и комбинацияФункции пула потоков, мы можем знать, чтоnewFixedThreadPool:
- Количество основных потоков coreSize и максимальное количество потоков maxPoolSizeТот же размер, nThreads.
- Время простоя равно 0, т.е.KeepAliveTime 0
- очередь блокировкипостроен без аргументовLinkedBlockingQueue
Для друзей, которые не очень хорошо знают характеристики пула потоков, вы можете прочитать мою статью,Требования к собеседованию: анализ пула потоков Java
Далее давайте взглянем на исходный код метода выполнения пула потоков execute.
Исходный код метода выполнения пула потоков execute
Исходный код выполнения и связанных с нимобъяснятьследующее:
public void execute(Runnable command) {
if (command == null)
throw new NullPointerException();
int c = ctl.get();
if (workerCountOf(c) < corePoolSize) { //步骤一:判断当前正在工作的线程是否比核心线程数量小
if (addWorker(command, true)) // 以核心线程的身份,添加到工作集合
return;
c = ctl.get();
}
//步骤二:不满足步骤一,线程池还在RUNNING状态,阻塞队列也没满的情况下,把执行任务添加到阻塞队列workQueue。
if (isRunning(c) && workQueue.offer(command)) {
int recheck = ctl.get();
//来个double check ,检查线程池是否突然被关闭
if (! isRunning(recheck) && remove(command))
reject(command);
else if (workerCountOf(recheck) == 0)
addWorker(null, false);
}
//步骤三:如果阻塞队列也满了,执行任务以非核心线程的身份,添加到工作集合
else if (!addWorker(command, false))
reject(command);
}
Глядя на приведенный выше код, мы можем найти, чтоaddWorker и workQueue.offer(команда)Возможно создание объекта. Тогда давайте сначала проанализируем метод addWorker.
Анализ исходного кода addWorker
Исходный код addWorker и соответствующие пояснения приведены ниже.
private boolean addWorker(Runnable firstTask, boolean core) {
retry:
for (;;) {
int c = ctl.get();
//获取当前线程池的状态
int rs = runStateOf(c);
//如果线程池状态是STOP,TIDYING,TERMINATED状态的话,则会返回false。
// 如果现在状态是SHUTDOWN,但是firstTask不为空或者workQueue为空的话,那么直接返回false
if (rs >= SHUTDOWN &&
! (rs == SHUTDOWN &&
firstTask == null &&
! workQueue.isEmpty()))
return false;
//自旋
for (;;) {
//获取当前工作线程的数量
int wc = workerCountOf(c);
//判断线程数量是否符合要求,如果要创建的是核心工作线程,判断当前工作线程数量是否已经超过coreSize,
// 如果要创建的是非核心线程,判断当前工作线程数量是否超过maximumPoolSize,是的话就返回false
if (wc >= CAPACITY ||
wc >= (core ? corePoolSize : maximumPoolSize))
return false;
//如果线程数量符合要求,就通过CAS算法,将WorkerCount加1,成功就跳出retry自旋
if (compareAndIncrementWorkerCount(c))
break retry;
c = ctl.get(); // Re-read ctl
if (runStateOf(c) != rs)
continue retry;
retry inner loop
}
}
//线程启动标志
boolean workerStarted = false;
//线程添加进集合workers标志
boolean workerAdded = false;
Worker w = null;
try {
//由(Runnable 构造Worker对象
w = new Worker(firstTask);
final Thread t = w.thread;
if (t != null) {
//获取线程池的重入锁
final ReentrantLock mainLock = this.mainLock;
mainLock.lock();
try {
//获取线程池状态
int rs = runStateOf(ctl.get());
//如果状态满足,将Worker对象添加到workers集合
if (rs < SHUTDOWN ||
(rs == SHUTDOWN && firstTask == null)) {
if (t.isAlive())
throw new IllegalThreadStateException();
workers.add(w);
int s = workers.size();
if (s > largestPoolSize)
largestPoolSize = s;
workerAdded = true;
}
} finally {
mainLock.unlock();
}
//启动Worker中的线程开始执行任务
if (workerAdded) {
t.start();
workerStarted = true;
}
}
} finally {
//线程启动失败,执行addWorkerFailed方法
if (! workerStarted)
addWorkerFailed(w);
}
return workerStarted;
}
процесс выполнения addWorker
вероятно суждениеСостояние пула потоков в порядке?, если все в порядке, оценка количества потоков в текущем заданииСоответствует ли он (меньше, чем coreSize/maximumPoolSize), если нет, не добавляйте, если удовлетворено, добавьте задачу выполнения в рабочий набор рабочих процессов и начните выполнение потока.
Взгляните еще раз на типы работников:
/**
* Set containing all worker threads in pool. Accessed only when
* holding mainLock.
*/
private final HashSet<Worker> workers = new HashSet<Worker>();
workers — это коллекция HashSet, которая управляется параметрами coreSize/maximumPoolSize, поэтому метод addWorker вызовет OOM? комбинироватьПример демо кода, coreSize=maximumPoolSize=10, если больше 10, то не добавится в воркеры, так что не является причиной парения памяти newFixedThreadPool. Тогда проблема должна заключаться в методе workQueue.offer(command). Для того, чтобы весь процесс был понятен, давайте нарисуем блок-схему выполнения выполнения.
Процесс выполнения метода пула потоков execute
В соответствии с приведенным выше анализом исходного кода execute и addWork мы рисуем блок-схему:
- Отправьте команду задачи, когда количество основных потоков, оставшихся в пуле потоков, меньше, чем количество потоков corePoolSize, вызовите метод addWorker, и пул потоков создаст основной поток для обработки отправленной задачи.
- Если количество основных потоков в пуле потоков заполнено, то есть количество потоков равно corePoolSize, вновь отправленная задача будет помещена в очередь задач workQueue для выполнения.
- Когда количество выживших нитей в пуле резьбы равно CorePoolSize, и рабочая задача очереди задач заполнена, определить, достигает количества потоков максимально, то есть ли максимальное количество потоков, а если нет, создавать неядерная поток для выполнения представленной задачи.
- Если текущее количество потоков достигает максимального размера пула и появляются новые задачи, политика отклонения будет использоваться напрямую.
Прочитав процесс выполнения execute, я предполагаю, что проблема с перегрузкой памятирабочая очередь заполнена. Затем проанализируйте исходный код очереди блокировки, чтобы раскрыть тайну проблемы с парящим объемом памяти.
Блокировка анализа исходного кода очереди
Диаграмма классов LinkedBlockingQueue
- LinkedBlockingQueue реализован с использованием односвязного списка, который имеет два узла, которые используются для хранения головных и хвостовых узлов соответственно, и атомарную переменную count с начальным значением 0, которая используется для записи количества элементов очереди.
- Есть также два экземпляра ReentrantLock, которые используются для управления атомарностью ввода и удаления элементов из очереди соответственно. управление может иметь только элементы одновременно.Один поток может получить блокировку, добавить элемент в конец очереди, а другие потоки должны ждать.
- Кроме того, notEmpty и notFull являются условными переменными, а внутри них находится очередь условий для хранения потоков, заблокированных при входе и выходе из очереди, по сути, это модель производитель-потребитель.
Нет ссылочного конструктора LinkedBlockingQueue
public LinkedBlockingQueue() {
this(Integer.MAX_VALUE);
}
public LinkedBlockingQueue(int capacity) {
if (capacity <= 0) throw new IllegalArgumentException();
this.capacity = capacity;
last = head = new Node<E>(null);
}
LinkedBlockingQueue конструктор без аргументов, конструкция по умолчаниюInteger.MAX_VALUE (настолько большой)Связанный список см. здесь, вы вспоминаете процесс выполнения, никогда не будет ли очередь блокировки заполнена, эта очередь не откажется прийти, и все задачи блокировки будут собраны под вашей командой. . . Это проблема парящего объема памяти, которая выходит на первый план?
Функция предложения LinkedBlockingQueue
В пуле потоков метод предложения используется для вставки очереди.Давайте посмотрим на операцию предложения очереди блокировки LinkedBlockingQueue.
public boolean offer(E e) {
//为空元素则抛出空指针异常
if (e == null) throw new NullPointerException();
final AtomicInteger count = this.count;
//如采当前队列满则丢弃将要放入的元素, 然后返回false
if (count.get() == capacity)
return false;
int c = -1;
//构造新节点,获取putLock独占锁
Node<E> node = new Node<E>(e);
final ReentrantLock putLock = this.putLock;
putLock.lock();
try {
//如采队列不满则进队列,并递增元素计数
if (count.get() < capacity) {
enqueue(node);
c = count.getAndIncrement();
//新元素入队后队列还有空闲空间,则
唤醒 notFull 的条件队列中一条阻塞线程
if (c + 1 < capacity)
notFull.signal();
}
} finally {
//释放锁
putLock.unlock();
}
if (c == 0)
signalNotEmpty();
return c >= 0;
}
предложить операциюВставьте элемент в конец очереди. Если очередь свободна, она вернет true после успешной вставки. Если очередь заполнена, текущий элемент будет отброшен, и будет возвращено значение false. Выдает Nul!PointerException, если элемент e имеет значение null. Кроме того, метод является неблокирующим.
Обнародованы результаты проблемы с быстрорастущей памятью
Пул потоков newFixedThreadPoolКоличество основных потоков фиксировано, он использует почтиНеограниченная очередь блокировки LinkedBlockingQueue. Когда основной поток израсходован, задача будет помещена в очередь блокировки.Если время выполнения задачи относительно велико и не освобождается, это вызоветВ очереди на блокировку накапливается все больше и больше задач, и, наконец, вызывают непрерывный рост использования памяти машины, что приводит к JVM OOM.
Ссылка и спасибо
- «Красота параллельного программирования на Java»
- Требования к собеседованию: анализ пула потоков Java
Личный публичный аккаунт
- Если вы хороший ребенок, который любит учиться, вы можете подписаться на мой официальный аккаунт, чтобы учиться и обсуждать вместе.
- Если вы считаете, что в этой статье есть какие-либо неточности, вы можете прокомментировать или подписаться на мой официальный аккаунт, пообщаться со мной в частном порядке, и все смогут учиться и прогрессировать вместе.