Анализ исходного кода — проблема увеличения памяти, вызванная пулом потоков newFixedThreadPool

Java

предисловие

Будут ли пулы потоков с неограниченными очередями вызывать всплески памяти? Интервьюеры часто задают этот вопрос.В этой статье будет проанализирована проблема увеличения памяти, вызванная пулом потоков 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:

Проблема с JVM-ООМВ целомСоздать слишком много объектов,в то же времяGC мусорСлишком поздно перерабатывать, так в чем причина?Причины 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, я предполагаю, что проблема с перегрузкой памятирабочая очередь заполнена. Затем проанализируйте исходный код очереди блокировки, чтобы раскрыть тайну проблемы с парящим объемом памяти.

Блокировка анализа исходного кода очереди

Вернитесь к конструктору newFixedThreadPool и обнаружите, что очередью блокировки является LinkedBlockingQueue, и этоБез параметра очередь LinkedBlockingQueue. Хорошо, тогда мы непосредственно анализируем исходный код LinkedBlockingQueue.

Диаграмма классов 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.

Ссылка и спасибо

Личный публичный аккаунт

  • Если вы хороший ребенок, который любит учиться, вы можете подписаться на мой официальный аккаунт, чтобы учиться и обсуждать вместе.
  • Если вы считаете, что в этой статье есть какие-либо неточности, вы можете прокомментировать или подписаться на мой официальный аккаунт, пообщаться со мной в частном порядке, и все смогут учиться и прогрессировать вместе.