JEP 499: Structured Concurrency (Fourth Preview)
Structured Concurrency (структурированная конкурентность), четвёртая версия Preview (предварительная версия)
| Автор | Ron Pressler & Alan Bateman |
| Ответственный | Alan Bateman |
| Тип | Feature |
| Область | SE |
| Статус | Closed / Delivered |
| Выпуск | 24 |
| Компонент | core-libs |
| Обсуждение | loom dash dev at openjdk dot org |
| Связан с | JEP 480: Structured Concurrency (Third Preview) |
| JEP 505: Structured Concurrency (Fifth Preview) | |
| Рецензенты | Paul Sandoz, Viktor Klang |
| Одобрен | Paul Sandoz |
| Создан | 2024/11/13 10:56 |
| Обновлён | 2025/02/11 19:43 |
| Задача | 8344096 |
Аннотация
Упростить конкурентное программирование с помощью API для Structured Concurrency. Structured Concurrency рассматривает группы связанных задач, которые выполняются в разных потоках, как единую единицу работы. Так упрощаются обработка ошибок и отмена, повышается надёжность и улучшается наблюдаемость. Это API в статусе Preview.
История
Structured Concurrency был предложен в JEP 428 и выпущен в JDK 19 как API в статусе Incubator (инкубационный модуль). Он был повторно выпущен в статусе Incubator в JEP 437 в JDK 20 с небольшим изменением: добавлено наследование Scoped Values (значения с ограниченной областью видимости; JEP 429). Впервые он вышел в статусе Preview в JDK 21 в JEP 453, где метод StructuredTaskScope::fork(...) стал возвращать Subtask, а не Future. Он повторно вышел в статусе Preview в JDK 22 в JEP 462 и в JDK 23 в JEP 480 без изменений.
Здесь мы предлагаем ещё раз выпустить API в статусе Preview в JDK 24 без изменений, чтобы было больше времени на отзывы по реальному использованию.
Цели
-
Продвигать стиль конкурентного программирования, который может устранить распространённые риски при отмене и завершении работы, например утечки потоков и задержки отмены.
-
Улучшить наблюдаемость конкурентного кода.
Что не является целью
-
Цель не в том, чтобы заменить какие-либо конструкции конкурентности из пакета
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-объектов, поэтому говорят, что задача присоединяет (join) свои подзадачи.
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 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 (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 24 --enable-preview Main.javaи запускайте её сjava --enable-preview Main; или -
При использовании средства запуска исходного кода запускайте программу с
java --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 такой:
-
Создайте область. Поток, который создаёт область, является её владельцем.
-
Используйте метод
fork(Callable), чтобы порождать подзадачи в области. -
В любой момент любая из подзадач или владелец области может вызвать метод
shutdown()области, чтобы отменить незавершённые подзадачи и запретить порождение новых подзадач. -
Владелец области присоединяет область, т. е. все её подзадачи, как единое целое. Владелец может вызвать метод
join()области, чтобы дождаться, пока все подзадачи либо завершатся (успешно или нет), либо будут отменены черезshutdown(). Либо он может вызвать методjoinUntil(java.time.Instant)области, чтобы ждать не дольше крайнего срока. -
После присоединения обработайте все ошибки в подзадачах и обработайте их результаты.
-
Закройте область, обычно неявно через
try-with-resources. Это завершает работу области, если она ещё не завершена, и ожидает завершения всех подзадач, которые были отменены, но ещё не завершились.
Каждый вызов fork(...) запускает новый поток для выполнения подзадачи; по умолчанию это поток Virtual Threads. Подзадача может создать собственную вложенную StructuredTaskScope, чтобы порождать свои подзадачи, и так образуется иерархия. Эта иерархия отражена в блочной структуре кода, которая ограничивает время жизни подзадач: гарантируется, что после закрытия области потоки всех подзадач завершены и ни один поток не остаётся после выхода из блока.
Любая подзадача в области, любые под-подзадачи во вложенной области, а также владелец области могут в любой момент вызвать метод shutdown() области, чтобы сообщить, что задача завершена, — даже пока другие подзадачи ещё выполняются. Метод shutdown() прерывает потоки, которые всё ещё выполняют подзадачи, и заставляет метод join() или joinUntil(Instant) вернуть управление. Поэтому все подзадачи следует писать так, чтобы они реагировали на прерывание. Новая подзадача, порождённая после вызова shutdown(), будет в состоянии UNAVAILABLE и не будет запущена. По сути, shutdown() — это конкурентный аналог оператора break в последовательном коде.
Вызов join() или joinUntil(Instant) внутри области обязателен. Если блок области завершается до ожидания (join), область дождётся завершения всех подзадач, а затем выбросит исключение.
Поток-владелец области может быть прерван до или во время ожидания. Например, он может быть подзадачей объемлющей области, которая была закрыта. В этом случае 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 напрямую, а один из двух подклассов, описанных в следующем разделе, которые реализуют политики закрытия. В других сценариях пользователи, скорее всего, будут писать собственные подклассы для реализации своих политик закрытия.
Политики закрытия
При работе с конкурентными подзадачами часто используют шаблоны с досрочным завершением (short-circuiting), чтобы избежать лишней работы. Например, иногда имеет смысл отменить все подзадачи, если одна из них завершилась с ошибкой (т. е. 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);
}
};
}
Собственные политики закрытия
Чтобы реализовать политики, отличные от политик ShutdownOnSuccess и ShutdownOnFailure, можно расширить StructuredTaskScope и переопределить его защищённый метод handleComplete(...). Подкласс может, например,
- собирать результаты успешно завершившихся подзадач и игнорировать подзадачи, завершившиеся с ошибкой,
- собирать исключения, когда подзадачи завершаются с ошибкой, или
- вызывать метод
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в существующие методы, принимающие этот тип, почти наверняка приводила бы к исключениям в большинстве ситуаций.