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

JEP 505: Structured Concurrency (Fifth Preview)

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

AuthorsAlan Bateman, Viktor Klang, & Ron Pressler
ОтветственныйAlan Bateman
ТипFeature
ОбластьSE
СтатусClosed / Delivered
Выпуск25
Компонентcore-libs
Обсуждениеloom dash dev at openjdk dot org
Связан сJEP 499: Structured Concurrency (Fourth Preview)
JEP 525: Structured Concurrency (Sixth Preview)
РецензентыPaul Sandoz
ОдобренPaul Sandoz
Создан2024/09/18 04:58
Обновлён2025/10/09 07:21
Задача8340343

Аннотация

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

История

Structured Concurrency находилась в статусе Incubator (инкубационный модуль) в JDK 19 в рамках JEP 428 и в JDK 20 в рамках JEP 437. В статусе Preview она вышла в JDK 21 в рамках JEP 453, при этом метод fork был изменён так, чтобы возвращать Subtask, а не Future. Повторно в статусе Preview она выходила в JDK 22 в рамках JEP 462, в JDK 23 в рамках JEP 480 и в JDK 24 в рамках JEP 499.

Мы предлагаем ещё раз выпустить API в статусе Preview в JDK 25 с рядом изменений API. В частности, StructuredTaskScope теперь открывается статическими фабричными методами, а не публичными конструкторами. Фабричный метод open без параметров покрывает типичный случай: он создаёт StructuredTaskScope, который ожидает, пока все подзадачи завершатся успешно или какая-либо подзадача завершится с ошибкой. Другие политики и результаты можно реализовать, передав подходящий Joiner одному из более гибких фабричных методов open.

Цели

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

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

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

  • Цель не в том, чтобы заменить какие-либо конкурентные конструкции из пакета 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 для каждой подзадачи и выполняет подзадачи конкурентно в соответствии со своей политикой планирования. Метод handle() ожидает результаты подзадач с помощью блокирующих вызовов методов get() их future-объектов, поэтому говорят, что задача присоединяет (join) свои подзадачи.

Response handle() throws ExecutionException, InterruptedException {
    Future<String> user = executor.submit(() -> findUser());
    Future<Integer> order = executor.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 и в блоке catch для задачи, завершившейся сбоем, вызывать методы cancel(boolean) future-объектов других задач. Также потребовалось бы использовать ExecutorService внутри оператора try-with-resources, как показано в примерах в JEP 444, поскольку 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 для нескольких потоков основана на двух идеях: у потока выполнения через блок кода есть чётко определённые точки входа и выхода, а время жизни операций вложено так же, как их синтаксическая вложенность в коде.

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

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

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

Описание

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

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

Response handle() throws InterruptedException {

    try (var scope = StructuredTaskScope.open()) {

        Subtask<String> user = scope.fork(() -> findUser());
        Subtask<Integer> order = scope.fork(() -> fetchOrder());

        scope.join();   // Join subtasks, propagating exceptions

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

    }

}

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

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

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

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

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

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

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

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

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

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

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

API StructuredTaskScope<T, R>, где T — тип результата задач, запущенных в области видимости, а R — тип результата метода join, можно кратко описать так:

public sealed interface StructuredTaskScope<T, R> extends AutoCloseable {

    public static <T> StructuredTaskScope<T, Void> open();
    public static <T, R> StructuredTaskScope<T, R> open(Joiner<? super T,
                                                               ? extends R> joiner);

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

    public R join() throws InterruptedException;

    public void close();

}

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

  1. Откройте новую область видимости, вызвав один из статических методов open. Поток, открывающий область видимости, становится её владельцем.

  2. Порождайте подзадачи в области видимости с помощью методов fork.

  3. Объедините все подзадачи области видимости как единое целое с помощью метода join. При этом может быть выброшено исключение.

  4. Обработайте итог.

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

В примере handle() фабричный метод без параметров open() создаёт и открывает StructuredTaskScope, реализующий политику завершения по умолчанию: сбой, если какая-либо подзадача завершилась сбоем. Можно реализовать и другие политики, как мы увидим ниже, передав подходящий Joiner в метод open с одним параметром.

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

Метод join должен вызываться потоком-владельцем области видимости изнутри этой области. Если выход из блока области видимости происходит до объединения, то область видимости отменяется, и владелец будет ждать в методе close завершения всех подзадач, прежде чем выбросить исключение.

После объединения владелец области видимости может обработать результаты подзадач с помощью объектов Subtask, возвращаемых методами fork. Метод Subtask::get выбрасывает исключение, если вызван до объединения.

Отмена

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

Чтобы отмена была возможна, подзадачи нужно писать так, чтобы при прерывании они завершались как можно скорее. Подзадачи, которые не реагируют на прерывания, например потому что блокируются в непрерываемых методах, могут бесконечно задерживать закрытие области видимости. Метод close всегда ждёт завершения потоков, выполняющих подзадачи, даже если область видимости отменена. Выполнение не может продолжиться дальше метода close, пока прерванные потоки не завершатся.

Scoped Values (значения с ограниченной областью видимости)

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

Структурное использование обеспечивается принудительно

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

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

Объекты Joiner

В примере handle(), если какая-либо подзадача завершается сбоем, метод join() выбрасывает исключение, и область видимости отменяется. Если все подзадачи завершаются успешно, метод join() завершается нормально и возвращает null. Это политика завершения по умолчанию.

Другие политики можно выбрать, создав StructuredTaskScope с подходящим StructuredTaskScoped.Joiner. Объект Joiner обрабатывает завершение подзадач и формирует результат для метода join(). В зависимости от joiner-объекта метод join() может вернуть результат, stream элементов или какой-либо другой объект.

Интерфейс Joiner объявляет фабричные методы, создающие joiner-объекты для некоторых типичных случаев. Например, фабричный метод anySuccessfulResultOrThrow() возвращает новый joiner-объект, который выдаёт результат любой подзадачи, завершившейся успешно:

<T> T race(Collection<Callable<T>> tasks) throws InterruptedException {
    try (var scope = StructuredTaskScope.open(Joiner.<T>anySuccessfulResultOrThrow())) {
        tasks.forEach(scope::fork);
        return scope.join();
    }
}

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

Фабричный метод allSuccessfulOrThrow() возвращает новый joiner-объект, который, когда все подзадачи завершаются успешно, выдаёт stream подзадач:

<T> List<T> runConcurrently(Collection<Callable<T>> tasks) throws InterruptedException {
    try (var scope = StructuredTaskScope.open(Joiner.<T>allSuccessfulOrThrow())) {
        tasks.forEach(scope::fork);
        return scope.join().map(Subtask::get).toList();
    }
}

Если одна или несколько подзадач завершаются с ошибкой, join() выбрасывает FailedException с исключением одной из неудачных подзадач в качестве причины. Политика завершения, реализуемая Joiner, та же, что и политика по умолчанию, реализуемая методом open() без параметров. Различаются они результатом: метод join() возвращает stream завершённых подзадач, а не null, поэтому этот joiner-объект подходит для случаев, когда все подзадачи возвращают результат одного и того же типа и когда объекты Subtask, возвращаемые методом fork, игнорируются. В этом примере элементы Subtask в stream отображаются в результат и собираются в список.

Интерфейс Joiner объявляет ещё три фабричных метода:

  • awaitAll(), который возвращает новый joiner-объект, просто ожидающий завершения всех подзадач, успешного или нет;

  • awaitAllSuccessfulOrThrow(), который возвращает новый joiner-объект, ожидающий успешного завершения всех подзадач; и

  • allUntil(Predicate<Subtask<? extends T>> isDone), который возвращает новый joiner-объект, который, когда все подзадачи завершаются успешно или же предикат для завершённой подзадачи возвращает true, отменяет объемлющую область видимости и выдаёт stream всех подзадач.

При использовании любого вида Joiner крайне важно создавать новый Joiner для каждого StructuredTaskScope. Объекты Joiner никогда не следует использовать в разных областях видимости задач или повторно использовать после закрытия области видимости.

Собственные объекты Joiner

Интерфейс Joiner можно реализовать напрямую, чтобы поддержать собственные политики завершения. У него два параметра типа: T — тип результата подзадач, выполняемых в области видимости, и R — тип результата метода join(). Интерфейс можно кратко описать так:

public interface Joiner<T, R> {
    public default boolean onFork(Subtask<? extends T> subtask);
    public default boolean onComplete(Subtask<? extends T> subtask);
    public R result() throws Throwable;
}

Метод onFork вызывается при запуске подзадачи, а метод onComplete — при завершении подзадачи. Оба метода возвращают boolean, указывающий, нужно ли отменить область видимости. Метод result вызывается, чтобы сформировать результат для метода join или же выбросить исключение, когда все подзадачи завершены или область видимости отменена. Если метод result выбрасывает исключение, то метод join выбросит FailedException с этим исключением в качестве причины.

Вот класс Joiner, который собирает результаты подзадач, завершившихся успешно, и игнорирует подзадачи, завершившиеся с ошибкой. Метод onComplete может вызываться несколькими потоками одновременно, поэтому он должен быть потокобезопасным. Метод result возвращает stream результатов задач.

class CollectingJoiner<T> implements Joiner<T, Stream<T>> {

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

    public boolean onComplete(Subtask<? extends T> subtask) {
        if (subtask.state() == Subtask.State.SUCCESS) {
            results.add(subtask.get());
        }
        return false;
    }

    public Stream<T> result() {
        return results.stream();
    }

}

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

<T> List<T> allSuccessful(List<Callable<T>> tasks) throws InterruptedException {
      try (var scope = StructuredTaskScope.open(new CollectingJoiner<T>())) {
          tasks.forEach(scope::fork);
          return scope.join().toList();
      }
  }

Обработка исключений

Способ обработки исключений зависит от сценария использования. Метод join() выбрасывает FailedException, когда область видимости считается завершившейся неудачно. В примере handle(), если подзадача завершается с ошибкой, выбрасывается FailedException с исключением неудачной подзадачи в качестве причины. В некоторых случаях может быть полезно добавить блок catch к оператору try-with-resources, чтобы обрабатывать исключения после закрытия области видимости:

try (var scope = StructuredTaskScope.open()) {
   ...
} catch (StructuredTaskScope.FailedException e) {
   Throwable cause = e.getCause();
   switch (cause) {
       case IOException ioe -> ..
       default -> ..
   }
}

Код обработки исключений может использовать оператор instanceof с Pattern Matching (сопоставление с образцом) (JEP 394), чтобы обрабатывать конкретные причины.

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

Конфигурация

В приведённом ранее кратком описании API StructuredTaskScope были показаны два статических метода open. Третий такой метод принимает Joiner вместе с функцией, которая может сформировать объект конфигурации, чтобы задать имя области видимости для целей мониторинга и управления, задать тайм-аут области видимости и задать фабрику потоков, которую методы fork области видимости будут использовать для создания потоков,

Вот изменённая версия метода runConcurrently, которая задаёт фабрику потоков и тайм-аут:

<T> List<T> runConcurrently(Collection<Callable<T>> tasks,
                            ThreadFactory factory,
                            Duration timeout)
    throws InterruptedException
{
    try (var scope = StructuredTaskScope.open(Joiner.<T>allSuccessfulOrThrow(),
                                              cf -> cf.withThreadFactory(factory)
                                                      .withTimeout(timeout))) {
        tasks.forEach(scope::fork);
        return scope.join().map(Subtask::get).toList();
    }
}

Метод fork в этой области видимости будет вызывать заданную фабрику потоков, чтобы создать поток для выполнения каждой подзадачи. Это может быть полезно, например, чтобы задать имя потока или другие его свойства.

timeout задаётся как java.time.Duration. Если тайм-аут истекает до или во время ожидания в методе join(), то область видимости отменяется, что отменяет все незавершённые подзадачи, и join() выбрасывает TimeoutException.

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

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

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

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

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

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

Расширить интерфейс ExecutorService

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

Возвращать Future из методов fork

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

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

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

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

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

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