Пулы потоков и Future
Создал пул с параметрами: базовых потоков 1, максимум 8, очередь безлимитная. Закинул 50 задач. Ожидание - пул вырастет до восьми. Реальность - в пуле остался ОДИН поток, а 35 задач лежат в очереди. С очередью на пять элементов тот же код тут же дал восемь потоков и пустую очередь.
Пулы и футуры - то, чем многопоточность реально пишется в проде, поэтому и спрашивают приземлённо: как настроишь пул, куда денется исключение из задачи, почему сервис встал, хотя всё вроде асинхронно.
// Формулировки: «fixed или cached?», «что будет с исключением в задаче?», «как правильно завершить пул?»
Как на самом деле растёт пул
Замер из начала объясняется одним правилом, которое почти никто не помнит: дополнительные потоки сверх базовых подключаются, только когда очередь ПЕРЕПОЛНЕНА. С безлимитной очередью она не переполняется никогда, поэтому пул навсегда остаётся размером с базовое число.
Отсюда и опасность готовых фабрик. Фиксированный пул создаётся с безлимитной очередью: задачи копятся часами, пока не кончится память, а метрика «потоков в пуле» при этом спокойная. Кэширующий пул наоборот - растёт без всякого предела и под постоянной нагрузкой заводит тысячи потоков.
Прод-пул создают явно: базовое и максимальное число потоков, ОГРАНИЧЕННАЯ очередь и осознанная политика отказа. Проверил и её: при переполнении по умолчанию прилетает RejectedExecutionException. Политика CallerRuns вместо отказа заставляет отправителя выполнить задачу самому - это естественный тормоз для производителя.
// Размер выбирают от профиля: для вычислительных задач около числа ядер, для ожидающих ввода-вывода заметно больше, потому что потоки в основном стоят.
new ThreadPoolExecutor(1, 8, 60, SECONDS, new LinkedBlockingQueue<>());
// 50 задач -> в пуле 1 поток, в очереди 35
new ThreadPoolExecutor(1, 8, 60, SECONDS, new ArrayBlockingQueue<>(5));
// те же 50 задач -> в пуле 8 потоков- правило роста
- потоки сверх базовых - только при переполненной очереди
- политика отказа
- что делать с задачей, когда всё занято
Куда девается исключение из задачи
Проверил на пуле из одного потока. Задача, отправленная через submit, бросила исключение - в консоли не появилось ничего. Ни стектрейса, ни строки в логе. Задача умерла молча, и узнать об этом можно только одним способом: вызвать get у полученного Future. Он отдал ExecutionException, а внутри - исходное исключение с его сообщением.
Тот же код через execute ведёт себя иначе: исключение дошло до обработчика необработанных исключений потока, и я увидел сообщение. Отсюда дисциплина - либо всегда забирать результат и обрабатывать, либо оборачивать тело задачи в try с логированием, либо использовать execute.
// Про отмену: cancel(true) шлёт прерывание, поэтому отмена работает только если задача к нему отзывчива - проверяет флаг и не глотает исключение прерывания. А завершение пула делают по шаблону: shutdown, затем ожидание с таймаутом, затем shutdownNow как запасной вариант. Незакрытый пул с обычными потоками держит виртуальную машину живой после выхода из main.
- ExecutionException
- обёртка, в которой get отдаёт исключение задачи
- мягкое завершение
- shutdown, ожидание с таймаутом, потом shutdownNow
CompletableFuture: конвейер и чужой пул
CompletableFuture превращает асинхронность в конвейер: thenApply преобразует результат, thenCompose присоединяет следующий асинхронный шаг (без него получилось бы будущее внутри будущего), thenCombine дожидается двух и сливает, allOf и anyOf работают с пачками.
Ошибки тоже часть конвейера. Метод exceptionally даёт запасное значение, handle обрабатывает результат и исключение одним обработчиком. И тут та же ловушка, что с Future: цепочка без терминальной обработки теряет исключение молча - я проверил, в консоли снова тишина. С exceptionally то же исключение превратилось в аккуратный запасной результат.
// Где всё это исполняется. Асинхронные варианты методов по умолчанию уходят в общий пул ForkJoinPool. Я напечатал его параллелизм на двухъядерной машине - ОДИН поток. Блокирующий вызов там останавливает буквально всё приложение, включая параллельные стримы и чужие асинхронные задачи. Блокирующее отправляют только в свой пул, вторым аргументом асинхронного метода.
- thenCompose
- присоединить следующий асинхронный шаг без вложенности
- общий пул
- дефолт асинхронных методов, на 2 ядрах это 1 поток
Как отвечать: «Как настроишь пул потоков для прод-сервиса?»
Явным ThreadPoolExecutor, а не готовыми фабриками - у них опасные дефолты. У фиксированного пула безлимитная очередь, и задачи копятся до нехватки памяти, у кэширующего нет предела на число потоков. Размер вывожу из профиля задач: для вычислительных около числа ядер, для ожидающих ввода-вывода заметно больше, потоки там в основном стоят. Очередь обязательно ограниченная - это обратное давление. И тут важная деталь, которую я проверял руками: дополнительные потоки сверх базовых подключаются только при ПЕРЕПОЛНЕННОЙ очереди. С безлимитной очередью пул с базовым размером один и максимумом восемь остался при одном потоке на пятидесяти задачах, а с очередью на пять сразу вырос до восьми. Политику отказа выбираю осознанно: по умолчанию это исключение, а CallerRuns заставляет отправителя выполнить задачу самому и естественно тормозит производителя. Плюс именованная фабрика потоков, чтобы дампы читались, и метрика на длину очереди - по ней видно приближение к пределу до того, как начнётся инцидент.
Сильный ответ: отказ от фабрик аргументирован конкретными дефолтами, размер выведен из профиля нагрузки, названо неочевидное правило роста пула с проверенными числами и добавлена эксплуатационная часть. Так отвечает человек, чей пул уже горел.
На чём валят
- −Задают максимум потоков и безлимитную очередь. Пул навсегда остаётся базового размера: у меня 1 поток на 50 задачах.
- −Берут кэширующий пул под постоянную нагрузку. Он растёт без предела до тысяч потоков.
- −Делают submit и не забирают результат. Исключение задачи не увидит никто, в логах тишина.
- −Пускают блокирующий вызов в общий пул через асинхронный метод. На двух ядрах там всего один поток.
- −Не закрывают пул. Обычные потоки держат виртуальную машину живой после выхода из main.
Проверьте себя
Пять вопросов из банка по этой подтеме. Всего их 16, остальные разбираются в тренажёре.
- Чем Callable отличается от Runnable?A)Callable выполняется синхронно в текущем потоке, а Runnable — обязательно в новом отдельном потокеB)Runnable возвращает значение через return, а Callable — только через побочный эффект в общее полеC)Callable возвращает результат и может бросить проверяемое исключениеD)Разницы нет: оба интерфейса объявляют метод run() и применяются в ExecutorService одинаково
показать ответ и разбор
+C)Callable возвращает результат и может бросить проверяемое исключение// разбор: Runnable.run() ничего не возвращает и не объявляет проверяемых исключений. Callable<V>.call() ВОЗВРАЩАЕТ значение V и может бросить checked-исключение. Поэтому для задач с результатом (или требующих проброса исключения) в ExecutorService.submit передают Callable и получают Future<V>. Runnable — для задач без результата. execute() принимает только Runnable, submit() — оба.
- Что делает future.get()?A)Мгновенно возвращает результат, а если он ещё не готов — отдаёт null и продолжает без ожиданияB)Запускает задачу на выполнение: до вызова get() отправленная в пул задача остаётся неактивнойC)Отменяет задачу и возвращает частичный результат, вычисленный к этому моменту в фоновом потокеD)Блокирует поток, пока результат не готов, и возвращает его (или бросает исключение задачи)
показать ответ и разбор
+D)Блокирует поток, пока результат не готов, и возвращает его (или бросает исключение задачи)// разбор: future.get() БЛОКИРУЕТ вызывающий поток до готовности результата и возвращает его; если задача бросила исключение — get() кинет ExecutionException (обёрнутое), при отмене — CancellationException, при прерывании — InterruptedException. Есть вариант с таймаутом get(t, unit). Блокирующая природа — минус: несколько последовательных get() сериализуют работу; для неблокирующей композиции берут CompletableFuture.
- Что даёт CompletableFuture по сравнению с обычным Future?A)Неблокирующую композицию: колбэки и цепочки (thenApply/thenCompose/thenCombine)B)Он обеспечивает, что задача выполнится быстрее, так как использует более приоритетный пул потоковC)Он отличается лишь тем, что get() у него не бросает исключений, упрощая обработку ошибокD)Он выполняет задачи строго последовательно, тогда как Future умеет запускать их параллельно
показать ответ и разбор
+A)Неблокирующую композицию: колбэки и цепочки (thenApply/thenCompose/thenCombine)// разбор: CompletableFuture — Future с возможностью КОМПОЗИЦИИ без блокирующего get(): навесить продолжение (thenApply — трансформировать, thenCompose — связать со следующей async-операцией, thenCombine — объединить две), обработать ошибку (exceptionally/handle), задать пул (…Async(…, executor)). Это декларативные асинхронные конвейеры вместо ручного ожидания. Основа реактивно-подобного кода на стандартной библиотеке.
- Когда брать thenCompose вместо thenApply у CompletableFuture?A)Когда нужно выполнить трансформацию строго в том же потоке, а не в отдельном пуле исполнителяB)Когда функция сама возвращает CompletableFuture (иначе получится вложенность)C)Когда результат — простое значение: thenCompose для примитивов, а thenApply только для объектовD)Когда требуется обработать исключение предыдущего шага — thenCompose ловит ошибки, а thenApply нет
показать ответ и разбор
+B)Когда функция сама возвращает CompletableFuture (иначе получится вложенность)// разбор: thenApply — как map: берёт результат и трансформирует в обычное значение (T -> U), результат оборачивается в CompletableFuture автоматически. thenCompose — как flatMap: применяется, когда функция САМА возвращает CompletableFuture (T -> CompletableFuture<U>), и «сплющивает» CompletableFuture<CompletableFuture<U>> в CompletableFuture<U>. Правило: цепляешь ещё одну async-операцию — thenCompose; просто преобразуешь готовое значение — thenApply.
- Чем shutdown() отличается от shutdownNow() у ExecutorService?A)Shutdown() мгновенно убивает все потоки, а shutdownNow() дожидается завершения всех текущих задачB)Оба метода одинаковы — просто shutdownNow() короче писать, а поведение у них идентичноеC)shutdown() доисполняет принятые задачи; shutdownNow() пытается прервать выполняемыеD)Shutdown() освобождает память пула, но потоки продолжают принимать и выполнять новые задачи как прежде
показать ответ и разбор
+C)shutdown() доисполняет принятые задачи; shutdownNow() пытается прервать выполняемые// разбор: shutdown() — мягкое завершение: пул перестаёт принимать НОВЫЕ задачи, но доисполняет уже принятые (в работе и в очереди), затем гасит потоки. shutdownNow() — жёсткое: пытается прервать (interrupt) выполняемые задачи и возвращает список тех, что не успели стартовать. После любого shutdown новые submit отклоняются. Обычно: shutdown() + awaitTermination(timeout), а если не уложился — shutdownNow().
дальше
Теорию прочитали. Навык ставится повторением
В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.