JEP 533: Structured Concurrency (Seventh Preview)
Structured Concurrency (структурированная конкурентность), седьмая версия Preview (предварительная версия)
| Authors | Alan Bateman, Viktor Klang, & Ron Pressler |
| Ответственный | Alan Bateman |
| Тип | Feature |
| Область | SE |
| Статус | Closed / Delivered |
| Выпуск | 27 |
| Компонент | core-libs |
| Обсуждение | loom dash dev at openjdk dot org |
| Трудоёмкость | S |
| Связан с | JEP 525: Structured Concurrency (Sixth Preview) |
| Рецензенты | Viktor Klang |
| Одобрен | Paul Sandoz |
| Создан | 2025/12/12 14:51 |
| Обновлён | 2026/08/19 06:20 |
| Задача | 8373610 |
Аннотация
Упростить конкурентное программирование, добавив API для Structured Concurrency. Structured Concurrency рассматривает группы связанных задач, выполняемых в разных потоках, как единые единицы работы. Это упрощает обработку ошибок и отмену, повышает надёжность и улучшает наблюдаемость. Это API в статусе Preview.
История
Structured Concurrency впервые появилась в виде модуля Incubator (инкубационный модуль) в JEP 428 (JDK 19), а затем в JEP 437 (JDK 20). Она стала API в статусе Preview в JEP 453 (JDK 21), при этом метод fork стал возвращать Subtask, а не Future. В этом виде она снова вышла в статусе Preview в JEP 462 (JDK 22), JEP 480 (JDK 23) и JEP 499 (JDK 24). Затем она снова вышла в статусе Preview в JEP 505 (JDK 25) с рядом изменений API. Самое заметное из них: публичные конструкторы StructuredTaskScope заменены статическими фабричными методами. Ещё раз в статусе Preview она вышла в JEP 525 (JDK 26) с несколькими небольшими изменениями.
Мы предлагаем ещё раз выпустить этот API в статусе Preview в JDK 27 со следующими изменениями:
-
У интерфейсов
StructuredTaskScopeиJoinerтеперь есть третий параметр типа,R_X. Это тип исключения, которое может выбросить методjoin()вStructuredTaskScope. -
Новый статический метод
openвStructuredTaskScopeреализует политику объединения по умолчанию и с помощью заданногоUnaryOperatorсоздаёт конфигурациюStructuredTaskScope. -
Фабричные методы
allSuccessfulOrThrow(),anySuccessfulOrThrow()иawaitAllSuccessfulOrThrow()вJoinerтеперь создают объекты joiner, с которымиjoin()выбрасываетExecutionException, если результатом является исключение. Новые перегрузки этих трёх методов позволяют указатьFunction, чтобы создавать другое исключение. -
Фабричный метод
awaitAll()вJoinerудалён. -
Метод
onTimeout()интерфейсаJoinerзаменён методомtimeout(). Когда область (scope) отменяется по тайм-ауту, этот метод либо создаёт результат, либо выбрасывает исключение. Если методtimeout()выбрасывает исключение, то выбрасывается исключение, причиной которого указанCancelledByTimeoutException.
Цели
-
Продвигать стиль конкурентного программирования, который может устранить распространённые риски при отмене и завершении работы, например утечки потоков и задержки отмены.
-
Улучшить наблюдаемость конкурентного кода.
Что не является целью
-
Цель не в том, чтобы заменить какие-либо конкурентные конструкции из пакета
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 ExecutionException, 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 27 --enable-preview Main.javaи запустите её сjava --enable-preview Main; или -
при использовании средства запуска исходного кода запустите программу с
java --enable-preview Main.java; или -
при использовании jshell запустите его с
jshell --enable-preview.
Использование StructuredTaskScope
У интерфейса StructuredTaskScope три параметра типа: T — тип результата подзадач, порождённых в области видимости, R — тип результата, возвращаемого join(), а R_X — тип исключения, которое может выбросить join().
public sealed interface StructuredTaskScope<T, R, R_X extends Throwable>
extends AutoCloseable
{
public static <T> StructuredTaskScope<T, Void, ExecutionException> open();
public static <T, R, R_X extends Throwable> StructuredTaskScope<T, R, R_X> open(
Joiner<? super T, ? extends R, R_X> joiner);
public <U extends T> Subtask<U> fork(Callable<? extends U> task);
public Subtask<? extends T> fork(Runnable task);
public R join() throws R_X, InterruptedException;
public void close();
}
Общий порядок работы кода с StructuredTaskScope такой:
-
Откройте новую область видимости, вызвав один из статических методов
open. Поток, открывающий область видимости, становится её владельцем. -
Порождайте подзадачи в области видимости с помощью методов
fork. -
Объедините все подзадачи области видимости как единое целое с помощью метода
join. При этом может быть выброшено исключение. -
Обработайте итог.
-
Закройте область видимости, обычно неявно через
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 506). Если владелец области видимости читает значение из привязанного ScopedValue, то каждая подзадача прочитает то же значение.
Структурное использование обеспечивается принудительно
Во время выполнения StructuredTaskScope навязывает конкурентным операциям структуру и порядок. Например, попытки вызвать метод fork из потока, который не является владельцем области видимости, завершатся исключением. Использование области видимости вне блока try-with-resources и возврат без вызова close(), или без соблюдения правильной вложенности вызовов close(), могут привести к тому, что методы области видимости выбросят StructureViolationException.
StructuredTaskScope не реализует интерфейсы ExecutorService и Executor, поскольку экземпляры этих интерфейсов обычно используются неструктурированно (см. ниже). Однако код, который использует ExecutorService, но выиграл бы от структуры, несложно перевести на StructuredTaskScope.
Объекты Joiner
В примере handle(), если какая-либо подзадача завершается сбоем, метод join выбрасывает исключение, и область видимости отменяется. Если все подзадачи завершаются успешно, метод join завершается нормально и возвращает null. Это политика завершения по умолчанию.
Другие политики можно выбрать, создав StructuredTaskScope с подходящим StructuredTaskScope.Joiner. Объект Joiner формирует итог для метода join. В зависимости от Joiner метод join может вернуть результат, список элементов или какой-то другой объект. Если итог — исключение, то тип выбрасываемого исключения зависит от Joiner.
Интерфейс Joiner объявляет фабричные методы для создания объектов Joiner для некоторых распространённых случаев. Например, фабричный метод anySuccessfulOrThrow() возвращает новый Joiner, который выдаёт результат любой успешно завершившейся подзадачи. Он выбрасывает ExecutionException, если все подзадачи завершились сбоем.
<T> T race(Collection<Callable<T>> tasks)
throws ExecutionException, InterruptedException
{
try (var scope = StructuredTaskScope.open(Joiner.<T>anySuccessfulOrThrow())) {
tasks.forEach(scope::fork);
return scope.join();
}
}
Как только одна подзадача завершается успешно, область видимости отменяется, при этом отменяются незавершённые подзадачи, а join() возвращает результат успешной подзадачи. Этот шаблон может быть полезен, например, в серверных приложениях, которым нужен результат от любого из набора дублирующих друг друга сервисов.
Фабричный метод allSuccessfulOrThrow() возвращает новый Joiner, который, когда все подзадачи завершились успешно, выдаёт список их результатов. Если одна или несколько подзадач завершились сбоем, он заставляет join() выбросить ExecutionException, причиной которого указано исключение одной из неудавшихся подзадач.
<T> List<T> runConcurrently(Collection<Callable<T>> tasks)
throws ExecutionException, InterruptedException {
try (var scope = StructuredTaskScope.open(Joiner.<T>allSuccessfulOrThrow())) {
tasks.forEach(scope::fork);
return scope.join();
}
}
Политика завершения, реализуемая этим Joiner, та же, что и политика по умолчанию, реализуемая методом без параметров open(). Они различаются итогом: метод join возвращает список результатов, а не null, поэтому этот Joiner подходит для случаев, когда все подзадачи возвращают результат одного типа, а объекты Subtask, возвращаемые методом fork, игнорируются.
Интерфейс Joiner объявляет дополнительные фабричные методы для создания объектов Joiner:
-
awaitAllSuccessfulOrThrow()возвращает новый Joiner, который ждёт успешного завершения всех подзадач; -
anySuccessfulOrThrow(Function),allSuccessfulOrThrow(Function)иawaitAllSuccessfulOrThrow(Function)эквивалентны соответствующим фабричным методам без параметров, за исключением того, что, когда итог — исключение, они заставляютjoin()выбросить исключение, созданное заданной функцией, поставляющей исключения; и -
allUntil(Predicate<? super Subtask<T>> isDone)возвращает новый Joiner, который, когда все подзадачи завершены или предикат для завершившейся подзадачи возвращаетtrue, отменяет область видимости и заставляет методjoinвернуть список всех подзадач.
При использовании любого вида Joiner крайне важно создавать новый Joiner для каждого StructuredTaskScope. Объекты Joiner никогда не следует использовать в разных областях видимости задач или повторно использовать после закрытия области видимости.
Собственные объекты Joiner
Интерфейс Joiner можно реализовать напрямую, чтобы поддержать собственные политики завершения. У него три параметра типа: T — тип результата подзадач, запущенных в области видимости, R — тип результата, возвращаемого join(), и R_X — тип исключения, которое может выбросить join().
public interface Joiner<T, R, R_X extends Throwable> {
public default boolean onFork(Subtask<T> subtask);
public default boolean onComplete(Subtask<T> subtask);
public R result() throws R_X;
public R timeout() throws R_X;
}
Метод onFork вызывается при запуске подзадачи, а метод onComplete — при завершении подзадачи. Оба метода возвращают boolean, чтобы указать, следует ли отменить область видимости. Метод result вызывается, чтобы получить итог метода join, то есть результат или исключение. Метод timeout вызывается, если область видимости открыта с тайм-аутом (см. ниже) и тайм-аут истекает до или во время вызова метода join.
Вот класс Joiner, который собирает результаты успешно завершившихся подзадач и игнорирует подзадачи, завершившиеся с ошибкой. Метод onComplete может вызываться несколькими потоками одновременно, поэтому он должен быть потокобезопасным. Метод result возвращает список результатов задач. Метод timeout выбрасывает исключение времени выполнения, а именно CompletionException, если область видимости открыта с тайм-аутом и тайм-аут истекает.
class CollectingJoiner<T> implements Joiner<T, List<T>, CompletionException> {
private final Queue<T> results = new ConcurrentLinkedQueue<>();
public boolean onComplete(Subtask<T> subtask) {
if (subtask.state() == Subtask.State.SUCCESS) {
results.add(subtask.get());
}
return false;
}
public List<T> result() {
return List.copyOf(results);
}
public List<T> timeout() {
throw new CompletionException(new CancelledByTimeoutException());
}
}
Эту собственную политику можно использовать так:
<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();
}
}
Обработка исключений
Способ обработки исключений зависит от сценария использования. Метод join выбрасывает исключение типа R_X, когда область видимости считается завершившейся неудачно. В примере handle(), где R_X — это ExecutionException, если подзадача завершается с ошибкой, выбрасывается ExecutionException, причиной которого служит исключение из этой подзадачи. В некоторых случаях может быть полезно добавить блок catch к оператору try-with-resources, чтобы обрабатывать исключения после закрытия области видимости:
try (var scope = StructuredTaskScope.open()) {
...
} catch (ExecutionException e) {
Throwable cause = e.getCause();
switch (cause) {
case IOException ioe -> ..
default -> ..
}
}
Код обработки исключений может использовать оператор instanceof с Pattern Matching (сопоставление с образцом) (JEP 394), чтобы обрабатывать конкретные причины.
Определённое исключение в подзадаче может приводить к возврату значения по умолчанию. В таких случаях может быть уместнее перехватить исключение в самой подзадаче и завершить её со значением по умолчанию в качестве результата, а не заставлять владельца области видимости обрабатывать исключение.
Конфигурация
В приведённом ранее обзоре API StructuredTaskScope были показаны два статических метода open. Есть ещё два метода open:
-
open(UnaryOperator)применяет заданный оператор к объекту конфигурации по умолчанию, чтобы задать имя области видимости для мониторинга и управления, тайм-аут области видимости и фабрику потоков, с помощью которой методыforkобласти видимости будут создавать потоки. -
open(Joiner, UnaryOperator)принимает и Joiner, и оператор для задания конфигурации.
Вот изменённая версия метода runConcurrently, которая задаёт фабрику потоков и тайм-аут:
<T> List<T> runConcurrently(Collection<Callable<T>> tasks,
ThreadFactory factory,
Duration timeout)
throws ExecutionException, InterruptedException
{
try (var scope = StructuredTaskScope.open(Joiner.<T>allSuccessfulOrThrow(),
cf -> cf.withThreadFactory(factory)
.withTimeout(timeout))) {
tasks.forEach(scope::fork);
return scope.join();
}
}
Метод fork в этой области видимости будет вызывать заданную фабрику потоков, чтобы создать поток для выполнения каждой подзадачи. Это может быть полезно, например, чтобы задать имя потока или другие его свойства.
timeout задаётся как java.time.Duration. В этом примере, если тайм-аут истекает до или во время ожидания в join(), область видимости отменяется, что отменяет все незавершённые подзадачи, и join() выбрасывает ExecutionException с CancelledByTimeoutException в качестве причины.
Наблюдаемость
Мы расширяем формат дампа потоков в 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.