сеньорчикОткрыть в Telegram
← вся теориятеория к собесу · Многопоточность в Java

Пулы потоков и 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, остальные разбираются в тренажёре.

  1. #executors_futures1 / 5
    Чем 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() — оба.

  2. #executors_futures2 / 5
    Что делает future.get()?
    A)Мгновенно возвращает результат, а если он ещё не готов — отдаёт null и продолжает без ожидания
    B)Запускает задачу на выполнение: до вызова get() отправленная в пул задача остаётся неактивной
    C)Отменяет задачу и возвращает частичный результат, вычисленный к этому моменту в фоновом потоке
    D)Блокирует поток, пока результат не готов, и возвращает его (или бросает исключение задачи)
    показать ответ и разбор
    +D)Блокирует поток, пока результат не готов, и возвращает его (или бросает исключение задачи)

    // разбор: future.get() БЛОКИРУЕТ вызывающий поток до готовности результата и возвращает его; если задача бросила исключение — get() кинет ExecutionException (обёрнутое), при отмене — CancellationException, при прерывании — InterruptedException. Есть вариант с таймаутом get(t, unit). Блокирующая природа — минус: несколько последовательных get() сериализуют работу; для неблокирующей композиции берут CompletableFuture.

  3. #executors_futures3 / 5
    Что даёт 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)). Это декларативные асинхронные конвейеры вместо ручного ожидания. Основа реактивно-подобного кода на стандартной библиотеке.

  4. #executors_futures4 / 5
    Когда брать 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.

  5. #executors_futures5 / 5
    Чем 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().

дальше

Теорию прочитали. Навык ставится повторением

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