JEP 462: Structured Concurrency (Second Preview)
Structured Concurrency (структурированная конкурентность), вторая версия Preview (предварительная версия)
| Автор | Ron Pressler & Alan Bateman |
| Ответственный | Alan Bateman |
| Тип | Feature |
| Область | SE |
| Статус | Closed / Delivered |
| Выпуск | 22 |
| Компонент | core-libs |
| Обсуждение | loom dash dev at openjdk dot org |
| Связан с | JEP 453: Structured Concurrency (Preview) |
| JEP 480: Structured Concurrency (Third Preview) | |
| Рецензенты | Paul Sandoz |
| Одобрен | Paul Sandoz |
| Создан | 2023/09/29 09:37 |
| Обновлён | 2025/02/27 17:47 |
| Задача | 8317302 |
Аннотация
Упростить конкурентное программирование, добавив API для Structured Concurrency. Structured Concurrency рассматривает группы связанных задач, выполняемых в разных потоках, как единую единицу работы. Так упрощается обработка ошибок и отмена, повышается надёжность и улучшается наблюдаемость. Это API в статусе Preview.
История
Structured Concurrency была предложена в JEP 428 и поставлена в JDK 19 как API в статусе Incubator (инкубационный модуль). Повторно в статусе Incubator она вышла в JDK 20 по JEP 437 с небольшим изменением: добавилось наследование Scoped Values (значения с ограниченной областью видимости, JEP 429). Впервые в статусе Preview она появилась в JDK 21 по JEP 453. При этом StructuredTaskScope::fork(...) стал возвращать Subtask, а не Future. Здесь мы предлагаем повторно выпустить этот API в статусе Preview в JDK 22 без изменений, чтобы получить больше отзывов.
Цели
-
Продвигать стиль конкурентного программирования, который может устранить распространённые риски при отмене и завершении работы, например утечки потоков и задержки отмены.
-
Улучшить наблюдаемость конкурентного кода.
Что не является целью
-
Цель не в том, чтобы заменить какие-либо конструкции для конкурентности из пакета
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 и в блоке catch для упавшей задачи вызывать методы cancel(boolean) у объектов future остальных задач. Кроме того, ExecutorService пришлось бы использовать внутри оператора try-with-resources, как показано в примерах в JEP 425, потому что Future не даёт способа дождаться отменённой задачи. Но сделать всё это правильно может быть очень непросто, и из-за этого часто труднее понять логический замысел кода. Отслеживать отношения между задачами и вручную добавлять обратно нужные связи для отмены между задачами — слишком много требовать от разработчиков.
Необходимость вручную согласовывать время жизни возникает потому, что ExecutorService и Future допускают произвольные схемы конкурентности. На участвующие потоки не накладывается никаких ограничений и никакого порядка. Один поток может создать ExecutorService, второй поток может отправлять в него работу, а потоки, которые выполняют эту работу, никак не связаны ни с первым, ни со вторым потоком. Более того, после того как поток отправил работу, результатов выполнения может ожидать совершенно другой поток. Любой код, у которого есть ссылка на Future, может присоединить его (то есть дождаться его результата, вызвав get()), даже код в потоке, отличном от того, который получил Future. По сути, подзадача, запущенная одной задачей, не обязана возвращаться к задаче, которая её отправила. Она может вернуться к любой из множества задач или даже ни к одной.
Поскольку ExecutorService и Future допускают такое неструктурированное использование, они не обеспечивают и даже не отслеживают отношения между задачами и подзадачами, хотя такие отношения распространены и полезны. Поэтому даже если подзадачи отправляются и присоединяются в одной и той же задаче, сбой одной подзадачи не может автоматически вызвать отмену другой: в приведённом выше методе handle() сбой fetchOrder() не может автоматически вызвать отмену findUser(). Future для fetchOrder() никак не связан с future для findUser(), и ни один из них не связан с потоком, который в итоге присоединит его через свой метод get(). Мы хотим не заставлять разработчиков управлять такой отменой вручную, а надёжно её автоматизировать.
Структура задач должна отражать структуру кода
В отличие от беспорядочного набора потоков в ExecutorService, выполнение однопоточного кода всегда обеспечивает иерархию задач и подзадач. Тело метода {...} соответствует задаче, а методы, вызываемые внутри этого блока, — подзадачам. Вызванный метод должен либо вернуть управление вызвавшему его методу, либо выбросить в него исключение. Он не может пережить вызвавший его метод, а также не может вернуть управление или выбросить исключение в другой метод. Таким образом, все подзадачи завершаются раньше задачи, каждая подзадача — дочерняя по отношению к своей родительской задаче, а время жизни каждой подзадачи относительно других подзадач и задачи определяется синтаксической блочной структурой кода.
Например, в этой однопоточной версии handle() отношение «задача — подзадача» очевидно из синтаксической структуры:
Response handle() throws IOException {
String theUser = findUser();
int theOrder = fetchOrder();
return new Response(theUser, theOrder);
}
Мы не запускаем подзадачу fetchOrder(), пока не завершится подзадача findUser(), успешно или сбоем. Если findUser() завершается сбоем, мы вообще не запускаем fetchOrder(), и задача handle() неявно завершается сбоем. Важно то, что подзадача может вернуться только к своей родительской задаче: это означает, что родительская задача может неявно считать сбой одной подзадачи сигналом отменить другие незавершённые подзадачи, а затем сама завершиться сбоем.
В однопоточном коде иерархия «задача — подзадача» во время выполнения овеществляется в стеке вызовов. Поэтому соответствующие отношения «родитель — потомок», которые определяют распространение ошибок, мы получаем бесплатно. При наблюдении за одним потоком иерархическое отношение очевидно: findUser() (а позже fetchOrder()) выглядят подчинёнными handle(). Поэтому легко ответить на вопрос «над чем сейчас работает handle()?»
Конкурентное программирование было бы проще, надёжнее и удобнее для наблюдения, если бы отношения «родитель — потомок» между задачами и их подзадачами были видны из синтаксической структуры кода, а также материализовались во время выполнения — так же, как в однопоточном коде. Синтаксическая структура очерчивала бы время жизни подзадач и позволяла бы представить во время выполнения иерархию между потоками, аналогичную стеку вызовов внутри потока. Такое представление обеспечило бы распространение ошибок и отмену, а также осмысленное наблюдение за конкурентной программой.
(В платформе Java уже есть API, задающий структуру конкурентным задачам, а именно java.util.concurrent.ForkJoinPool — механизм выполнения, на котором работают параллельные потоки данных (parallel streams). Однако этот API предназначен для вычислительно ёмких задач, а не для задач, связанных с вводом-выводом.)
Structured Concurrency
Structured Concurrency — это подход к конкурентному программированию, сохраняющий естественную связь между задачами и подзадачами, благодаря чему конкурентный код становится более читаемым, сопровождаемым и надёжным. Термин «Structured Concurrency» ввёл Martin Sústrik, а популяризировал Nathaniel J. Smith. На проектирование обработки ошибок в Structured Concurrency повлияли идеи из других языков, например иерархические супервизоры Erlang.
Structured Concurrency исходит из простого принципа:
Если задача разделяется на конкурентные подзадачи, то все они возвращаются в одно и то же место, а именно в блок кода задачи.
В Structured Concurrency подзадачи работают от имени задачи. Задача ожидает результатов подзадач и следит за их сбоями. Как и в приёмах структурного программирования для кода в одном потоке, сила Structured Concurrency для нескольких потоков основана на двух идеях: (1) чётко определённые точки входа в блок кода и выхода из него для потока выполнения и (2) строгая вложенность времени жизни операций, повторяющая их синтаксическую вложенность в коде.
Поскольку точки входа в блок кода и выхода из него чётко определены, время жизни конкурентной подзадачи ограничено синтаксическим блоком её родительской задачи. Поскольку время жизни подзадач одного уровня вложено во время жизни их родительской задачи, о них можно рассуждать и управлять ими как единым целым. Поскольку время жизни родительской задачи, в свою очередь, вложено во время жизни её родителя, среда выполнения может овеществить иерархию задач в виде дерева — конкурентного аналога стека вызовов одного потока. Так код может применять политики, например крайние сроки, ко всему поддереву задач, а инструменты наблюдения могут показывать подзадачи подчинёнными их родительским задачам.
Structured Concurrency отлично сочетается с Virtual Threads — легковесными потоками, реализованными в JDK. Много потоков Virtual Threads используют один и тот же поток операционной системы, поэтому их может быть очень много. Помимо того что потоков Virtual Threads может быть много, они достаточно дёшевы, чтобы представлять любую конкурентную единицу поведения, даже поведение, связанное с вводом-выводом. Это значит, что серверное приложение может с помощью Structured Concurrency обрабатывать одновременно тысячи или миллионы входящих запросов: оно может выделить новый поток Virtual Threads под задачу обработки каждого запроса, а когда задача разветвляется, отправляя подзадачи на конкурентное выполнение, — выделить новый поток Virtual Threads под каждую подзадачу. Внутри отношение «задача — подзадача» овеществляется в виде дерева: каждый поток Virtual Threads хранит ссылку на своего единственного родителя, подобно тому как кадр стека вызовов ссылается на своего единственного вызывающего.
Итак, Virtual Threads дают потоки в изобилии. Structured Concurrency может правильно и надёжно согласовывать их работу и позволяет инструментам наблюдения показывать потоки так, как их понимает разработчик. API для Structured Concurrency в JDK упростит создание сопровождаемых, надёжных и наблюдаемых серверных приложений.
Описание
Основной класс 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 21 --enable-preview Main.javaи запускайте её сjava --enable-preview Main; или -
При использовании средства запуска исходного кода запускайте программу с
java --source 21 --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 для области, то есть для всех её подзадач, как для единого целого. Владелец может вызвать метод
join()области, чтобы дождаться, пока все подзадачи либо завершатся (успешно или нет), либо будут отменены черезshutdown(). Кроме того, он может вызвать методjoinUntil(java.time.Instant)области, чтобы ждать не дольше заданного срока. -
После join обработайте ошибки в подзадачах и обработайте их результаты.
-
Закройте область, обычно неявно, с помощью
try-with-resources. Это завершает работу области, если она ещё не завершена, и ожидает, пока завершатся подзадачи, которые были отменены, но ещё не завершились.
Каждый вызов fork(...) запускает новый поток для выполнения подзадачи; по умолчанию это поток Virtual Threads. Подзадача может создать собственную вложенную StructuredTaskScope, чтобы порождать свои подзадачи, и так образуется иерархия. Эта иерархия отражена в блочной структуре кода, которая ограничивает время жизни подзадач: гарантируется, что после закрытия области все потоки подзадач завершены, и при выходе из блока ни один поток не остаётся.
Любая подзадача в области, любые под-подзадачи во вложенной области и владелец области могут в любой момент вызвать метод shutdown() области, чтобы сообщить, что задача завершена, — даже если другие подзадачи ещё выполняются. Метод shutdown() прерывает потоки, которые ещё выполняют подзадачи, и приводит к возврату из метода join() или joinUntil(Instant). Поэтому все подзадачи следует писать так, чтобы они реагировали на прерывание. Новая подзадача, порождённая после вызова shutdown(), будет в состоянии UNAVAILABLE и не будет запущена. По сути, shutdown() — конкурентный аналог оператора break в последовательном коде.
Вызов join() или joinUntil(Instant) внутри области обязателен. Если выход из блока области произойдёт до join, область дождётся завершения всех подзадач, а затем выбросит исключение.
Поток-владелец области может быть прерван до join или во время него. Например, он может быть подзадачей объемлющей области, работа которой была завершена. В этом случае join() и joinUntil(Instant) выбросят исключение, поскольку продолжать нет смысла. Затем оператор try-with-resources завершит работу области, что отменит все подзадачи и дождётся их завершения. В результате отмена задачи автоматически распространяется на её подзадачи. Если срок метода joinUntil(Instant) истечёт до того, как подзадачи завершатся или будет вызван shutdown(), он выбросит исключение, и оператор try-with-resources снова завершит работу области.
Когда join() завершается успешно, каждая из подзадач либо завершилась успешно, либо завершилась с ошибкой, либо была отменена, поскольку работа области была завершена.
После join владелец области обрабатывает подзадачи, завершившиеся с ошибкой, и результаты успешно завершившихся подзадач; обычно это делает политика завершения (см. ниже). Результат успешно завершившейся задачи можно получить методом Subtask.get(). Метод get() никогда не блокируется; он выбрасывает IllegalStateException, если его по ошибке вызвали до join или если подзадача не завершилась успешно.
Подзадачи, порождённые в области, наследуют привязки ScopedValue (JEP 446). Если владелец области читает значение из привязанного ScopedValue, то каждая подзадача прочитает то же самое значение.
Если владелец области сам является подзадачей существующей области, т. е. был создан как порождённая подзадача, то эта область становится родительской для новой области. Таким образом, области и подзадачи образуют дерево.
Структурное использование StructuredTaskScope обеспечивается во время выполнения. Например, попытки вызвать fork(Callable) из потока, который не входит в древовидную иерархию области, — т. е. не является владельцем, подзадачей или подзадачей во вложенных областях (под-подзадачей), — завершатся исключением. Использование области вне блока try-with-resources и возврат без вызова close() или без соблюдения правильной вложенности вызовов close() могут привести к тому, что методы области выбросят StructureViolationException.
StructuredTaskScope накладывает структуру и порядок на конкурентные операции. Поэтому он не реализует интерфейсы ExecutorService и Executor, поскольку экземпляры этих интерфейсов обычно используются неструктурированно (см. ниже). Тем не менее код, который использует ExecutorService, но выиграл бы от структуры, несложно перевести на StructuredTaskScope.
На практике в большинстве случаев использования StructuredTaskScope класс StructuredTaskScope будет применяться не напрямую, а через один из двух подклассов, описанных в следующем разделе, которые реализуют политики завершения. В других сценариях пользователи, вероятно, будут писать собственные подклассы для реализации своих политик завершения.
Политики завершения
При работе с конкурентными подзадачами часто используют шаблоны с коротким замыканием, чтобы не выполнять лишнюю работу. Иногда имеет смысл, например, отменить все подзадачи, если одна из них завершилась с ошибкой (т. е. invoke all), или, наоборот, если одна из них завершилась успешно (т. е. invoke any). Два подкласса StructuredTaskScope, ShutdownOnFailure и ShutdownOnSuccess, поддерживают эти шаблоны с помощью политик, которые завершают работу области, когда первая подзадача завершается с ошибкой или успешно соответственно.
Кроме того, политики завершения предоставляют централизованные методы для обработки исключений и, возможно, успешных результатов. Это соответствует духу Structured Concurrency, согласно которому вся область рассматривается как единое целое.
Вот StructuredTaskScope с политикой завершения при ошибке (она использована и в примере handle() выше), которая конкурентно выполняет набор задач и завершается с ошибкой, если с ошибкой завершается любая из них:
<T> List<T> runAll(List<Callable<T>> tasks)
throws InterruptedException, ExecutionException {
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
List<? extends Supplier<T>> suppliers = tasks.stream().map(scope::fork).toList();
scope.join()
.throwIfFailed(); // Propagate exception if any subtask fails
// Here, all tasks have succeeded, so compose their results
return suppliers.stream().map(Supplier::get).toList();
}
}
Вот StructuredTaskScope с политикой завершения при успехе, которая возвращает результат первой успешно завершившейся подзадачи:
<T> T race(List<Callable<T>> tasks, Instant deadline)
throws InterruptedException, ExecutionException, TimeoutException {
try (var scope = new StructuredTaskScope.ShutdownOnSuccess<T>()) {
for (var task : tasks) {
scope.fork(task);
}
return scope.joinUntil(deadline)
.result(); // Throws if none of the subtasks completed successfully
}
}
Как только одна подзадача завершается успешно, эта область автоматически завершает работу и отменяет незавершённые подзадачи. Задача завершается с ошибкой, если с ошибкой завершаются все подзадачи или если истекает заданный срок. Этот шаблон может быть полезен, например, в серверных приложениях, которым нужен результат от любого из набора дублирующих сервисов.
Эти две политики завершения предоставляются в готовом виде, но разработчики могут создавать собственные политики, абстрагирующие другие шаблоны (см. ниже).
Обработка результатов
После join и централизованной обработки исключений политикой завершения (например, с помощью ShutdownOnFailure::throwIfFailed) владелец области может обработать результаты подзадач с помощью объектов Subtask, возвращённых вызовами fork(...), если их не обрабатывает политика (например, ShutdownOnSuccess::result()).
Как правило, единственный метод Subtask, который будет вызывать владелец области, — метод get(). Все остальные методы Subtask обычно используются только в реализации метода handleComplete(...) собственных политик завершения (см. ниже). Более того, мы рекомендуем объявлять переменные, ссылающиеся на Subtask, возвращённый fork(...), с типом, например, Supplier<String>, а не Subtask<String> (если, конечно, вы не решите использовать var). Если политика завершения сама обрабатывает результаты подзадач — как в случае ShutdownOnSuccess, — то объектов Subtask, возвращаемых fork(...), следует вообще избегать, а метод fork(...) рассматривать так, как будто он возвращает void. Подзадачи должны возвращать в качестве результата всю информацию, которую владелец области должен обработать после централизованной обработки исключений политикой.
Если владелец области вместо использования политики завершения сам обрабатывает исключения подзадач, чтобы получить составной результат, то исключения можно возвращать из подзадач как значения. Например, вот метод, который выполняет список задач параллельно и возвращает список завершённых Future, содержащих соответствующий успешный или исключительный результат каждой задачи:
<T> List<Future<T>> executeAll(List<Callable<T>> tasks)
throws InterruptedException {
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
List<? extends Supplier<Future<T>>> futures = tasks.stream()
.map(task -> asFuture(task))
.map(scope::fork)
.toList();
scope.join();
return futures.stream().map(Supplier::get).toList();
}
}
static <T> Callable<Future<T>> asFuture(Callable<T> task) {
return () -> {
try {
return CompletableFuture.completedFuture(task.call());
} catch (Exception ex) {
return CompletableFuture.failedFuture(ex);
}
};
}
Собственные политики завершения
StructuredTaskScope можно расширить и переопределить его защищённый метод handleComplete(...), чтобы реализовать политики, отличные от политик ShutdownOnSuccess и ShutdownOnFailure. Подкласс может, например,
- собирать результаты успешно завершившихся подзадач и игнорировать подзадачи, завершившиеся с ошибкой,
- собирать исключения, когда подзадачи завершаются с ошибкой, или
- вызывать метод
shutdown(), чтобы завершить работу и пробудитьjoin(), когда возникает некоторое условие.
Когда подзадача завершается, даже после вызова shutdown(), об этом сообщается методу handleComplete(...) в виде Subtask:
public sealed interface Subtask<T> extends Supplier<T> {
enum State { SUCCESS, FAILED, UNAVAILABLE }
State state();
Callable<? extends T> task();
T get();
Throwable exception();
}
Метод handleComplete(...) вызывается для подзадач, которые завершились успешно (состояние SUCCESS) или неуспешно (состояние FAILED) до вызова shutdown(). Метод get() можно вызывать, только если подзадача находится в состоянии SUCCESS, а метод exception() — только если подзадача находится в состоянии FAILED; вызов get() или exception() в других ситуациях приведёт к тому, что они выбросят IllegalStateException. Состояние UNAVAILABLE означает одно из следующего: (1) подзадача порождена, но ещё не завершилась; (2) подзадача завершилась после завершения работы области или (3) подзадача порождена после завершения работы области и поэтому не запускалась. Метод handleComplete(...) никогда не вызывается для подзадачи в состоянии UNAVAILABLE.
Подкласс обычно определяет методы, которые делают результаты, состояние или другой итог доступными коду, выполняющемуся после возврата из метода join(). Подкласс, который собирает результаты и игнорирует подзадачи, завершившиеся с ошибкой, может определить метод, возвращающий коллекцию результатов. Подкласс, реализующий политику завершения работы при ошибке подзадачи, может определить метод для получения исключения первой подзадачи, завершившейся с ошибкой.
Вот пример подкласса StructuredTaskScope, который собирает результаты успешно завершившихся подзадач. Он определяет метод results(), с помощью которого основная задача получает результаты.
class MyScope<T> extends StructuredTaskScope<T> {
private final Queue<T> results = new ConcurrentLinkedQueue<>();
MyScope() { super(null, Thread.ofVirtual().factory()); }
@Override
protected void handleComplete(Subtask<? extends T> subtask) {
if (subtask.state() == Subtask.State.SUCCESS)
results.add(subtask.get());
}
@Override
public MyScope<T> join() throws InterruptedException {
super.join();
return this;
}
// Returns a stream of results from the subtasks that completed successfully
public Stream<T> results() {
super.ensureOwnerAndJoined();
return results.stream();
}
}
Эту собственную политику можно использовать так:
<T> List<T> allSuccessful(List<Callable<T>> tasks) throws InterruptedException {
try (var scope = new MyScope<T>()) {
for (var task : tasks) scope.fork(task);
return scope.join()
.results().toList();
}
}
Сценарии fan-in
Примеры выше касались сценариев fan-out, в которых нужно управлять несколькими конкурентными исходящими операциями ввода-вывода. StructuredTaskScope полезен и в сценариях fan-in, в которых нужно управлять несколькими конкурентными входящими операциями ввода-вывода. В таких сценариях мы обычно создаём заранее неизвестное число подзадач в ответ на входящие запросы.
Вот пример сервера, который порождает подзадачи для обработки входящих соединений внутри StructuredTaskScope:
void serve(ServerSocket serverSocket) throws IOException, InterruptedException {
try (var scope = new StructuredTaskScope<Void>()) {
try {
while (true) {
var socket = serverSocket.accept();
scope.fork(() -> handle(socket));
}
} finally {
// If there's been an error or we're interrupted, we stop accepting
scope.shutdown(); // Close all active connections
scope.join();
}
}
}
С точки зрения конкурентности этот сценарий отличается не столько направлением запросов, сколько длительностью и количеством задач. Здесь, в отличие от предыдущих примеров, время работы владельца области не ограничено — он остановится, только когда будет прерван. Число подзадач также неизвестно, поскольку они порождаются динамически в ответ на внешние события.
Все подзадачи обработки соединений создаются внутри области, поэтому их назначение легко увидеть в дампе потоков, где они будут показаны как дочерние для владельца области. Также легко завершить работу всего сервиса как единого целого.
Наблюдаемость
Мы расширяем новый формат дампа потоков в JSON, добавленный в JEP 444, чтобы показать группировку потоков в иерархию, которую задаёт StructuredTaskScope:
$ jcmd <pid> Thread.dump_to_file -format=json <file>
JSON-объект каждой области видимости содержит массив потоков, порождённых в этой области, вместе с их трассировками стека. Поток-владелец области обычно заблокирован в методе join в ожидании завершения подзадач. Дамп потоков показывает древовидную иерархию, которую задаёт Structured Concurrency, поэтому по нему легко увидеть, что делают потоки подзадач. JSON-объект области также содержит ссылку на родительскую область, так что структуру программы можно восстановить по дампу.
API com.sun.management.HotSpotDiagnosticsMXBean также можно использовать для создания таких дампов потоков — напрямую или косвенно, через платформенный MBeanServer и локальный или удалённый инструмент JMX.
Почему fork(...) не возвращает Future?
Когда API StructuredTaskScope был в статусе Incubator, метод fork(...) возвращал Future. Так API казался привычнее: 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существующим методам, принимающим этот тип, почти наверняка приводила бы к исключениям в большинстве ситуаций.