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

JEP 428: Structured Concurrency (Incubator)

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

AuthorsAlan Bateman, Ron Pressler
ОтветственныйAlan Bateman
ТипFeature
ОбластьJDK
СтатусClosed / Delivered
Выпуск19
Компонентcore-libs
Обсуждениеloom dash dev at openjdk dot java dot net
РецензентыAlex Buckley, Brian Goetz
Создан2021/11/15 15:01
Обновлён2023/06/08 16:05
Задача8277129

Аннотация

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

Цели

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

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

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

  • Мы не стремимся заменить какие-либо конструкции конкурентности из пакета 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. На проектирование обработки ошибок в Structured Concurrency повлияли идеи из других языков, например иерархические супервизоры Erlang.

В Structured Concurrency подзадачи работают от имени задачи. Задача ожидает результатов подзадач и отслеживает их сбои. Как и в приёмах структурного программирования для кода в одном потоке, сила Structured Concurrency для нескольких потоков основана на двух идеях: (1) чётко определённые точки входа и выхода для потока выполнения через блок кода и (2) строгая вложенность времени жизни операций, отражающая их синтаксическую вложенность в коде.

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

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

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

Описание

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

Вот пример 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() возвращает управление, все ответвлённые подзадачи либо завершены (успешно или нет), либо отменены. Их результаты или исключения можно получить без дополнительной блокировки через методы resultNow() или exceptionNow() их future. (Эти методы выбрасывают 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.

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

Зависимости