openjdk.ruOpenJDK на русском

JEP 480: Structured Concurrency (Third Preview)

Structured Concurrency (структурированная конкурентность), третья версия Preview (предварительная версия)

АвторRon Pressler & Alan Bateman
ОтветственныйAlan Bateman
ТипFeature
ОбластьSE
СтатусClosed / Delivered
Выпуск23
Компонентcore-libs
Обсуждениеloom dash dev at openjdk dot org
Связан сJEP 462: Structured Concurrency (Second Preview)
JEP 499: Structured Concurrency (Fourth Preview)
РецензентыPaul Sandoz
ОдобренPaul Sandoz
Создан2024/04/22 12:11
Обновлён2025/02/25 16:34
Задача8330818

Аннотация

Упростить конкурентное программирование, добавив API для Structured Concurrency. Structured Concurrency рассматривает группы связанных задач, выполняющихся в разных потоках, как единую единицу работы, тем самым упрощая обработку ошибок и отмену, повышая надёжность и улучшая наблюдаемость. Это API в статусе Preview.

История

Structured Concurrency была предложена в JEP 428 и поставлена в JDK 19 как API в статусе Incubator (инкубационный модуль). Повторно в статусе Incubator API был представлен в JEP 437 в JDK 20 с небольшим изменением: добавлено наследование Scoped Values (значения с ограниченной областью видимости; JEP 429). Впервые в статусе Preview API появился в JDK 21 в JEP 453: метод StructuredTaskScope::fork(...) стал возвращать Subtask вместо Future. Повторно в статусе Preview он вышел в JDK 22 в JEP 462 без изменений. Здесь мы предлагаем ещё раз выпустить API в статусе Preview в JDK 23 без изменений, чтобы получить больше отзывов.

Цели

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

  • Улучшить наблюдаемость конкурентного кода.

Что не является целью

  • Целью не является замена каких-либо конструкций конкурентности из пакета java.util.concurrent, таких как ExecutorService и Future.

  • Целью не является определение окончательного API Structured Concurrency для платформы Java. Другие конструкции Structured Concurrency могут быть определены сторонними библиотеками или в будущих выпусках JDK.

  • Цель не в том, чтобы определить способ обмена потоками данных между потоками (т. е. каналы). Возможно, мы предложим это в будущем.

  • Цель не в том, чтобы заменить существующий механизм прерывания потоков новым механизмом отмены потоков. Возможно, мы предложим это в будущем.

Мотивация

Разработчики справляются со сложностью, разбивая задачи на несколько подзадач. В обычном однопоточном коде подзадачи выполняются последовательно. Однако если подзадачи достаточно независимы друг от друга и аппаратных ресурсов достаточно, то общую задачу можно выполнить быстрее (т. е. с меньшей задержкой), выполняя подзадачи конкурентно. Например, задача, объединяющая результаты нескольких операций ввода-вывода, выполнится быстрее, если каждая операция ввода-вывода выполняется конкурентно в собственном потоке. Virtual Threads (виртуальные потоки; JEP 444) делают экономически оправданным выделение отдельного потока под каждую такую операцию ввода-вывода, но управление огромным числом потоков, которое может при этом возникнуть, остаётся сложной задачей.

Неструктурированная конкурентность с ExecutorService

API java.util.concurrent.ExecutorService, появившийся в Java 5, помогает разработчикам выполнять подзадачи конкурентно.

Например, вот метод handle(), представляющий задачу в серверном приложении. Он обрабатывает входящий запрос, отправляя две подзадачи в ExecutorService. Одна подзадача выполняет метод findUser(), а другая — метод fetchOrder(). ExecutorService сразу возвращает Future для каждой подзадачи и выполняет подзадачи конкурентно в соответствии с политикой планирования Executor. Метод handle() ожидает результатов подзадач с помощью блокирующих вызовов методов get() их future-объектов, поэтому говорят, что задача присоединяет свои подзадачи.

Response handle() throws ExecutionException, InterruptedException {
    Future<String>  user  = esvc.submit(() -> findUser());
    Future<Integer> order = esvc.submit(() -> fetchOrder());
    String theUser  = user.get();   // Join findUser
    int    theOrder = order.get();  // Join fetchOrder
    return new Response(theUser, theOrder);
}

Поскольку подзадачи выполняются конкурентно, каждая из них может завершиться успешно или сбоем независимо от других. (Сбой в этом контексте означает выброс исключения.) Часто задача, такая как handle(), должна завершиться сбоем, если сбоем завершилась любая из её подзадач. Разобраться во времени жизни потоков при сбое может оказаться на удивление сложно:

  • Если findUser() выбрасывает исключение, то handle() выбросит исключение при вызове user.get(), но fetchOrder() продолжит выполняться в своём потоке. Это утечка потока, которая в лучшем случае расходует ресурсы впустую, а в худшем — поток fetchOrder() будет мешать другим задачам.

  • Если поток, выполняющий handle(), прерывается, прерывание не распространяется на подзадачи. Оба потока, findUser() и fetchOrder(), утекут и продолжат выполняться даже после сбоя handle().

  • Если findUser() выполняется долго, а fetchOrder() тем временем завершается сбоем, то handle() будет без необходимости ждать findUser(), блокируясь на user.get(), вместо того чтобы отменить эту подзадачу. Только после того как findUser() завершится и user.get() вернёт управление, order.get() выбросит исключение, из-за чего handle() завершится сбоем.

В каждом случае проблема в том, что наша программа логически структурирована отношениями «задача — подзадача», но эти отношения существуют только в голове разработчика.

Из-за этого не только больше возможностей для ошибок, но и сложнее диагностировать и устранять такие ошибки. Например, средства наблюдения, такие как дампы потоков, покажут handle(), findUser() и fetchOrder() в стеках вызовов не связанных между собой потоков, без какого-либо указания на отношение «задача — подзадача».

Можно попытаться сделать лучше, явно отменяя другие подзадачи при возникновении ошибки, например оборачивая задачи в try-finally и вызывая методы cancel(boolean) future-объектов других задач в блоке catch для задачи, завершившейся с ошибкой. Также понадобилось бы использовать ExecutorService внутри оператора try-with-resources, как показано в примерах в JEP 425, потому что Future не даёт способа дождаться отменённой задачи. Но всё это бывает очень сложно сделать правильно, и из-за этого часто труднее понять логический замысел кода. Отслеживать отношения между задачами и вручную добавлять обратно нужные связи отмены между задачами — это слишком многого требовать от разработчиков.

Необходимость вручную координировать время жизни вызвана тем, что ExecutorService и Future допускают неограниченные шаблоны конкурентности. На участвующие потоки не накладывается никаких ограничений и никакого порядка. Один поток может создать ExecutorService, второй — отправить в него работу, а потоки, выполняющие эту работу, никак не связаны ни с первым, ни со вторым потоком. Более того, после того как поток отправил работу, результатов выполнения может ожидать совсем другой поток. Любой код, у которого есть ссылка на Future, может присоединить его (т. е. дождаться его результата, вызвав get()), даже код в потоке, отличном от того, который получил Future. По сути, подзадача, запущенная одной задачей, не обязана возвращаться к задаче, которая её отправила. Она может вернуться к любой из множества задач — или ни к одной.

Поскольку ExecutorService и Future допускают такое неструктурированное использование, они не обеспечивают и даже не отслеживают отношения между задачами и подзадачами, хотя такие отношения распространены и полезны. Поэтому даже когда подзадачи отправляются и присоединяются в одной и той же задаче, сбой одной подзадачи не может автоматически вызвать отмену другой: в методе handle() выше сбой fetchOrder() не может автоматически вызвать отмену findUser(). Future-объект для fetchOrder() никак не связан с future-объектом для findUser(), и ни один из них не связан с потоком, который в итоге присоединит его через его метод get(). Вместо того чтобы требовать от разработчиков управлять такой отменой вручную, мы хотим надёжно её автоматизировать.

Структура задач должна отражать структуру кода

В отличие от вольного набора потоков при использовании ExecutorService, выполнение однопоточного кода всегда обеспечивает иерархию задач и подзадач. Блок тела {...} метода соответствует задаче, а методы, вызываемые внутри блока, — подзадачам. Вызванный метод должен либо вернуть управление вызвавшему его методу, либо выбросить в него исключение. Он не может пережить вызвавший его метод и не может вернуть управление или выбросить исключение в другой метод. Таким образом, все подзадачи завершаются раньше задачи, каждая подзадача является дочерней по отношению к своей родительской задаче, а время жизни каждой подзадачи относительно других подзадач и задачи определяется синтаксической блочной структурой кода.

Например, в этой однопоточной версии handle() отношение «задача — подзадача» видно из синтаксической структуры:

Response handle() throws IOException {
    String theUser  = findUser();
    int    theOrder = fetchOrder();
    return new Response(theUser, theOrder);
}

Мы не запускаем подзадачу fetchOrder(), пока не завершится подзадача findUser(), успешно или сбоем. Если findUser() завершается сбоем, мы вообще не запускаем fetchOrder(), и задача handle() неявно завершается сбоем. Важно то, что подзадача может вернуться только к своей родительской задаче: это означает, что родительская задача может неявно считать сбой одной подзадачи сигналом отменить другие незавершённые подзадачи, а затем сама завершиться сбоем.

В однопоточном коде иерархия «задача — подзадача» материализуется во время выполнения в стеке вызовов. Поэтому соответствующие отношения «родитель — потомок», которые управляют распространением ошибок, мы получаем бесплатно. При наблюдении за одним потоком иерархическое отношение очевидно: findUser() (а затем fetchOrder()) выглядят подчинёнными handle(). Поэтому легко ответить на вопрос: «Над чем handle() работает сейчас?»

Конкурентное программирование было бы проще, надёжнее и удобнее для наблюдения, если бы отношения «родитель — потомок» между задачами и их подзадачами были видны из синтаксической структуры кода, а также материализовались во время выполнения — так же, как в однопоточном коде. Синтаксическая структура очерчивала бы время жизни подзадач и позволяла бы представить во время выполнения иерархию между потоками, аналогичную стеку вызовов внутри потока. Такое представление обеспечило бы распространение ошибок и отмену, а также осмысленное наблюдение за конкурентной программой.

(В платформе Java уже есть API для придания структуры конкурентным задачам, а именно java.util.concurrent.ForkJoinPool — механизм выполнения, лежащий в основе параллельных потоков данных (parallel streams). Однако этот API рассчитан на вычислительно ёмкие задачи, а не на задачи, связанные с вводом-выводом.)

Structured Concurrency

Structured Concurrency — это подход к конкурентному программированию, сохраняющий естественную связь между задачами и подзадачами, благодаря чему конкурентный код становится более читаемым, сопровождаемым и надёжным. Термин «Structured Concurrency» ввёл Martin Sústrik, а популяризировал Nathaniel J. Smith. На проектирование обработки ошибок в Structured Concurrency повлияли идеи из других языков, например иерархические супервизоры Erlang.

Structured Concurrency основывается на простом принципе:

Если задача разделяется на конкурентные подзадачи, то все они возвращаются в одно и то же место, а именно в блок кода задачи.

В Structured Concurrency подзадачи работают от имени задачи. Задача ожидает результатов подзадач и следит за их сбоями. Как и в случае приёмов структурного программирования для кода в одном потоке, сила Structured Concurrency для нескольких потоков основана на двух идеях: (1) чётко определённые точки входа и выхода для потока выполнения через блок кода и (2) строгая вложенность времени жизни операций, отражающая их синтаксическую вложенность в коде.

Поскольку точки входа и выхода блока кода чётко определены, время жизни конкурентной подзадачи ограничено синтаксическим блоком её родительской задачи. Поскольку время жизни подзадач одного уровня вложено во время жизни их родительской задачи, их можно рассматривать и ими можно управлять как единым целым. Поскольку время жизни родительской задачи, в свою очередь, вложено во время жизни её родителя, среда выполнения может материализовать иерархию задач в виде дерева — конкурентного аналога стека вызовов одного потока. Так код может применять политики, например крайние сроки, ко всему поддереву задач, а инструменты наблюдения могут показывать подзадачи как подчинённые их родительским задачам.

Structured Concurrency отлично сочетается с Virtual Threads — легковесными потоками, реализованными в JDK. Многие Virtual Threads разделяют один и тот же поток операционной системы, что позволяет иметь очень большое число потоков Virtual Threads. Помимо того что их может быть много, Virtual Threads достаточно дёшевы, чтобы представлять любую конкурентную единицу поведения, даже поведение, связанное с вводом-выводом. Это означает, что серверное приложение может использовать Structured Concurrency для одновременной обработки тысяч или миллионов входящих запросов: оно может выделить новый поток Virtual Threads под задачу обработки каждого запроса, а когда задача разветвляется, отправляя подзадачи на конкурентное выполнение, — выделить новый поток Virtual Threads под каждую подзадачу. За кулисами отношение «задача — подзадача» материализуется в дерево за счёт того, что каждый поток Virtual Threads хранит ссылку на своего единственного родителя, подобно тому как кадр в стеке вызовов ссылается на единственный вызвавший его кадр.

Итак, Virtual Threads дают изобилие потоков. Structured Concurrency может корректно и надёжно их координировать и позволяет инструментам наблюдения отображать потоки так, как их понимает разработчик. Наличие в JDK API для Structured Concurrency упростит создание сопровождаемых, надёжных и наблюдаемых серверных приложений.

Описание

Основной класс API Structured Concurrency — StructuredTaskScope в пакете java.util.concurrent. Этот класс позволяет разработчикам структурировать задачу как семейство конкурентных подзадач и координировать их как единое целое. Подзадачи выполняются в собственных потоках: их по отдельности порождают, а затем присоединяют как единое целое и, возможно, отменяют как единое целое. Успешные результаты или исключения подзадач собираются и обрабатываются родительской задачей. StructuredTaskScope ограничивает время жизни подзадач чёткой лексической областью видимости, в которой происходит всё взаимодействие задачи с её подзадачами — порождение, присоединение, отмена, обработка ошибок и объединение результатов.

Вот приведённый ранее пример handle(), переписанный с использованием StructuredTaskScope (ShutdownOnFailure объясняется ниже):

Response handle() throws ExecutionException, InterruptedException {
    try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
        Supplier<String>  user  = scope.fork(() -> findUser());
        Supplier<Integer> order = scope.fork(() -> fetchOrder());

        scope.join()            // Join both subtasks
             .throwIfFailed();  // ... and propagate errors

        // Here, both subtasks have succeeded, so compose their results
        return new Response(user.get(), order.get());
    }
}

В отличие от исходного примера, здесь легко понять время жизни участвующих потоков: при любых условиях оно ограничено лексической областью видимости, а именно телом оператора try-with-resources. Кроме того, использование StructuredTaskScope обеспечивает ряд ценных свойств:

  • Обработка ошибок с коротким замыканием — если любая из подзадач findUser() или fetchOrder() завершается сбоем, другая отменяется, если она ещё не завершилась. (Этим управляет политика завершения, реализованная в ShutdownOnFailure; возможны и другие политики.)

  • Распространение отмены — если поток, выполняющий handle(), прерывается до или во время вызова join(), обе подзадачи автоматически отменяются, когда поток выходит из области.

  • Ясность — у приведённого кода чёткая структура: подготовить подзадачи, дождаться, пока они завершатся или будут отменены, а затем решить, завершиться ли успешно (и обработать результаты дочерних задач, которые уже закончены) или со сбоем (подзадачи уже закончены, так что очищать больше нечего).

  • Наблюдаемость — дамп потоков, описанный ниже, наглядно показывает иерархию задач: потоки, выполняющие findUser() и fetchOrder(), отображаются как дочерние элементы области видимости.

StructuredTaskScope — это API в статусе Preview, по умолчанию отключённый

Чтобы использовать API StructuredTaskScope, нужно включить Preview-API следующим образом:

  • Скомпилируйте программу с javac --release 23 --enable-preview Main.java и запускайте её с java --enable-preview Main; или

  • При использовании средства запуска исходного кода запускайте программу с java --source 23 --enable-preview Main.java; или

  • При использовании jshell запускайте его с jshell --enable-preview.

Использование StructuredTaskScope

API StructuredTaskScope выглядит так:

public class StructuredTaskScope<T> implements AutoCloseable {

    public <U extends T> Subtask<U> fork(Callable<? extends U> task);
    public void shutdown();

    public StructuredTaskScope<T> join() throws InterruptedException;
    public StructuredTaskScope<T> joinUntil(Instant deadline)
        throws InterruptedException, TimeoutException;
    public void close();

    protected void handleComplete(Subtask<? extends T> handle);
    protected final void ensureOwnerAndJoined();

}

Общий порядок работы кода с StructuredTaskScope такой:

  1. Создайте область. Поток, создающий область, является её владельцем.

  2. Используйте метод fork(Callable), чтобы порождать подзадачи в области.

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

  4. Владелец области присоединяется (join) к области, т. е. ко всем её подзадачам, как к единому целому. Владелец может вызвать метод join() области, чтобы дождаться, пока все подзадачи либо завершатся (успешно или нет), либо будут отменены через shutdown(). Он также может вызвать метод joinUntil(java.time.Instant) области, чтобы ждать не дольше заданного крайнего срока.

  5. После присоединения обработайте все ошибки в подзадачах и обработайте их результаты.

  6. Закройте область, обычно неявно, через try-with-resources. При этом работа области останавливается, если она ещё не остановлена, и выполняется ожидание завершения всех подзадач, которые были отменены, но ещё не завершились.

Каждый вызов fork(...) запускает новый поток для выполнения подзадачи; по умолчанию это поток Virtual Threads. Подзадача может создать собственную вложенную StructuredTaskScope, чтобы порождать свои подзадачи, и так образуется иерархия. Эта иерархия отражается в блочной структуре кода, которая ограничивает время жизни подзадач: гарантируется, что после закрытия области все потоки подзадач завершены, и при выходе из блока не остаётся ни одного потока.

Любая подзадача в области, любые подподзадачи во вложенной области и владелец области могут в любой момент вызвать метод shutdown() области, чтобы обозначить, что задача выполнена, — даже когда другие подзадачи ещё выполняются. Метод shutdown() прерывает потоки, которые ещё выполняют подзадачи, и заставляет метод join() или joinUntil(Instant) вернуть управление. Поэтому все подзадачи следует писать так, чтобы они реагировали на прерывание. Новая подзадача, порождённая после вызова shutdown(), будет находиться в состоянии UNAVAILABLE и не будет запущена. По сути, shutdown() — это аналог оператора break из последовательного кода для конкурентного кода.

Вызов join() или joinUntil(Instant) внутри области обязателен. Если выход из блока области происходит до присоединения, область дождётся завершения всех подзадач, а затем выбросит исключение.

Поток — владелец области может быть прерван как до присоединения, так и во время него. Например, он может быть подзадачей объемлющей области, работа которой была остановлена. В этом случае join() и joinUntil(Instant) выбросят исключение, поскольку продолжать нет смысла. Затем оператор try-with-resources остановит работу области, что отменит все подзадачи и дождётся их завершения. В результате отмена задачи автоматически распространяется на её подзадачи. Если крайний срок метода joinUntil(Instant) истекает до того, как подзадачи завершатся или будет вызван shutdown(), метод выбросит исключение, и оператор try-with-resources снова остановит работу области.

Когда join() завершается успешно, каждая из подзадач либо завершилась успешно, либо завершилась сбоем, либо была отменена, поскольку работа области была остановлена.

После присоединения владелец области обрабатывает подзадачи, завершившиеся сбоем, и обрабатывает результаты успешно завершившихся подзадач; обычно это делает политика остановки (см. ниже). Результат успешно завершившейся задачи можно получить методом Subtask.get(). Метод get() никогда не блокируется; он выбрасывает IllegalStateException, если его по ошибке вызвали до присоединения или если подзадача не завершилась успешно.

Подзадачи, порождённые в области, наследуют привязки ScopedValue (JEP 446). Если владелец области читает значение из привязанного ScopedValue, каждая подзадача прочитает то же самое значение.

Если владелец области сам является подзадачей существующей области, т. е. был создан как порождённая подзадача, то эта область становится родительской для новой области. Таким образом, области и подзадачи образуют дерево.

Структурированное использование StructuredTaskScope обеспечивается во время выполнения. Например, попытки вызвать fork(Callable) из потока, который не входит в древовидную иерархию области — т. е. не является владельцем, подзадачей или подзадачей во вложенной области (подподзадачей), — завершатся исключением. Использование области вне блока try-with-resources и возврат без вызова close() или без соблюдения правильной вложенности вызовов close() может привести к тому, что методы области выбросят StructureViolationException.

StructuredTaskScope навязывает конкурентным операциям структуру и порядок. Поэтому он не реализует интерфейсы ExecutorService и Executor, так как экземпляры этих интерфейсов обычно используются неструктурированно (см. ниже). Однако код, который использует ExecutorService, но выиграл бы от структуры, несложно перевести на StructuredTaskScope.

На практике в большинстве случаев StructuredTaskScope будет использоваться не через класс StructuredTaskScope напрямую, а через один из двух подклассов, описанных в следующем разделе, которые реализуют политики остановки. В других сценариях пользователи, скорее всего, будут писать собственные подклассы для реализации своих политик остановки.

Политики остановки

При работе с конкурентными подзадачами часто используют паттерны с досрочным завершением, чтобы не выполнять ненужную работу. Например, иногда имеет смысл отменить все подзадачи, если одна из них завершилась сбоем (т. е. invoke all), или, наоборот, если одна из них завершилась успешно (т. е. invoke any). Два подкласса StructuredTaskScope, ShutdownOnFailure и ShutdownOnSuccess, поддерживают эти паттерны политиками, которые останавливают работу области, когда первая подзадача завершается сбоем или успешно соответственно.

Кроме того, политики остановки предоставляют централизованные методы для обработки исключений и, возможно, успешных результатов. Это соответствует духу Structured Concurrency, согласно которому вся область рассматривается как единое целое.

Вот StructuredTaskScope с политикой остановки при сбое (она также использовалась в примере handle() выше), который выполняет набор задач конкурентно и завершается сбоем, если сбоем завершается любая из них:

<T> List<T> runAll(List<Callable<T>> tasks) 
        throws InterruptedException, ExecutionException {
    try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
        List<? extends Supplier<T>> suppliers = tasks.stream().map(scope::fork).toList();
        scope.join()
             .throwIfFailed();  // Propagate exception if any subtask fails
        // Here, all tasks have succeeded, so compose their results
        return suppliers.stream().map(Supplier::get).toList();
    }
}

Вот StructuredTaskScope с политикой остановки при успехе, который возвращает результат первой успешно завершившейся подзадачи:

<T> T race(List<Callable<T>> tasks, Instant deadline) 
        throws InterruptedException, ExecutionException, TimeoutException {
    try (var scope = new StructuredTaskScope.ShutdownOnSuccess<T>()) {
        for (var task : tasks) {
            scope.fork(task);
        }
        return scope.joinUntil(deadline)
                    .result();  // Throws if none of the subtasks completed successfully
    }
}

Как только одна подзадача завершается успешно, работа этой области автоматически останавливается, а незавершённые подзадачи отменяются. Задача завершается сбоем, если сбоем завершаются все подзадачи или если истекает заданный крайний срок. Этот паттерн может быть полезен, например, в серверных приложениях, которым нужен результат от любого из набора дублирующих друг друга сервисов.

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

Обработка результатов

После присоединения и централизованной обработки исключений политикой остановки (например, с помощью ShutdownOnFailure::throwIfFailed) владелец области может обработать результаты подзадач, используя объекты Subtask, возвращённые вызовами fork(...), если эти результаты не обрабатывает политика (например, ShutdownOnSuccess::result()).

Как правило, единственный метод Subtask, который вызывает владелец области, — это метод get(). Все остальные методы Subtask обычно используются только в реализации метода handleComplete(...) пользовательских политик остановки (см. ниже). Более того, мы рекомендуем объявлять переменные, ссылающиеся на Subtask, возвращённый fork(...), с типом, например, Supplier<String>, а не Subtask<String> (если, конечно, вы не решили использовать var). Если политика остановки сама обрабатывает результаты подзадач — как в случае ShutdownOnSuccess, — то объектов Subtask, возвращаемых fork(...), следует вовсе избегать, а метод fork(...) рассматривать так, будто он возвращает void. Подзадачи должны возвращать в качестве результата всю информацию, которую владелец области должен обработать после централизованной обработки исключений политикой.

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

<T> List<Future<T>> executeAll(List<Callable<T>> tasks)
        throws InterruptedException {
    try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
    	  List<? extends Supplier<Future<T>>> futures = tasks.stream()
    	      .map(task -> asFuture(task))
     	      .map(scope::fork)
     	      .toList();
    	  scope.join();
    	  return futures.stream().map(Supplier::get).toList();
    }
}

static <T> Callable<Future<T>> asFuture(Callable<T> task) {
   return () -> {
       try {
           return CompletableFuture.completedFuture(task.call());
       } catch (Exception ex) {
           return CompletableFuture.failedFuture(ex);
       }
   };
}

Пользовательские политики остановки

StructuredTaskScope можно расширить и переопределить его защищённый метод handleComplete(...), чтобы реализовать политики, отличные от политик ShutdownOnSuccess и ShutdownOnFailure. Подкласс может, например,

  • собирать результаты успешно завершившихся подзадач и игнорировать подзадачи, завершившиеся сбоем,
  • собирать исключения, когда подзадачи завершаются сбоем, или
  • вызывать метод shutdown(), чтобы остановить работу и заставить join() проснуться при наступлении некоторого условия.

Когда подзадача завершается, даже после вызова shutdown(), о ней сообщается методу handleComplete(...) в виде Subtask:

public sealed interface Subtask<T> extends Supplier<T> {
    enum State { SUCCESS, FAILED, UNAVAILABLE }

    State state();
    Callable<? extends T> task();
    T get();
    Throwable exception();
}

Метод handleComplete(...) вызывается для подзадач, которые завершились либо успешно (состояние SUCCESS), либо неуспешно (состояние FAILED) до вызова shutdown(). Метод get() можно вызывать, только если подзадача находится в состоянии SUCCESS, а метод exception() — только если подзадача находится в состоянии FAILED; вызов get() или exception() в других ситуациях приведёт к тому, что они выбросят IllegalStateException. Состояние UNAVAILABLE означает одно из следующего: (1) подзадача была порождена, но ещё не завершилась; (2) подзадача завершилась после остановки или (3) подзадача была порождена после остановки и поэтому не запускалась. Метод handleComplete(...) никогда не вызывается для подзадачи в состоянии UNAVAILABLE.

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

Вот пример подкласса StructuredTaskScope, который собирает результаты успешно завершившихся подзадач. Он определяет метод results(), которым основная задача получает результаты.

class MyScope<T> extends StructuredTaskScope<T> {

    private final Queue<T> results = new ConcurrentLinkedQueue<>();

    MyScope() { super(null, Thread.ofVirtual().factory()); }

    @Override
    protected void handleComplete(Subtask<? extends T> subtask) {
        if (subtask.state() == Subtask.State.SUCCESS)
            results.add(subtask.get());
    }

    @Override
    public MyScope<T> join() throws InterruptedException {
        super.join();
        return this;
    }

    // Returns a stream of results from the subtasks that completed successfully
    public Stream<T> results() {
        super.ensureOwnerAndJoined();
        return results.stream();
    }

}

Эту собственную политику можно использовать так:

<T> List<T> allSuccessful(List<Callable<T>> tasks) throws InterruptedException {
    try (var scope = new MyScope<T>()) {
        for (var task : tasks) scope.fork(task);
        return scope.join()
                    .results().toList();
    }
}

Сценарии fan-in

Приведённые выше примеры были посвящены сценариям fan-out, в которых управляется несколько конкурентных исходящих операций ввода-вывода. StructuredTaskScope полезен и в сценариях fan-in, в которых управляется несколько конкурентных входящих операций ввода-вывода. В таких сценариях мы обычно создаём заранее неизвестное число подзадач в ответ на входящие запросы.

Вот пример сервера, который порождает подзадачи для обработки входящих соединений внутри StructuredTaskScope:

void serve(ServerSocket serverSocket) throws IOException, InterruptedException {
    try (var scope = new StructuredTaskScope<Void>()) {
        try {
            while (true) {
                var socket = serverSocket.accept();
                scope.fork(() -> handle(socket));
            }
        } finally {
            // If there's been an error or we're interrupted, we stop accepting
            scope.shutdown();  // Close all active connections
            scope.join();
        }
    }
}

С точки зрения конкурентности этот сценарий отличается не столько направлением запросов, сколько длительностью и количеством задач. Здесь, в отличие от предыдущих примеров, длительность работы владельца области не ограничена — он остановится, только когда его прервут. Количество подзадач тоже неизвестно, поскольку они порождаются динамически в ответ на внешние события.

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

Наблюдаемость

Мы расширяем новый формат дампа потоков в JSON, добавленный в JEP 444, чтобы показывать группировку потоков в иерархию, которую задаёт StructuredTaskScope:

$ jcmd <pid> Thread.dump_to_file -format=json <file>

JSON-объект каждой области содержит массив потоков, порождённых в этой области, вместе с их трассировками стека. Поток-владелец области, как правило, будет заблокирован в методе join в ожидании завершения подзадач; дамп потоков позволяет легко увидеть, чем заняты потоки подзадач, поскольку показывает древовидную иерархию, которую задаёт Structured Concurrency. JSON-объект области также содержит ссылку на родительскую область, чтобы по дампу можно было восстановить структуру программы.

API com.sun.management.HotSpotDiagnosticsMXBean также можно использовать для создания таких дампов потоков — напрямую или косвенно, через платформенный MBeanServer и локальный или удалённый инструмент JMX.

Почему fork(...) не возвращает Future?

Когда API StructuredTaskScope был в статусе Incubator, метод fork(...) возвращал Future. Это создавало ощущение привычности: fork(...) напоминал существующий метод ExecutorService::submit. Однако, поскольку StructuredTaskScope предназначен для использования иначе, чем ExecutorService, — структурированным образом, как описано выше, — использование Future вносило больше путаницы, чем ясности.

  • Привычное использование Future подразумевает вызов его метода get(), который блокируется до тех пор, пока не станет доступен результат. Но в контексте StructuredTaskScope такое использование Future не просто не рекомендуется, а контрпродуктивно. Структурированные объекты Future следует опрашивать только после возврата из join(), когда уже известно, что они завершены или отменены, и использовать при этом следует не привычный метод get(), а новый метод resultNow(), который никогда не блокируется.

  • Некоторые разработчики задавались вопросом, почему fork(...) не возвращает вместо этого более функциональные объекты CompletableFuture. Поскольку Future, возвращаемый fork(...), следует использовать только тогда, когда уже известно, что он завершён, CompletableFuture не дал бы никаких преимуществ: его расширенные возможности полезны только для незавершённых future. Кроме того, CompletableFuture спроектирован для парадигмы асинхронного программирования, тогда как StructuredTaskScope поощряет блокирующую парадигму.

    Словом, Future и CompletableFuture спроектированы так, чтобы давать степени свободы, которые в Structured Concurrency контрпродуктивны.

  • Суть Structured Concurrency в том, чтобы рассматривать несколько задач, выполняемых в разных потоках, как единую единицу работы, тогда как Future полезен в основном тогда, когда несколько задач рассматриваются как отдельные задачи. Область должна блокироваться только один раз, ожидая результатов своих подзадач, а затем централизованно обрабатывать исключения. Поэтому в подавляющем большинстве случаев единственным методом, который следовало вызывать у Future, возвращённого из fork(...), был resultNow(). Это заметно отличалось от обычного использования Future, и интерфейс Future отвлекал от правильного использования в этом контексте.

В текущем API Subtask::get() ведёт себя в точности так же, как вёл себя Future::resultNow(), когда API был в статусе Incubator.

Альтернативы

  • Расширить интерфейс ExecutorService. Мы создали прототип реализации этого интерфейса, которая всегда обеспечивает структуру и ограничивает круг потоков, которые могут отправлять задачи. Однако мы сочли этот подход проблематичным, поскольку большинство случаев использования ExecutorService (и его родительского интерфейса Executor) в JDK и в экосистеме не являются структурированными. Повторное использование того же API для гораздо более ограниченной концепции неизбежно приведёт к путанице. Например, передача структурированного экземпляра ExecutorService в существующие методы, принимающие этот тип, почти наверняка приводила бы к выбросу исключений в большинстве ситуаций.