Executors и пулы потоков
Фреймворк Executor, Runnable против Callable, подбор размера пула, политики отказа, CompletableFuture и ForkJoinPool.
8 вопросов
JuniorТеорияОчень частоЧто даёт фреймворк Executor по сравнению с new Thread(...).start()?
Что даёт фреймворк Executor по сравнению с new Thread(...).start()?
Он отделяет отправку задачи от потока, который её выполняет. Пул переиспользует небольшой набор потоков вместо потока на задачу, ставит работу в очередь, пока они заняты, возвращает Future на результат и добавляет управление через shutdown/awaitTermination.
Типичные ошибки
- ✗Думать, что пул создаёт новый поток на задачу, а не переиспользует фиксированный набор
- ✗Считать, что Executor лишь добавляет
Futureи ничего не меняет в потоках - ✗Забывать про
shutdownпула, оставляя утечку недемонных потоков
Уточняющие вопросы
- →Что станет с потоками пула, если никогда не вызвать
shutdown? - →Почему фиксированный пул безопаснее неограниченного под нагрузкой?
MiddleКодОчень частоКомпозиция двух зависимых асинхронных этапов через CompletableFuture
Композиция двух зависимых асинхронных этапов через CompletableFuture
CompletableFuture — это Future, который можно завершить самому и выстроить в цепочку этапов. thenApply преобразует готовое значение обычной функцией; thenCompose присоединяет этап, сам возвращающий CompletableFuture, и разворачивает результат, поэтому зависимые асинхронные вызовы идут по порядку без блокирующего get().
Типичные ошибки
- ✗Вызывать
get()между этапами, превращая асинхронную цепочку обратно в блокирующий код - ✗Применять
thenApply, когда функция возвращает future, получая вложенныйCompletableFuture - ✗Считать, что
CompletableFutureнельзя завершить вручную или скомпоновать
Уточняющие вопросы
- →Когда вы возьмёте
thenComposeвместоthenApply? - →Чем
thenCombineотличается отthenComposeпри соединении двух future?
JuniorТеорияЧастоЧем отличаются Runnable и Callable, и submit() против execute()?
Чем отличаются Runnable и Callable, и submit() против execute()?
Callable<V> возвращает значение и может бросить проверяемое исключение; Runnable — ни то, ни другое. Ловушка: submit() прячет брошенное исключение в Future и раскрывает его лишь при get(); execute() пропускает его в обработчик неперехваченных исключений потока.
Типичные ошибки
- ✗Считать, что задача из
submit()сама логирует своё исключение - ✗Думать, что
Runnableможет вернуть значение или бросить проверяемое исключение - ✗Ждать, что
execute()иsubmit()сообщают о сбоях одинаково
Уточняющие вопросы
- →Почему задача, отправленная через
submit(), будто падает молча? - →Как направить неперехваченные исключения задач в единый обработчик?
MiddleТеорияИногдаКак work-stealing в ForkJoinPool планирует подзадачи?
Как work-stealing в ForkJoinPool планирует подзадачи?
У каждого рабочего потока своя двусторонняя очередь (deque). fork() кладёт подзадачу в собственный deque владельца, и владелец снимает её LIFO — сперва самую свежую, «горячую» в кэше. Простаивающий поток крадёт с другого конца deque занятого потока, забирая самую старую и крупную задачу, чем минимизирует борьбу с владельцем. Центральной очереди, которая стала бы узким местом, нет.
Типичные ошибки
- ✗Представлять одну общую центральную очередь вместо deque у каждого потока
- ✗Думать, что владелец и вор берут работу с одного и того же конца deque
- ✗Считать, что поток, заблокированный в
join(), просто простаивает, а не выполняет другую работу
Уточняющие вопросы
- →Почему владелец снимает задачи LIFO, а вор крадёт с противоположного конца?
- →Что делает
ManagedBlocker, когда задаче вForkJoinPoolприходится блокироваться?
MiddleТеорияИногдаЧто происходит, когда очередь задач ThreadPoolExecutor заполнена?
Что происходит, когда очередь задач ThreadPoolExecutor заполнена?
Сначала executor растит пул до maximumPoolSize; и лишь когда очередь заполнена и этот максимум достигнут, задача уходит в RejectedExecutionHandler. Встроенных политик четыре: AbortPolicy (бросает RejectedExecutionException — по умолчанию), CallerRunsPolicy (выполняет задачу на отправившем потоке, тормозя производителя), DiscardPolicy и DiscardOldestPolicy.
Типичные ошибки
- ✗Ждать, что
execute()заблокируется на полной очереди, а не откажет в задаче - ✗Считать, что
maximumPoolSizeдостигается до заполнения очереди, а не после - ✗Полагать, что политика по умолчанию молча отбрасывает задачу, а не бросает исключение
Уточняющие вопросы
- →Почему
CallerRunsPolicyработает как backpressure для потока-производителя? - →При неограниченной очереди какая политика отказа вообще срабатывает?
MiddleДизайнИногдаВаша команда держит Spring Boot-сервис в контейнере с 8 vCPU. Два вида работы делят один ThreadPoolExecutor: пересжатие изображений — чистая нагрузка на CPU — и выгрузка заказов, которая около 90% времени простаивает в блокировке на медленном стороннем HTTP API. Сейчас это Executors.newFixedThreadPool(8) с неограниченной LinkedBlockingQueue. Под нагрузкой выгрузки стоят в очереди минутами, а загрузка CPU держится около 30%; один всплеск трафика довёл JVM до OutOfMemoryError. На design review обоснуйте, какую конфигурацию пула вы предлагаете — core и maximum size, тип и границу очереди, как разделить два вида работы — и оправдайте каждое число характером задач.
Ваша команда держит Spring Boot-сервис в контейнере с 8 vCPU. Два вида работы делят один ThreadPoolExecutor: пересжатие изображений — чистая нагрузка на CPU — и выгрузка заказов, которая около 90% времени простаивает в блокировке на медленном стороннем HTTP API. Сейчас это Executors.newFixedThreadPool(8) с неограниченной LinkedBlockingQueue. Под нагрузкой выгрузки стоят в очереди минутами, а загрузка CPU держится около 30%; один всплеск трафика довёл JVM до OutOfMemoryError. На design review обоснуйте, какую конфигурацию пула вы предлагаете — core и maximum size, тип и границу очереди, как разделить два вида работы — и оправдайте каждое число характером задач.
Размер выбирается по доле блокировки, а не по единому правилу: потоков ≈ ядра × (1 + ожидание/работа). Пересжатию, упирающемуся в CPU, нужно около availableProcessors() потоков; выгрузке, заблокированной ~90% времени, — примерно вдесятеро больше. Разнесите их по разным пулам, чтобы ни один не морил голодом другой. Неограниченную очередь замените ограниченной — неограниченная молча вбирает всплеск в кучу, что и дало OutOfMemoryError, — и задайте явную политику отказа.
Типичные ошибки
- ✗Брать размер пула по числу ядер независимо от того, сколько его задачи блокируются
- ✗Оставлять очередь неограниченной, так что всплеск уходит в кучу вместо отказа
- ✗Делить один пул между CPU-задачами и блокирующими, позволяя одним морить голодом другие
Уточняющие вопросы
- →При неограниченной очереди когда вообще срабатывает
maximumPoolSize? - →Как выбрать между отказом в задаче и торможением отправляющего потока?
SeniorДизайнИногдаСервис приёма событий принимает webhook-вызовы по HTTP и передаёт каждый в ThreadPoolExecutor на обогащение — два блокирующих вызова примерно по 400 мс — перед записью в Kafka. Производители внешние и при любой ошибке агрессивно повторяют запрос. В штатном режиме сервис держит 2000 событий/с; во время переигрывания upstream'а темп ненадолго доходит до 20 000 событий/с. Сейчас у пула 64 потока и неограниченная LinkedBlockingQueue: при последнем переигрывании очередь выросла до миллионов записей, латентность поднялась до минут, куча заполнилась, и процесс умер, потеряв всё, что стояло в очереди. Спроектируйте стратегию backpressure для следующего релиза: где проходит граница, что делает executor при насыщении и что HTTP-слой отвечает производителю, которого не может обслужить. Обоснуйте, почему нельзя просто увеличить пул.
Сервис приёма событий принимает webhook-вызовы по HTTP и передаёт каждый в ThreadPoolExecutor на обогащение — два блокирующих вызова примерно по 400 мс — перед записью в Kafka. Производители внешние и при любой ошибке агрессивно повторяют запрос. В штатном режиме сервис держит 2000 событий/с; во время переигрывания upstream'а темп ненадолго доходит до 20 000 событий/с. Сейчас у пула 64 потока и неограниченная LinkedBlockingQueue: при последнем переигрывании очередь выросла до миллионов записей, латентность поднялась до минут, куча заполнилась, и процесс умер, потеряв всё, что стояло в очереди. Спроектируйте стратегию backpressure для следующего релиза: где проходит граница, что делает executor при насыщении и что HTTP-слой отвечает производителю, которого не может обслужить. Обоснуйте, почему нельзя просто увеличить пул.
Ограничьте буфер и верните давление производителю, а не в кучу. Задайте ограниченную очередь под ту латентность, которую готовы терпеть, и выберите политику отказа, сигнализирующую о насыщении: AbortPolicy с отображением в HTTP 429/503 и Retry-After либо CallerRunsPolicy, тормозящую принимающий поток. Увеличение пула не помогает: работа упирается в I/O и ограничена вызовами вниз по цепочке, поэтому лишние потоки лишь углубляют очередь. Надёжный буфер — это брокер, а не куча.
Типичные ошибки
- ✗Добавлять потоки, чтобы поглотить всплеск, чей реальный потолок — зависимость вниз по цепочке
- ✗Оставлять очередь неограниченной, превращая проблему темпа в
OutOfMemoryError - ✗Принимать работу, которую не можешь обслужить, вместо сигнала о насыщении производителю
Уточняющие вопросы
- →Почему
CallerRunsPolicyтормозит производителя, а не просто переносит работу? - →Когда надёжным буфером должен быть брокер, а не собственная очередь executor'а?
SeniorТеорияИногдаЧто вызывает голодание потоков в ограниченном пуле?
Что вызывает голодание потоков в ограниченном пуле?
Все потоки пула заблокированы в ожидании работы, которую может выполнить только этот же пул, — прогрессировать некому: взаимоблокировка по исчерпанию. Классика: задача отправляет подзадачу в собственный пул и блокируется на Future.get(); при достаточном их числе все потоки стоят на get(), а подзадачи вечно ждут в очереди. Долгие блокирующие I/O или один медленный потребитель, захвативший общий пул, морят его голодом так же. Лечение: не делать join на своём пуле и разносить нагрузки по разным пулам.
Типичные ошибки
- ✗Путать голодание с просто маленьким пулом и лечить его увеличением размера
- ✗Отправлять подзадачу в собственный пул и затем блокироваться на
Future.get() - ✗Считать, что
Future.get()на время ожидания возвращает рабочий поток в пул
Уточняющие вопросы
- →Почему
ForkJoinPool.join()избегает голодания, которое обычный пул ловит на вложенности? - →Как разнесение нагрузок по разным пулам сдерживает одного медленного потребителя?