Создавать новый поток на каждую задачу — дорого и опасно. 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);
}
}
Запустить
Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →
Параметры работают вместе:
- Пока число потоков меньше
corePoolSize— создаётся новый поток. - Если
corePoolSizeдостигнут, задача кладётся в очередь. - Если очередь заполнена и потоков меньше
maximumPoolSize— создаётся ещё один поток. - Если очередь заполнена и потоков уже
maximumPoolSize— срабатывает политика отказа.
Ещё один параметр, о котором узнают поздно: keepAliveTime по умолчанию действует только на потоки сверх ядра. Потоки ядра живут вечно, даже если задач не было неделю, и держат свои стеки и ThreadLocal-значения. Для редко нагружаемого пула это лишняя память; pool.allowCoreThreadTimeOut(true) распространяет таймаут и на них, и тогда простаивающий пул съёживается до нуля потоков, поднимая их заново при первой задаче.
Здесь же прячется ловушка newFixedThreadPool. У него corePoolSize равен maximumPoolSize, поэтому шаг 3 не срабатывает никогда: потоков сверх ядра не появится. А очередь у него неограниченная, так что и до шага 4 дело не доходит. Отсюда поведение, которое со стороны выглядит загадочно: под нагрузкой fixed-пул не разгоняется и ни от чего не отказывается — он просто молча копит задачи, пока не кончится память.
Пул с одним потоком ядра, очередью на одно место и максимумом в два потока. Задача 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 работы.
Закон Литтла на числах статьи: 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 прерывает потоки и возвращает невыполненное списком.
Глубже: 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снимает её с расписания навсегда.
Что почитать дальше
- Потоки и процессы в Java — что такое поток и почему он дорогой.
- CompletableFuture: асинхронные цепочки — как строить зависимые асинхронные шаги без ручного ожидания
Future. - Виртуальные потоки — альтернатива пулам для IO-bound задач в Java 21.
- Потокобезопасные коллекции — какие структуры данных безопасно использовать между потоками пула.