Java Concurrency: потоки, асинхронность и виртуальные потоки.

Я уже почти год назад законспектировал основные способы сделать код на Java подходящим для многопоточности. Барьеры, локи и синхронизация – всё это было, а вот главный вопрос, откуда берутся потоки, кто и как делает многопоточность в приложении, я тогда не задавал. Поэтому в этой статье опишу способы создавать и выполнять потоки в джаве, а не частные случаи их взаимодействия.

Начну поэтому с самого потока (и немного с задачи). Задача и поток – разные понятия. Задача представляет собой код, который необходимо выполнить, а поток – механизм, который этот код выполняет.

Здесь мне на ум приходит ещё классификация Флинна, которая для архитектуры процессоров выделяет SISD (Single Instruction, Single Data), SIMD (Single Instruction, Multiple Data), MISD и MIMD. Прямой связи между этими понятиями нет, но мне кажется, что аппаратная архитектура здесь несколькок перекликается с программной архитектурой.

Потоки и задачи в Java

В джаве задача описывается интерфейсом Runnable, с единственным методом run(). Он является точкой входа для выполнения задачи.

Связать задачу с потоком можно двумя способами: реализовать интерфейс Runnable и передать этот объект в конструктор Thread, либо унаследоваться от класса Thread и переопределить метод run(). Первый вариант считается более изящным, потому что он отделяет описание задачи от механизма её выполнения.

После создания поток запускается методом start(). Если просто вызвать метод run() – код выполнится синхронно в том же потоке. А вот start() создаёт новый поток выполнения, который уже самостоятельно вызывает run().

У потока есть имя, есть приоритет (который не всегда может повлиять на ситуацию), и они могут быть deamon (так сказать, фоновыми) или non-deamon (обычные потоки не дают JVM завершить работу, пока они ещё выполняются). Есть несколько состояний, и важно понимать, что завершённый поток уже нельзя «возродить» и запустить заново. Даже если мы хотим повторить тот же код, который он только что прогнал. Придётся делать новый поток, и тут уже виднеется проблема, которую решают сервисы.

Выполнение задач с ExecutorService и Future

Вместо того, чтобы каждый раз создавать новые потоки, можно переиспользовать уже созданные после того, как они справились с текущей задачей. Именно на этой идее основаны пулы потоков (Thread Pool).

Пул содержит фиксированное или изменяемое количество потоков. Новый задачи помещаются в очередь, откуда они попадают на выполнение к свободным потокам. Получается очередь на обслуживание на почте: окошек фиксированное число, а посетителей может быть как и больше, так и меньше. Задачи могут выполняться или одни и те же по запросам от разных клиентов, или очень разные – потокам без разницы, что они делают.

Но кто-то должен заведовать всей этой логикой: распределять задачи по потокам, хранить весь пул и т. д. Для этого используют Executor и его наследников. Сам интерфейс пока предлагает только один метод execute() для выполнения задачи и может выполнять её в текущем потоке, но уже конкретные реализации позволяют управлять пулами.

public interface Executor {
void execute(Runnable command);
}

Но уже его наследник – интерфейс ExecutorService обещает намного больше, например запускать несколько задач сразу, отменять их выполнение и возвращать результат выполнения во Future – чтобы получать значение, необходимо код выполнять в методе submit().

Итак execute() принимает задачу и передаёт её на выполнение, но не позволяет узнать, когда она завершится и был ли получен какой-либо результат. Поэтому на практике часто используют именно ExecutorService и егоsubmit().

Но интерфейс Runnable подходит только для задач, которые ничего не возвращают. Если необходимо получить результат работы, используется интерфейс Callable. Вместо run() здесь метод call().

Future<String> future = executorService.submit(callable);

Задача помещается в очередь пула, а метод сразу возвращает объект Future, не дожидаясь окончания вычислений. То есть такой код может быть не просто многопоточным (executor исполняет задачи в отдельных потоках), но и асинхронным (мы можем не ждать результата вычислений и продолжать работу в текущем потоке).

Объект Future позволяет узнать информацию о задаче. Его основные методы:

isDone() – завершилась ли задача;
isCancelled() – была ли задача отменена;
cancel() – запросить отмену выполнения;
get() – получить результат;
get(timeout, unit) – дождаться результата не дольше указанного времени.

Метод get() возвращает результат вычислений, но если задача ещё не завершилась, вызывающий поток будет заблокирован до окончания её выполнения. То есть по механике в асинхронном коде он работает, как join() или барьер. Поэтому часто вызов get() слишком рано лишает программу преимуществ асинхронности.

Runnable тоже можно отправлять в ExecutorService, так как метод submit() умеет работать и с н им. Но тогда get() вернёт просто null, так как Runnable не может вернуть свой результат. Если null не хочется, то можно добавить вторым аргументом значение, которое вернётся после выполнения задачи.

Future<?> futureNull = executor.submit(runnable);
Future<String> futureDone = executor.submit(runnable, "Done");

Итого после завершения задачи возможны несколько вариантов:

  • возвращается вычисленный результат (или null);
  • выбрасывается ExecutionException, если сама задача завершилась с ошибкой;
  • выбрасывается InterruptedException, если поток, ожидающий результат, был прерван;
  • выбрасывается CancellationException, если задача была отменена.

Важно понимать, что ExecutionException лишь оборачивает настоящее исключение, возникшее внутри задачи. Получить исходную причину можно через метод getCause().

Отменить выполнение задачи можно с методом cancel():

future.cancel(boolean mayInterruptIfRunning)

Если задача ещё находится в очереди на выполнение, она просто удаляется оттуда и не будет выполнена. Если же задача уже выполняется, параметр mayInterruptIfRunning определяет, следует ли отправить рабочему потоку сигнал interrupt(). Это не принудительная остановка – задача должна самостоятельно реагировать на прерывание. Если код игнорирует interrupt(), выполнение продолжится.

После отмены вызов get() всегда приводит к возникновению CancellationException, даже если задача проигнорировала запрос отмены и успела завершиться самостоятельно.

Реализации ExecutorService

Реализация этого интерфейса ThreadPoolExecutor используется чаще всего. Можно её создавать через вспомогательный класс Executors:

ExecutorService executor = Executors.newFixedThreadPool(4);

Или можно самостоятельно:

ExecutorService executor = new ThreadPoolExecutor(4,8, 60,TimeUnit.SECONDS, new ArrayBlockingQueue <>(100));

Здесь аргументы по порядку:

  • corePoolSize – базовое количество рабочих потоков, которое пул старается использовать до помещения новых задач в очередь. Эти потоки обычно создаются по мере необходимости; заранее создать их можно отдельно.
  • maximumPoolSize – максимальное число потоков,
  • keepAliveTime – через сколько времени можно удалить «лишний» поток, который уже не используется
  • TimeUnit.SECONDS – в каких единицах измеряется keepAliveTime.
  • И очередь, куда складывать задачи, если занято corePoolSize потоков.

Пока потоков меньше corePoolSize, создаются новые потоки. После достижения corePoolSize задачи отправляются в очередь. Новые потоки сверх corePoolSize до maximumPoolSize создаются только тогда, когда очередь не может принять задачу (в примере мы разрешили класть туда 100 задач). После достижения maximumPoolSize и заполнения очереди задача отклоняется – мы получаем RejectedExecutionException .

То есть задачи не обязательно выполняются по порядку. Если четыре основных потока заняты первыми четырьмя задачами, и в очереди уже лежат 100 задач, а потом приходит ещё одна, 105-я по счёту, то ThreadPoolExecutor снова попытается положить её в очередь, но workQueue.offer(task) уже вернёт false, и тогда она выполнится на дополнительном потоке – раньше всех ста из очереди. В этом особенность ThreadPoolExecutor и асихнронности – мы не полагаемся на порядок выполнения задач. Но если задачи сверх ёмкости очереди не приходят, то принцип FIFO для задач в очереди всё же сохраняется – оттуда они достаются по порядку.

Реализации очередей BlockingQueue для хранения задач

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

while (true) {
Runnable task = queue.poll();
if (task != null) {
task.run();
}
}

LinkedBlockingQueue хранит задачи в связанном списке. Позволяет удобно добавлять задачи, может быть ограниченной или практически неограниченной, подходит для очереди FIFO. Но может забирать много памяти, если задачи поступают слишком быстро, а мы не утсановили ограничение:

BlockingQueue<Runnable> queue = new LinkedBlockingQueue<>();

ArrayBlockingQueue – ограниченная очередь на основе массива. Поэтому размер всегда ограничен и потребление памяти предсказуемое. Но не гибкая структура – надо заранее определиться с ёмкостью.

BlockingQueue<Runnable> queue = new ArrayBlockingQueue<>(100);

SynchronousQueue – очередь, которая ничего не хранит и почти сразу направляет задачи на выполнение. Или создаётся новый поток или задача отклоняется – поэтому тут очень важно установить разумный максимальный размер пула.

BlockingQueue<Runnable> queue = new SynchronousQueue<>();

PriorityBlockingQueue – выполняет задачи не обязательно в порядке поступления, а по приоритету.

BlockingQueue<Runnable> queue = new PriorityBlockingQueue<>();


Но у обычного Runnable нет приоритета, поэтому нужен специальный класс, который будет хранить приоритет и позволять их сравнивать. Например:

public class PriorityTask implements Runnable, Comparable<PriorityTask> {
private final int priority;
public PriorityTask(int priority) {
this.priority = priority;
}
@Override
public void run() {
...
}
@Override
public int compareTo(PriorityTask other) {
return Integer.compare(other.priority, this.priority);
}
}

DelayQueue<DelayedTask> – выдаёт задачу только после истечения задержки. Но с обычным ThreadPoolExecutor преимуществ не будет, но он используется в другой реализации – ScheduledThreadPoolExecutor, которая поддерживает параметр задержки.

ScheduledExecutorService executor = Executors.newScheduledThreadPool(2);
executor.schedule(
() -> System.out.println("Через 5 секунд"),
5,
TimeUnit.SECONDS
);


И периодическая задача c помощью метода scheduleAtFixedRate().

executor.scheduleAtFixedRate(
this::checkNotifications,
0,
10,
TimeUnit.MINUTES
);

Распареллеливание вычислений с ForkJoinPool

Но вернёмся к реализациям ExecutorService. Существует также другой наследник – ForkJoinPool . Он используется для рекурсивных задач и распараллеливания вычислений – с ним можно одну задачу разделить на более мелкие независимые подзадачи, выполнить их параллельно, а затем объединить результаты.

Название отражает принцип его работы:
Fork – большая задача рекурсивно разбивается на несколько маленьких.
Join – после выполнения всех подзадач их результаты объединяются в итоговый результат.

Такой подход хорош для сортировки больших массивов, обработки изображений, рекурсивного обхода деревьев и т. д. Например, он применяется в Arrays.parallelSort().

Главное отличие ForkJoinPool от ThreadPoolExecutor в организации очередей задач. В обычном ThreadPoolExecutor все рабочие потоки получают задачи из одной общей очереди. Если один поток освободился, он просто берёт следующую задачу из этой очереди.

В ForkJoinPool у каждого рабочего потока есть собственная двусторонняя очередь задач (deque). Пока поток создаёт новые подзадачи, он помещает их в свою очередь и сам их обрабатывает. Если другой поток закончил работу раньше и простаивает, он может украсть часть задач из очереди занятого потока. Это называется Work Stealing (кража работы). Благодаря этому нагрузка перераспределяется между потоками, и никто из них не простаивает подолгу.

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

class SumTask extends RecursiveTask<Long> {
private final int[] array;
private final int from;
private final int to;
@Override
protected Long compute() {
if (to - from < 1000) {
long sum = 0;
for (int i = from; i < to; i++) {
sum += array[i];
}
return sum;
}
int middle = (from + to) / 2;
SumTask left = new SumTask(array, from, middle);
SumTask right = new SumTask(array, middle, to);
left.fork();
long rightResult = right.compute();
long leftResult = left.join();
return leftResult + rightResult;
}
}

Здесь мы просто суммируем элементы большого массива. Если их немного, то делаем это сразу, а иначе разделяем задачу на две. Одна поleft.fork() помещается в очередь, а другая половина (right.compute()) продолжает выполняться в текущем потоке. И так далее, рекурсивно.

Ещё существуют реализации:

SingleThreadExecutor – один поток, задачи выполняются строго по очереди.

CachedThreadPool – при высокой нагрузке создаёт новые потоки, а простаивающие позже удаляет.

Жизненный цикл ExecutorService

В отличие от обычного потока, который завершается после выхода из метода run(), потоки в ExecutorService продолжают существовать, даже если в очереди больше нет задач. Они остаются в ожидании следующей работы. Если ничего не предпринять, такие потоки могут продолжать существовать бесконечно, удерживая ресурсы и, в случае обычных (non-daemon) потоков, не позволяя JVM завершить работу. Именно поэтому после завершения использования пула его необходимо явно закрывать.

Для завершения часто используют метод shutdown(). В этом случае новые задачи уже нельзя отправить на выполнение, а вот старые – даже те, которые находятся в очереди, будут выполнены. То есть это мягкое завершение, которое не гарантирует быстроту, но более безопасно.

Тогда жизненный цикл выходит таким:

RUNNING ---- shutdown() ----> SHUTDOWN ----> TERMINATED

Если нужно в текущем потоке подождать завершения потоков из пула, то можно использовать awaitTermination().

ExecutorService executor = Executors.newFixedThreadPool(2);
executor.submit(task1);
executor.submit(task2);
executor.submit(task3);
executor.shutdown();
System.out.println("Main finished");

В этом случае вывод из главного потока произойдёт до выполнения всех заданий. Но если надо подождать их выполнения (сделать некоторую сихнронность), то используем awaitTermination():

executor.shutdown();
executor.awaitTermination(1, TimeUnit.MINUTES);
System.out.println("Main finished");


Этот метод возвращает boolean: true означает, что пул успел завершиться до истечения таймаута.

executor.shutdown();
boolean terminated = executor.awaitTermination(1, TimeUnit.MINUTES);
if (!terminated) {
executor.shutdownNow();
}

Для более жёсткого завершения используют shutdownNow(). Этот метод не только запрещает новые задачи, но и удаляет задачи из очереди, посылает interrupt всем выполняющимся задачам и возвращает список задач, которые так и не были запущены.

Но и shutdownNow() не убивает поток мгновенно – он только только отправляет interrupt, а что там дальше в потоке произойдёт – не его забота. То есть внутри выполняемой задачи должна быть какая-то реакция на прерывание, иначе поток продолжит выполнять свои дела:

try {
Thread.sleep(10000);
} catch (InterruptedException e) {
System.out.println("Меня прервали");
return;
}


Ну или должен проверять флаг прерывания:

while (!Thread.currentThread().isInterrupted()) {
doWork();
}
if (Thread.interrupted()) {
return;
}

Ну и вернёмся к аналогии с почтой, которую я предложил выше. Милая работница вызовет методshutdown() и будет разворачивать всех тех, ко приходит прямо к закрытию, но тех, кто уже стоял в очереди, обслужит.

А вот другая работница вызоветshutdownNow() и просто уйдёт в положенное время, возможно даже не доделав текущую операцию. Придётся её потворять на следующий день.

Или отрегаирует на такое «прерывание», всё же обслужив последнего клиента, потом ещё приберётся на рабочем месте и только после этого уйдёт домой – но клиентам в очереди всё равно шанса не оставит.

Цепочки действий с CompletableFuture

Обычный Future умеет довольно мало. Он чаще всего просто говорит, выполнена ли задача и позволяет получить результат. Но если нужно делать какие-то ещё операции, то задача снова устроить асинхронное выполнение ложится на разработчика.

CompletableFuture позволяет описать цепочку действий, которая автоматически продолжится после завершения предыдущей задачи.

CompletableFuture
.supplyAsync(this::loadUser)
.thenApply(this::calculateStatistics)
.thenAccept(this::sendEmail);

Основные методы в цепочке:

supplyAsync() – запускает задачу, которая возвращает результат. Аналог executor.submit(callable).
thenApply() – преобразует результат. Похоже на map() у Stream.
thenAccept() – конец цепочки, ничего не возвращает, но вызывает метод.
exceptionally() – обрабатывает ошибку.

Под капотом у CompletableFuture ForkJoinPool. И поэтому, чтобы главный поток ждал выполнения всех действий, нужно вызывать метод join():

System.out.println("1");
CompletableFuture<String> future =
CompletableFuture.supplyAsync(() -> {
sleep(3000);
return "Hello";
});
System.out.println("2");
String result = future.join();
System.out.println(result);

Этот пример сначала выведет 1, потом 2, а потом уже поздоровается Hello из CompletableFuture.

Но можно указать свой вариант:

ExecutorService executor = Executors.newFixedThreadPool(4);
CompletableFuture<User> future =
CompletableFuture.supplyAsync(
this::loadUser,
executor
);

Потоки в Spring

В Spring потоки встречаются постоянно, даже если разработчик явно не создаёт Thread. Например:

@Async
public void sendEmail() {
}

Метод с @Async выполняется через Spring-managed Executor. По умолчанию или по собственной конфигурации Spring передаёт задачу в пул потоков. Можно настроить пул:

@Bean
public Executor taskExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(4);
executor.setMaxPoolSize(8);
executor.setQueueCapacity(100);
executor.setThreadNamePrefix("async-");
executor.initialize();
return executor;
}

Здесь ThreadPoolTaskExecutor – обёртка над Java-пулом потоков. Внутри он обычно использует ThreadPoolExecutor.

Также потоки используются в @Scheduled, обработке HTTP-запросов, Spring Batch, обработчиках событий, Kafka listeners. И как раз про Kafka будет следующая статья.

Виртугальные потоки в Java

И перед тем, как перейти к Kafka, нужно обсудить относительную новинку в языке – виртуальные потоки. До Java 21 потоки создавались только операционной системой (это были Platform Threads), но это довольно дорого: нужен собственный стек (1-2 МБ), а ОС должна хранить информацию о потоке. Ещё было ограничение на количество, поэтому создать тысячи потоков, как мы можем сделать при выполнении какой-нибудь рекурсивной задачи, было не всегда возможно.

Поэтому теперь существуют Virtual Threads, которые создаются JVM и не требуют столько ресурсов. Их преимущество так же в том, что JVM снимает виртуальный поток с платформенного (они называются Carrier Thread), если тот долго простаивает. Раньше бы поток со sleep() просто зря занимал поток ОС, а теперь на него можно назначить другой виртуальный поток. Этот процесс называется монтированием (mount) и размонтированием (unmount).

Вычислительная способность процессора не растёт, но может уменьшиться время его простаивания: JVM будет занимать его разными виртуальными потоками. Они особенно полезны, когда поток занимается базами данными и ждёт ответ, читает файлы, ожидает сетевой запрос или сообщение от другого сервиса. В общем, виртуальные потоки хорошо подходят для I/O-задач.

У виртуальных потоков существует одно ограничение, известное как pinning. Если виртуальный поток выполняет блокирующую операцию внутри блока synchronized, JVM не сможет выполнить unmount. В результате Carrier Thread окажется «прикреплён» к этому виртуальному потоку и будет простаивать вместе с ним.

Это снижает масштабируемость приложения, поэтому в современном коде вместо synchronized всё чаще используют более гибкие средства синхронизации из пакета java.util.concurrent, например ReentrantLock. По крайней мере, на длительных блокирующих операциях.

public void download(Socket socket) throws IOException {
synchronized (lock) {
socket.getInputStream().read(); // ждём данные
}
}

Внутри JDK виртуальный поток реализован отдельным классом VirtualThread, но он не публичный, поэтому просто написать new VirtualThread() нельзя. Вместо этого используются специальные фабричные методы:

Thread.startVirtualThread(() -> doWork());
Thread.ofVirtual().start(() -> doWork());

Если приложение уже использует ExecutorService, переход на виртуальные потоки может потребовать минимальных изменений:

try (ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor()) {
executor.submit(this::doWork);
}

Несмотря на слово Executor, здесь нет привычного пула потоков. Каждая отправленная задача получает собственный виртуальный поток.

Но теперь, когда виртуальные потоки можно быстро создавать и завершать, хранить и переиспользовать их уже нет смысла. Получается, что стандартный пул больше и не нужен, но Executor всё ещё остаётся актуальным. Зачем тогда вообще использовать ExecutorService, если можно просто вызвать Thread.startVirtualThread()?

ExecutorService – прежде всего абстракция для отправки задач на выполнение, управления их жизненным циклом и получения результатов через Future. Конкретная реализация уже сама решает, как выполнять эти задачи.

Классический ThreadPoolExecutor хранит ограниченное количество дорогих платформенных потоков и многократно использует каждый из них для выполнения разных объектов Runnable и Callable. А Executors.newVirtualThreadPerTaskExecutor() действует иначе: для каждой отправленной задачи он создаёт отдельный виртуальный поток и завершает его после окончания работы.

ThreadPoolExecutor:
task1 ─┐
task2 ─┼─> один и тот же Platform Thread
task3 ─┘
VirtualThreadPerTaskExecutor:
task1 ─> Virtual Thread 1 ─┐
task2 ─> Virtual Thread 2 ─┼─> Carrier Threads
task3 ─> Virtual Thread 3 ─┘

И там и там происходит удачная попытка экономить ресурсы и переиспользовать существующие платформенные потоки. В одном случае они используются для разных задач, в другом – для выпонления виртуальных потоков, которые уже получают по задаче. Просто теперь между задачей и платформенным потоком появился ещё один уровень абстракции – виртуальный поток.

Оставить комментарий