Я уже почти год назад законспектировал основные способы сделать код на 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. Например:
@Asyncpublic void sendEmail() {}
Метод с @Async выполняется через Spring-managed Executor. По умолчанию или по собственной конфигурации Spring передаёт задачу в пул потоков. Можно настроить пул:
@Beanpublic 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 Threadtask3 ─┘
VirtualThreadPerTaskExecutor:task1 ─> Virtual Thread 1 ─┐task2 ─> Virtual Thread 2 ─┼─> Carrier Threadstask3 ─> Virtual Thread 3 ─┘
И там и там происходит удачная попытка экономить ресурсы и переиспользовать существующие платформенные потоки. В одном случае они используются для разных задач, в другом – для выпонления виртуальных потоков, которые уже получают по задаче. Просто теперь между задачей и платформенным потоком появился ещё один уровень абстракции – виртуальный поток.