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

JEP 437: Structured Concurrency (Second Incubator)

Structured Concurrency (структурированная конкурентность), вторая версия Incubator (инкубационный модуль)

AuthorsAlan Bateman, Ron Pressler
ОтветственныйAlan Bateman
ТипFeature
ОбластьJDK
СтатусClosed / Delivered
Выпуск20
Компонентcore-libs
Обсуждениеloom dash dev at openjdk dot org
Связан сJEP 453: Structured Concurrency (Preview)
РецензентыAlex Buckley
ОдобренBrian Goetz
Создан2022/10/28 12:41
Обновлён2023/06/08 16:58
Задача8296037

Аннотация

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

История

Structured Concurrency была предложена в JEP 428 и вошла в JDK 19 как API в статусе Incubator. Этот JEP предлагает повторно выпустить API в статусе Incubator, без изменений, в JDK 20, чтобы получить больше отзывов и больше опыта работы с этой возможностью.

Единственное изменение в повторно выпущенном API: StructuredTaskScope обновлён и теперь поддерживает наследование Scoped Values (значения с ограниченной областью видимости) (JEP 429) потоками, созданными в области задачи. Так упрощается совместное использование неизменяемых данных несколькими потоками.

Цели

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

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

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

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

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

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

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

Мотивация

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Structured Concurrency

Structured Concurrency — это подход к многопоточному программированию, который сохраняет читаемость, сопровождаемость и наблюдаемость однопоточного кода. Он воплощает принцип:

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

Термин «Structured Concurrency» ввёл Martin Sústrik, а популяризировал Nathaniel J. Smith. Идеи из других языков, например иерархические супервизоры 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 гарантирует их правильное и надёжное согласование и позволяет инструментам наблюдения отображать потоки так, как их понимает разработчик. Наличие в JDK API для Structured Concurrency улучшило бы сопровождаемость, надёжность и наблюдаемость серверных приложений.

Описание

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

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

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

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

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

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

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

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

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

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

Как и ExecutorService.submit(...), метод StructuredTaskScope.fork(...) принимает Callable и возвращает Future. Однако, в отличие от ExecutorService, возвращённый future не предназначен для присоединения через его метод get() или отмены через его метод cancel(). Все ответвления в области, напротив, предполагается присоединять или отменять как единое целое. Два новых метода Future, resultNow() и exceptionNow(), предназначены для использования после завершения подзадач, например после вызова scope.join().

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

Общая схема работы кода, использующего StructuredTaskScope, выглядит так:

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

  2. Ответвите в области конкурентные подзадачи.

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

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

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

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

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

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

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

Когда join() возвращает управление, все ответвления либо завершились (успешно или нет), либо отменены. Их результаты или исключения можно получить без какой-либо дополнительной блокировки через методы их future resultNow() или exceptionNow(). (Эти методы выбрасывают IllegalStateException, если вызваны до завершения future.)

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

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

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

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

StructuredTaskScope находится в Incubator-модуле, который по умолчанию исключён

В примерах выше используется API StructuredTaskScope, поэтому, чтобы запустить их на JDK XX, нужно добавить модуль jdk.incubator.concurrent, а также включить возможности в статусе Preview (предварительная версия), чтобы включить Virtual Threads:

  • Скомпилируйте программу с javac --release XX --enable-preview --add-modules jdk.incubator.concurrent Main.java и запускайте её с java --enable-preview --add-modules jdk.incubator.concurrent Main; или

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

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

Политики завершения работы

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

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

<T> List<T> runAll(List<Callable<T>> tasks) throws Throwable {
    try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
        List<Future<T>> futures = tasks.stream().map(scope::fork).toList();
        scope.join();
        scope.throwIfFailed(e -> e);  // Propagate exception as-is if any fork fails
        // Here, all tasks have succeeded, so compose their results
        return futures.stream().map(Future::resultNow).toList();
    }
}

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

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

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

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

Сценарии 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 425, чтобы показать группировку потоков в иерархию, которую задаёт StructuredTaskScope:

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

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

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

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

Зависимости