В коде ниже цепочка создаётся до завершения исходного этапа. Какой поток вправе выполнить тело thenApply?
import java.util.concurrent.CompletableFuture;
class Demo {
public static void main(String[] args) throws Exception {
CompletableFuture<String> source = new CompletableFuture<>();
CompletableFuture<String> result = source.thenApply(value -> {
System.out.println(Thread.currentThread().getName());
return value.toUpperCase();
});
Thread producer = new Thread(() -> source.complete("ok"), "producer");
producer.start();
producer.join();
result.join();
}
}
thenApply не обязан создавать отдельный поток. В данном сценарии обработчик обычно выполняется потоком, который завершает source, то есть producer, но API не следует использовать как гарантию конкретного имени потока.
Для принудительного выполнения через асинхронный исполнитель применяют thenApplyAsync. Это важно, потому что обычный thenApply может выполнить тяжёлую или блокирующую работу прямо внутри потока, завершающего предыдущий этап.
CompletableFuture появился в Java 8 как средство построения цепочек асинхронных операций без ручной координации потоков, флагов готовности и блокирующих ожиданий. Одной из задач было отделить описание этапов вычисления от непосредственного управления потоками.
Поэтому API разделяет обычные методы вроде thenApply и асинхронные варианты вроде thenApplyAsync. Такое разделение позволяет выбрать между минимальными накладными расходами и передачей продолжения исполнителю.
Вызов thenApply часто ошибочно трактуют как обязательный запуск вычисления в новом потоке. Если предыдущий этап завершается в потоке HTTP-запроса, потоке общего пула или потоке, удерживающем важный ресурс, продолжение может выполниться там же.
Неверный выбор приводит к блокировке рабочего потока, снижению пропускной способности и взаимному влиянию несвязанных задач. Особенно опасны внутри продолжения сетевые вызовы, операции с диском, ожидание join() или длительные вычисления.
Для неасинхронных методов CompletableFuture обработчик может выполнить поток, который завершает предыдущий этап, либо другой поток, вызывающий операцию завершения. Если этап уже завершён к моменту регистрации обработчика, продолжение может выполниться непосредственно в потоке, вызывающем thenApply.
Следовательно, точнее говорить не «thenApply выполняется в потоке producer», а «thenApply не гарантирует выделенный поток и допускает выполнение в потоке завершения». В показанном коде source незавершён при регистрации, поэтому при обычном завершении через source.complete обработчик выполняется в рамках завершения producer.
thenApplyAsync передаёт обработчик исполнителю. В перегрузке без явного Executor используется стандартный исполнитель CompletableFuture, обычно основанный на ForkJoinPool.commonPool(), а точный выбор следует проверять с учётом документации и среды выполнения. В перегрузке с Executor поток выполнения контролируется явно.
Асинхронный вариант не делает операцию автоматически безопасной, быстрой или неблокирующей: он лишь меняет способ запуска продолжения. Если исполнитель имеет мало потоков, блокирующие этапы всё равно могут исчерпать его ресурсы.
Для CPU-bound работы подходит ограниченный пул, соответствующий вычислительным ресурсам. Для блокирующих операций обычно выделяют отдельный исполнитель, чтобы ожидание не задерживало независимые вычисления; размер пула и политика очереди должны учитывать нагрузку и необходимость ограничения backpressure.
Сервис получает запрос, создаёт CompletableFuture и в thenApply выполняет обращение к медленной файловой системе. Предыдущий этап завершается в потоке веб-сервера, поэтому этот поток начинает ждать диск и перестаёт обслуживать новые запросы.
Вариант с оставлением thenApply имеет мало накладных расходов, но связывает продолжение с потоком завершения и плохо изолирует блокирующую работу. Вариант с thenApplyAsync без собственного исполнителя проще, однако общая очередь может быть перегружена другими задачами.
Практически выбирают отдельный ограниченный Executor для блокирующего этапа, а для коротких преобразований оставляют thenApply. В результате веб-потоки быстрее освобождаются, а нагрузку на внешнюю систему можно контролировать размером пула и очереди.
Что произойдёт, если source уже завершён до вызова thenApply?
Продолжение может выполниться синхронно в потоке, который регистрирует обработчик. Поэтому даже вызов thenApply в обычном пользовательском потоке способен непосредственно запустить тяжёлую работу, если исходный этап уже готов. Нельзя определять модель выполнения только по месту, где был создан source.
Гарантирует ли thenApplyAsync новый поток для каждого продолжения?
Нет. Он гарантирует передачу задачи асинхронному исполнителю, но исполнитель может переиспользовать существующий поток, поставить задачу в очередь или ограничить параллелизм. Слово Async означает асинхронную передачу выполнения, а не обязательное создание нового потока.
Что изменится, если продолжение завершится исключением?
Исключение из функции thenApply делает результирующий CompletableFuture завершённым ошибочно; оно не обязано немедленно выброситься в потоке, где возникло. Ошибку можно обработать через exceptionally, handle или whenComplete, а вызов join() обычно сообщит её как CompletionException с исходной причиной.
Это означает, что обработка ошибок является частью цепочки, а не только локальным try-catch вокруг регистрации продолжения. Если исключительный результат не проверить и не дождаться, ошибка может остаться незамеченной.