← назад к разделу

Создавать новый поток на каждую задачу — дорого и опасно. ExecutorService решает эту проблему: держит готовые потоки и распределяет между ними работу.

Обязательно

Проблема: зачем вообще пул

Поток в Java — дорогой объект. Само создание занимает несколько десятков микросекунд, и это ещё не главная беда: под каждый поток резервируется адресное пространство под стек — на 64-битной машине обычно мегабайт, а на macOS с ARM-процессором два. Если приложение создаёт по потоку на каждый HTTP-запрос или фоновую задачу, при нагрузке это приводит к двум проблемам:

  • Перерасход памяти. У каждого потока свой стек, по умолчанию около мегабайта; тысяча потоков — это гигабайт адресного пространства, зарезервированного под одни только стеки (реально занятой памяти меньше — страницы выделяются по мере использования).
  • Перегрузка планировщика. Операционной системе приходится постоянно переключаться между тысячами потоков, и на полезную работу времени остаётся всё меньше. Обычно упираются в это раньше, чем в память.

Пул потоков — это набор заранее созданных потоков, которые ожидают задачи в очереди. Задача поступает → свободный поток берёт её → выполняет → возвращается ждать следующую. Создание потока происходит один раз, а не при каждом запросе.

ExecutorService: базовый интерфейс

ExecutorService — главный интерфейс для управления пулом. Он расширяет Executor, добавляя возможность отправлять задачи с результатом и управлять жизненным циклом.

Задачу передают двумя способами:

живой пример

import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;

public class PoolBasics {
    public static void main(String[] args) throws Exception {
        ExecutorService pool = Executors.newFixedThreadPool(2);

        pool.execute(() -> System.out.println("execute: результата нет, поток " + Thread.currentThread().getName()));

        Future<Integer> receipt = pool.submit(() -> {
            Thread.sleep(200);
            return 42;
        });
        System.out.println("квитанция на руках, задача ещё считает: isDone=" + receipt.isDone());
        System.out.println("get() дождался результата: " + receipt.get());

        pool.shutdown();
        System.out.println("пул остановился за секунду: " + pool.awaitTermination(1, TimeUnit.SECONDS));
    }
}
Запустить

Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →

Callable<V> отличается от Runnable тем, что может вернуть значение и выбросить проверяемое исключение. Future<V> — это «квитанция»: задача ещё выполняется, но результат можно получить позже через future.get(), и этот вызов блокирует вызывающий поток до готовности.

В выводе isDone=false — не гарантия, а следствие Thread.sleep(200) внутри задачи: за то время, что главный поток печатал строчку, задача заведомо не успела. Без этой паузы результат был бы как повезёт.

Когда задач сразу несколько и нужны все результаты, вместо горсти submit берут invokeAll(...): он принимает коллекцию Callable, ждёт, пока отработают все, и возвращает список готовых Future в том же порядке. У него есть и вариант с таймаутом, invokeAll(tasks, 5, TimeUnit.SECONDS): по истечении срока оставшиеся задачи отменяются, а в их Future лежит отмена — на get() прилетит CancellationException. Родственный invokeAny(...) возвращает результат первой успешно завершившейся задачи и отменяет остальные: так опрашивают несколько реплик или зеркал, когда нужен любой ответ побыстрее.

Что возвращает Future, кроме значения: get(2, TimeUnit.SECONDS) ждёт с потолком и бросает TimeoutException, а не висит вечно, — в сервисном коде берут именно его. cancel(true) переводит задачу в отменённое состояние и прерывает выполняющий её поток; cancel(false) только не даёт ей стартовать, если она ещё в очереди. Важно, что прерывание это просьба (см. статью про потоки): задача, которая не проверяет флаг и не вызывает блокирующих методов, доработает до конца, сколько её ни отменяй.

Куда девается исключение задачи

Самое неприятное в пулах — не настройки, а то, что ошибка задачи может исчезнуть бесследно. Поведение зависит от того, каким методом её отправили.

ExecutorService pool = Executors.newFixedThreadPool(2);

pool.execute(() -> { throw new IllegalStateException("execute"); });   // стек напечатается
Future<?> f = pool.submit(() -> { throw new IllegalStateException("submit"); });  // тишина
Thread.sleep(200);
System.out.println("задача провалилась: " + f.isDone());               // true, и всё

execute: исключение доходит до потока пула, тот умирает, обработчик по умолчанию печатает стек в System.err (мимо вашего логгера), а пул молча заводит новый поток взамен. Увидеть это можно, но только если кто-то читает System.err.

submit: исключение перехватывается и кладётся в Future. Нигде не печатается, поток не умирает. Если результат никому не нужен и get() никто не вызывает — а именно так обычно и бывает с фоновыми задачами, — сбой исчезает совсем: задача «выполнена», в логах пусто, данные не обработаны. Это одна из самых частых причин «задача просто не отработала, и никто не заметил».

Лечится двумя способами. Либо задача сама оборачивает своё тело в try/catch и логирует всё, что поймала, — самый надёжный вариант, он не зависит от способа отправки. Либо у пула переопределяют afterExecute (ThreadPoolExecutor даёт такой хук) и достают исключение из Future там.

ThreadFactory: имена потоков и обработчик

По умолчанию потоки пула называются pool-1-thread-3, и в дампе на пять пулов это пять одинаковых имён без подсказки, чей поток где. Имя задаёт ThreadFactory — заодно там же ставят обработчик ошибок и признак демона:

ThreadFactory factory = r -> {
    Thread t = new Thread(r, "report-" + counter.incrementAndGet());
    t.setUncaughtExceptionHandler((th, e) -> log.error("поток {} умер", th.getName(), e));
    t.setDaemon(false);
    return t;
};
ExecutorService pool = Executors.newFixedThreadPool(4, factory);

Понятное имя окупается на первом же разборе зависания: строка report-3 в дампе сразу говорит, какой пул встал. То же самое умеют готовые фабрики из библиотек (ThreadFactoryBuilder у Guava, CustomizableThreadFactory у Spring), а в Spring имя обычно задают одной настройкой — spring.task.execution.thread-name-prefix.

Фабрики Executors и их ограничения

Класс Executors предоставляет готовые фабрики, и у двух самых ходовых есть дыра, которая проявляется только под нагрузкой. У newFixedThreadPool очередь не ограничена (LinkedBlockingQueue без лимита): если задачи поступают быстрее, чем обрабатываются, очередь растёт до исчерпания памяти. У newCachedThreadPool нет верхнего предела числа потоков: при всплеске нагрузки система может создать тысячи потоков. Поэтому в продакшене берут явный ThreadPoolExecutor, а фабрики остаются для тестов и скриптов:

МетодПоведение
newFixedThreadPool(n)ровно n потоков, очередь не ограничена
newCachedThreadPool()потоки создаются по требованию, живут 60 с после простоя
newSingleThreadExecutor()один поток, задачи строго по очереди
newScheduledThreadPool(n)для задач с задержкой и по расписанию

У последнего есть своя тихая ловушка, которую стоит знать до того, как на нём построят фоновую работу. scheduleAtFixedRate(task, 0, 1, MINUTES) повторяет задачу по расписанию, scheduleWithFixedDelay — с паузой от конца до начала (разница та же, что у @Scheduled в Spring). А вот что происходит при ошибке: исключение, вылетевшее из периодической задачи, снимает её с расписания навсегда. Не логируется, не перезапускается — просто больше никогда не выполняется. Сервис при этом жив и здоров, и обнаруживается это через неделю по тому, что отчёты перестали приходить. Поэтому тело периодической задачи всегда оборачивают целиком:

scheduler.scheduleWithFixedDelay(() -> {
    try {
        rebuildReport();
    } catch (Exception e) {          // ловим всё: иначе расписание умрёт молча
        log.error("пересборка отчёта упала", e);
    }
}, 0, 1, TimeUnit.HOURS);

ThreadPoolExecutor: полный контроль

Это тот же пул, но опасные умолчания вынесены в параметры конструктора. Пул в примере намеренно тесный: так все четыре исхода видны на четырёх задачах.

живой пример

import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;

public class PoolRouting {
    public static void main(String[] args) throws Exception {
        ThreadPoolExecutor pool = new ThreadPoolExecutor(
                1, 2,                          // corePoolSize — постоянные потоки, maximumPoolSize — предел
                30, TimeUnit.SECONDS,          // keepAliveTime — сколько «лишний» поток живёт без работы
                new ArrayBlockingQueue<>(1),   // ограниченная очередь: одно место
                new ThreadPoolExecutor.CallerRunsPolicy()); // политика отказа

        for (int i = 1; i <= 4; i++) {
            int number = i;
            pool.execute(() -> {
                System.out.println("задача " + number + " — поток " + Thread.currentThread().getName());
                try {
                    Thread.sleep(300);
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                }
            });
        }

        pool.shutdown();
        pool.awaitTermination(3, TimeUnit.SECONDS);
    }
}
Запустить

Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →

Параметры работают вместе:

  1. Пока число потоков меньше corePoolSize — создаётся новый поток.
  2. Если corePoolSize достигнут, задача кладётся в очередь.
  3. Если очередь заполнена и потоков меньше maximumPoolSize — создаётся ещё один поток.
  4. Если очередь заполнена и потоков уже maximumPoolSize — срабатывает политика отказа.

Ещё один параметр, о котором узнают поздно: keepAliveTime по умолчанию действует только на потоки сверх ядра. Потоки ядра живут вечно, даже если задач не было неделю, и держат свои стеки и ThreadLocal-значения. Для редко нагружаемого пула это лишняя память; pool.allowCoreThreadTimeOut(true) распространяет таймаут и на них, и тогда простаивающий пул съёживается до нуля потоков, поднимая их заново при первой задаче.

Здесь же прячется ловушка newFixedThreadPool. У него corePoolSize равен maximumPoolSize, поэтому шаг 3 не срабатывает никогда: потоков сверх ядра не появится. А очередь у него неограниченная, так что и до шага 4 дело не доходит. Отсюда поведение, которое со стороны выглядит загадочно: под нагрузкой fixed-пул не разгоняется и ни от чего не отказывается — он просто молча копит задачи, пока не кончится память.

отправитель пул: ядро 1, очередь 1, максимум 2 поток 1 · задача 1 очередь — одно место в очереди · задача 2 место до максимума поток 2 · задача 3 задача 1 — ядро пусто задача 2 — ядро занято задача 3 — очередь полна задача 4 — мест нет CallerRunsPolicyвыполняет сам отправитель

Пул с одним потоком ядра, очередью на одно место и максимумом в два потока. Задача 1 попадает в поток ядра, задача 2 — в очередь, задача 3 поднимает второй поток, а задаче 4 места нет: срабатывает политика отказа, и при CallerRunsPolicy её выполняет тот поток, который её отправил.

Порядок строк в выводе от запуска к запуску меняется — потоки работают параллельно, — а маршрут задач один и тот же. Задача 4 напечатает main: пул её не принял, и CallerRunsPolicy выполнила её прямо в потоке отправителя. В сервисе числа другие — четыре потока ядра, очередь на 200, максимум восемь, — но маршрут тот же.

Политики отказа

RejectedExecutionHandler определяет, что делать с задачей, которую некуда поставить:

  • AbortPolicy (по умолчанию) — бросает RejectedExecutionException.
  • CallerRunsPolicy — задачу выполняет сам поток, который её отправил.
  • DiscardPolicy — задача молча отбрасывается.
  • DiscardOldestPolicy — из очереди удаляется самая старая задача, новая встаёт на её место.

Для большинства сервисов CallerRunsPolicy — разумный выбор: вместо потери задач система сама замедляет приём. Замедляет, правда, не абстрактно: задачу выполняет тот поток, который её отправил, а в Spring это обычно HTTP-поток Tomcat — тот самый, которому положено принимать запросы. Пока он занят чужой фоновой работой, приём стоит. И вторая деталь: после shutdown() эта политика задачу не выполняет, а молча отбрасывает.

Сколько потоков создавать

Универсальной формулы нет, но есть два полюса:

Задачи, упирающиеся в процессор (CPU-bound: вычисления, обработка данных без ожидания): размер пула ≈ числу ядер процессора. Лишние потоки только создают конкуренцию за CPU.

int cores = Runtime.getRuntime().availableProcessors();
ExecutorService cpuPool = Executors.newFixedThreadPool(cores);

Задачи, ждущие ввода-вывода (IO-bound: запросы к БД, HTTP-вызовы, чтение файлов): потоки большую часть времени ждут ответа. Пул может быть в несколько раз больше числа ядер — пока один поток ждёт IO, другие работают. Типичный ориентир «2–4× ядер» — лишь отправная точка.

Точнее считает закон Литтла: число задач, одновременно находящихся в работе, равно интенсивности их поступления, умноженной на среднее время обслуживания. Если сервис делает 600 HTTP-вызовов в секунду, а каждый ждёт ответа в среднем 0,4 с, то в полёте одновременно 600 × 0,4 = 240 вызовов — и пулу нужно около 240 потоков, чтобы очередь перед ним не росла. С восемью потоками «по два на ядро» остальные 230 задач будут стоять в очереди, латентность вырастет до секунд, а CPU останется свободным: потоки не считают, а ждут сеть. Число ядер здесь ни при чём: оно задаёт размер пула только для CPU-bound работы.

пул на 240 600 вызовов/с 240 в полёте очередь пуста пул на 8 600 вызовов/с 8 в работе 232 в очереди

Закон Литтла на числах статьи: 600 вызовов в секунду по 0,4 секунды ожидания дают 240 задач в полёте, и пул на 240 потоков держит очередь пустой, а пул на восемь оставляет 232 задачи ждать при свободном процессоре.

Корректное завершение пула

Пул — это ресурс, его нужно закрывать. Если пул не остановить, JVM не завершится, пока живы потоки пула (если они не демонические).

pool.shutdown();

try {
    // ждём не дольше 10 секунд
    if (!pool.awaitTermination(10, TimeUnit.SECONDS)) {
        pool.shutdownNow();
        // ещё подождём прерывания
        pool.awaitTermination(5, TimeUnit.SECONDS);
    }
} catch (InterruptedException e) {
    pool.shutdownNow();
    Thread.currentThread().interrupt(); // восстанавливаем флаг прерывания
}

Разница между методами:

  • shutdown() — «мягкая» остановка: новые задачи не принимаются, уже поставленные в очередь и выполняющиеся — доделываются.
  • shutdownNow() — «жёсткая»: вызывает interrupt() у всех потоков пула и возвращает список задач, которые не успели выполниться. Задача должна реагировать на прерывание, иначе поток продолжит работу.

С Java 21 у этого ритуала есть короткая запись: ExecutorService реализует AutoCloseable, и его close() делает ровно то же самое — перестаёт принимать задачи и ждёт, пока доработают уже принятые.

try (ExecutorService pool = Executors.newFixedThreadPool(4)) {
    pool.submit(() -> handle(request));
}   // на выходе из блока пул закроется сам и дождётся задач

Одна оговорка: close() ждёт без ограничения по времени. Если задача зависла, зависнет и выход из блока — там, где нужен потолок ожидания, остаётся развёрнутый вариант с awaitTermination.

shutdown() новые не берём очередь доработает потоки завершились shutdownNow() новые не берём очередь вернули потокам interrupt

Две ветки остановки на одной шкале: обе перестают принимать новые задачи, но shutdown доводит очередь до конца, а shutdownNow прерывает потоки и возвращает невыполненное списком.

Дополнительно: при первом чтении можно пропустить

Глубже: ThreadLocal в пуле: зачем, где течёт, чем заменитьрасширенное

ThreadLocal<T> даёт каждому потоку свою копию значения: set в одном потоке невидим другому. На нём стоят MDC в логах, контекст безопасности, текущая транзакция Spring: всё, что нужно «по пути запроса» без протаскивания параметром через сорок методов.

В пуле у него две беды. Поток живёт долго и обслуживает тысячи задач, а значение, которое задача положила и не убрала, достаётся следующей задаче: запрос другого пользователя увидит чужой userId в логах или чужую роль. Вторая беда медленнее: значение это объект, поток держит на него ссылку, и если значения тяжёлые, а потоков сотни, память течёт до перезапуска. Поэтому правило: set в начале задачи, remove в finally, всегда.

static final ThreadLocal<String> USER = new ThreadLocal<>();
void handle(Request r) {
    USER.set(r.userId());
    try { service.process(); } finally { USER.remove(); }
}

С виртуальными потоками проблема меняется: их миллионы, и миллион копий одного ThreadLocal это уже не мелочь. Для них в Java 21+ есть ScopedValue: значение привязано не к потоку, а к блоку кода (ScopedValue.where(USER, id).run(() -> ...)), оно неизменяемо, наследуется дочерними задачами в структурированной конкурентности и исчезает само по выходе из блока, убирать нечего. Задача «ThreadLocal в пуле потоков» в тренажёре как раз про забытый remove.

Глубже: ForkJoinPool и parallelStream: общий пул и кража работырасширенное

IntStream.range(0, n).parallel() выглядит как бесплатное ускорение, а в двух статьях фазы «общий пул» назван источником бед. Вот что за ним стоит. ForkJoinPool это пул для задач, которые делятся на подзадачи: большой массив режется пополам, половины пополам, и так до порога, а результаты собираются обратно. У каждого потока своя очередь задач, и свободный поток крадёт задачи из чужой очереди (work stealing), поэтому пул хорошо загружает все ядра на неравномерной работе.

Параллельные стримы и CompletableFuture.supplyAsync без своего исполнителя работают в одном общем экземпляре, ForkJoinPool.commonPool(), размером в число ядер минус один. Отсюда две беды. Одна блокирующая операция внутри parallelStream (запрос к базе, HTTP) занимает поток общего пула на всё время ожидания; несколько таких стримов в разных частях приложения исчерпывают пул, и тормозит всё, что его делит, включая чужой код. И пул общий для всех запросов сервера: параллельный стрим на двухстах потоках Tomcat не даёт двести параллелизмов, он выстраивает всех в очередь к семи потокам.

Правила короткие. parallel() только для чистых вычислений над большими коллекциями в памяти, без ввода-вывода внутри. Ввод-вывод в параллель через свой ExecutorService с понятным размером и именем, а с Java 21 через виртуальные потоки. Если параллельный стрим всё же нужен со своим пулом, его запускают внутри new ForkJoinPool(8).submit(() -> list.parallelStream()...).get(): стрим использует пул, из которого вызван.

Глубже: однопоточная модель: один поток на сущностьрасширенное

Всё в фазе про защиту общего состояния. Есть третий путь, кроме замков и неизменяемости: сделать так, чтобы состояние трогал ровно один поток, а остальные присылали ему сообщения.

Самая простая форма: Executors.newSingleThreadExecutor() владеет объектом, все операции над ним отправляются задачами в этот исполнитель и выполняются по очереди. Замков нет, гонок нет, порядок гарантирован. Масштабируется это партиционированием по ключу: заказы с чётным номером идут в один однопоточный исполнитель, с нечётным в другой; в общем виде executors[hash(key) % n]. Одна сущность всегда обрабатывается одним потоком, разные сущности параллельно.

private final ExecutorService[] lanes = IntStream.range(0, 8)
    .mapToObj(i -> Executors.newSingleThreadExecutor()).toArray(ExecutorService[]::new);

void apply(OrderEvent e) {
    lanes[Math.floorMod(e.orderId().hashCode(), lanes.length)].submit(() -> handle(e));
}

Ровно так устроены партиции Kafka с одним потребителем на партицию, обработчики событий в Node.js и актор-модель: актор это объект с почтовым ящиком, который обрабатывает сообщения по одному. Цена: операция, которой нужны две сущности из разных партиций, не может взять обе «под одним замком», и её проектируют как обмен сообщениями с промежуточными состояниями. Зато исчезает целый класс ошибок, и в системах, где событий много, а связей между сущностями мало, это обычный выбор.

Коротко

  • Создавать поток на каждую задачу расточительно — пул переиспользует потоки.
  • ExecutorService — главный интерфейс; execute для Runnable, submit для Callable/Future.
  • Фабрики Executors удобны для прототипов, но неограниченные очередь или число потоков опасны в продакшене: для сервисов берут ThreadPoolExecutor с явными параметрами.
  • Маршрут задачи: поток ядра → очередь → поток до максимума → политика отказа.
  • CPU-bound: пул ≈ числу ядер; IO-bound: пул считают по закону Литтла, а не по числу ядер.
  • Всегда завершать пул через shutdown() + awaitTermination() + shutdownNow() при таймауте.
  • ThreadLocal в пуле убирают в finally, иначе значение достаётся чужому запросу; для виртуальных потоков есть ScopedValue.
  • parallelStream и supplyAsync без исполнителя делят один общий ForkJoinPool на всё приложение: только чистые вычисления, ввод-вывод в свой пул.
  • Третий путь кроме замков и неизменяемости: один поток на сущность, партиционирование по ключу, как партиции Kafka и акторы.
  • Исключение из execute печатает стек в System.err и пересоздаёт поток, из submit — ложится в Future и пропадает, если get() не вызвать: тело задачи оборачивают в try/catch; invokeAll ждёт всех, invokeAny берёт первый успешный, get(timeout) и cancel(true) задают потолок и отмену. ThreadFactory даёт потокам понятные имена и обработчик ошибок; allowCoreThreadTimeOut(true) распускает простаивающее ядро; исключение из периодической задачи ScheduledExecutorService снимает её с расписания навсегда.

Что почитать дальше