← назад к разделу

Классическая боль параллельного кода: метод запустил три подзадачи, одна упала — а две другие продолжают жечь ресурсы, и отменить их некому. Или наоборот: метод вернулся, а «осиротевшая» подзадача всё ещё жива и пишет в лог через минуту после ответа клиенту. Структурированная конкурентность решает это одним принципом: время жизни подзадач ограничено блоком кода, который их породил. В Java этот принцип воплощает StructuredTaskScope.

без области метод handle() граница метода fetchUser ✗ fetchOrders метод вернулся жива и пишет в лог в области scope граница области fetchUser ✗ fetchOrdersотменена выход: живых нет

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

Обязательно

Проблема: неструктурированный fan-out

Соберём профиль пользователя из двух источников на ExecutorService:

Future<User> userF = executor.submit(() -> fetchUser(id));
Future<List<Order>> ordersF = executor.submit(() -> fetchOrders(id));
User user = userF.get();          // а если fetchOrders уже упал?
List<Order> orders = ordersF.get(); // узнаем только здесь

Здесь три хронические болезни. Утечка работы: если fetchUser бросил исключение, get() пробрасывает его наверх — но fetchOrders никто не отменил, он доработает впустую. Потерянные ошибки: упади fetchOrders первым — мы не узнаем об этом, пока не дойдём до его get(). Оборванная связь: по стеку вызовов не видно, что эти задачи — дети нашего метода; поток-родитель и подзадачи живут каждый своей жизнью, и отмена родителя детей не касается.

CompletableFuture с allOf часть проблем сглаживает, но отмену «братьев» при первой ошибке и жёсткую привязку к породившему блоку всё равно приходится собирать руками.

Идея: конкурентность с блочной структурой

Аналогия: когда-то поток управления прыгал по goto куда угодно; структурное программирование ввело блоки — вход сверху, выход снизу, — и код стало возможно читать. Структурированная конкурентность делает то же с потоками: подзадачи живут строго внутри области (scope), и выйти из области нельзя, пока все подзадачи не завершены — успехом, ошибкой или отменой.

Из этого автоматически следуют гарантии: ни одна подзадача не переживёт родителя; ошибки собираются в одном месте; отмена распространяется сверху вниз.

Такую границу можно выстроить и на обычном API, без preview-флагов, — правда, руками:

живой пример

import java.util.concurrent.*;

public class ScopeDemo {
    static String user() { throw new IllegalStateException("нет связи"); }

    static String orders() throws InterruptedException {
        Thread.sleep(200);
        System.out.println("  fetchOrders доработал впустую");
        return "orders";
    }

    public static void main(String[] args) throws Exception {
        ExecutorService loose = Executors.newVirtualThreadPerTaskExecutor();
        Future<String> f = loose.submit(ScopeDemo::user);
        loose.submit(ScopeDemo::orders);
        try { f.get(); } catch (ExecutionException e) { System.out.println("без области: вышли с ошибкой"); }
        Thread.sleep(400);

        try (ExecutorService scope = Executors.newVirtualThreadPerTaskExecutor()) {
            var done = new ExecutorCompletionService<String>(scope);
            Future<String> a = done.submit(ScopeDemo::user);
            Future<String> b = done.submit(ScopeDemo::orders);
            try { done.take().get(); } catch (ExecutionException e) {
                System.out.println("в области: ошибка — отменяем остальных");
                a.cancel(true);
                b.cancel(true);
            }
        }
        System.out.println("вышли из блока: живых подзадач нет");
    }
}
Запустить

Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →

Первая половина оставила подзадачу дописывать работу в пустоту. Вторая дождалась всех и на первой же ошибке отменила остальных — но посмотрите, чего это стоило. close() у ExecutorService подзадачи только дожидается, отменять он их не умеет, поэтому cancel(true) пришлось расставить вручную и не забыть ни одной. Граница получилась та же, а держится она на внимательности автора. StructuredTaskScope даёт её конструкцией: отменять там руками нечего.

StructuredTaskScope: API Java 25

Response handle(long id) throws InterruptedException {
    try (var scope = StructuredTaskScope.<Object>open()) {
        Subtask<User> user = scope.fork(() -> fetchUser(id));
        Subtask<List<Order>> orders = scope.fork(() -> fetchOrders(id));

        scope.join();   // ждём всех; при ошибке любой — отмена остальных и исключение

        return new Response(user.get(), orders.get());
    } // выход из try гарантирует: живых подзадач не осталось
}

Разбор:

  • open() создаёт область; try-with-resources гарантирует уборку — как бы блок ни завершился, все подзадачи будут дождаты или отменены. Параметр типа в <Object>open() — это общий тип результатов подзадач области. Здесь подзадачи возвращают разное (User и List<Order>), общий предок у них только Object — отсюда и такая запись. Когда все подзадачи возвращают одно и то же, пишут конкретный тип: StructuredTaskScope.<User>open().
  • fork(...) запускает подзадачу в собственном виртуальном потоке и возвращает Subtask — «билет» за результатом.
  • join() — единственная точка ожидания. Политика по умолчанию: все успешны — едем дальше; любая упала — остальные отменяются, и join() бросает FailedException с причиной внутри.
  • subtask.get() после join — просто забрать готовое значение, никаких блокировок.

Что ловить на стороне вызывающего

join() объявляет два исключения, и различать их приходится в каждом реальном использовании.

try (var scope = StructuredTaskScope.<Object>open(Joiner.allSuccessfulOrThrow(),
        cf -> cf.withTimeout(Duration.ofSeconds(2)))) {
    var user = scope.fork(() -> fetchUser(id));
    var orders = scope.fork(() -> fetchOrders(id));
    scope.join();
    return new Response(user.get(), orders.get());
} catch (StructuredTaskScope.FailedException e) {
    Throwable cause = e.getCause();          // настоящая причина отказа подзадачи
    throw new ProfileUnavailableException(cause);
} catch (TimeoutException e) {
    throw new ProfileTimeoutException();     // не уложились в дедлайн области
} catch (InterruptedException e) {
    Thread.currentThread().interrupt();      // прервали нас самих
    throw new ProfileInterruptedException();
}

Три исхода — три разных решения. FailedException означает «подзадача упала», и причина лежит в getCause(), как у ExecutionException (лишний слой обёртки — плата за то, что ошибка приехала из другого потока). TimeoutException означает «никто не упал, просто не успели»; подзадачи к этому моменту уже отменены областью. InterruptedException объявлен потому, что join() — обычное блокирующее ожидание: прервать могут и сам поток-владелец, например при остановке приложения; тогда область отменяет подзадачи и выходит, а флаг прерывания надо восстановить.

Когда политика — Joiner.awaitAll(), исключений не будет вовсе: область честно дождётся всех, и разбирать исходы придётся самому, по каждой подзадаче:

scope.join();
for (Subtask<Price> t : subtasks) {
    switch (t.state()) {
        case SUCCESS -> results.add(t.get());
        case FAILED -> log.warn("источник отказал", t.exception());
        case UNAVAILABLE -> log.warn("подзадача не запускалась или отменена");
    }
}

state() даёт один из трёх исходов, get() законен только у SUCCESS (иначе IllegalStateException), exception() — только у FAILED. Это и есть способ собрать частичный результат: три источника из пяти ответили, и этого достаточно.

Утечка работы, потерянные ошибки и осиротевшие задачи исчезают не дисциплиной, а конструкцией: их больше негде оформить.

allSuccessfulOrThrow одна упала остальных отменить FailedException anySuccessfulResultOrThrow первая успела остальных отменить результат из join awaitAll ждём всех никого не отменяем разбор по state()

Три политики ожидания рядом: первые две при первом же исходе отменяют остальные подзадачи, а awaitAll доводит до конца всех и оставляет разбор исходов вызывающему.

Статус и практика

StructuredTaskScope — preview-возможность (JEP 505, пятое preview в Java 25): включается флагом --enable-preview, и API уже менялся — в Java 21–24 вместо open() были классы ShutdownOnFailure/ShutdownOnSuccess (они ещё встречаются в кодовых базах тех лет). Идея стабильна, сигнатуры — сверяйте со своим JDK.

Если ваш базовый JDK — 21, ни один фрагмент выше не скомпилируется: там та же задача записывается иначе, а fork возвращает не Subtask, а Future.

// Java 21, --enable-preview
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
    Future<User> user = scope.fork(() -> fetchUser(id));
    Future<List<Order>> orders = scope.fork(() -> fetchOrders(id));

    scope.join();          // дождаться всех
    scope.throwIfFailed();  // упала любая — пробросить её ошибку

    return new Response(user.resultNow(), orders.resultNow());
}

Смысл тот же: область, ожидание в одной точке, отмена остальных при первой ошибке. Разошлись имена, а не идея.

На сегодня: в библиотеках и коде, который живёт годами, — осторожно; в сервисах на свежем JDK — хороший кандидат везде, где пишут CompletableFuture.allOf или пачку submit/get. С одной оговоркой: возможность предварительная и требует --enable-preview, а этот флаг в прод пускают далеко не везде. Пока он не разрешён, это тема для следующего обновления, а не для сегодняшней переделки.

Дополнительно: при первом чтении можно пропустить

Глубже: Joiner — политики ожиданиярасширенное

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

// Первый успешный результат — остальных отменить (гонка зеркал/реплик):
try (var scope = StructuredTaskScope.open(
        Joiner.<String>anySuccessfulResultOrThrow())) {
    scope.fork(() -> fetchFrom(mirrorA));
    scope.fork(() -> fetchFrom(mirrorB));
    return scope.join();   // сам возвращает результат победителя
}

// Дождаться ВСЕХ и разобрать исходы вручную:
try (var scope = StructuredTaskScope.open(Joiner.<Widget>awaitAll())) { ... }

Плюс конфигурация области — например, общий дедлайн:

StructuredTaskScope.open(Joiner.allSuccessfulOrThrow(),
        cf -> cf.withTimeout(Duration.ofSeconds(2)));
// не успели — все подзадачи отменяются, join() бросает TimeoutException

Таймаут на группу целиком — то, что с фьючерами собиралось из костылей, здесь одна строка.

0 мс fork: цены 0 мс fork: остатки 0 мс fork: отзывы 900 мс цены готовы 1500 мс остатки готовы 2000 мс дедлайн области 2000 мс отзывы отменены, TimeoutException

Дедлайн стоит на группе целиком, а не на каждой подзадаче: на отсечке область сама отменяет тех, кто не успел, и вызывающий получает TimeoutException.

Ещё деталь: значения ScopedValue наследуются подзадачами автоматически — контекст запроса (пользователь, trace id) течёт в fan-out без ручного проброса.

Глубже: вложенные области и отменарасширенное

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

try (var outer = StructuredTaskScope.<Report>open()) {
    outer.fork(() -> {
        try (var inner = StructuredTaskScope.<Row>open()) {   // подзадача открыла свою область
            inner.fork(() -> loadRows(a));
            inner.fork(() -> loadRows(b));
            inner.join();
            return buildReport(...);
        }
    });
    outer.join();
}

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

И структура дерева видна снаружи: дамп через jcmd <pid> Thread.dump_to_file -format=json показывает виртуальные потоки сгруппированными по областям, с именем метода, который область открыл, — то есть по дампу читается, какая подзадача чьей дочерью является. С россыпью submit такого не увидеть: там все потоки пула равноправны и связи между ними нет.

И честная граница применимости: область требует, чтобы подзадачи не пережили породивший блок. Значит, всё, что должно продолжаться после ответа клиенту, — фоновая отправка письма, прогрев кэша, дозапись метрик — областью не оформляют: выход из блока их отменит. Для такой работы остаются CompletableFuture и обычный исполнитель, которым владеет приложение, а не запрос.

Коротко

  • Проблема неструктурированного fan-out: утечки работы, потерянные ошибки, осиротевшие подзадачи, отмена руками.
  • Принцип: подзадачи живут внутри области; выход из области = все подзадачи завершены. Это структурное программирование, применённое к потокам.
  • API: open() → fork() → join() → get(); try-with-resources гарантирует уборку; ошибка одной подзадачи отменяет остальных.
  • Joiner задаёт политику: все успешны / первый успешный / дождаться всех; конфигурация добавляет дедлайн на группу.
  • ScopedValue наследуется подзадачами — контекст течёт сам. Статус: preview (JEP 505), API в 21–24 отличался.
  • join() бросает FailedException (причина в getCause()), TimeoutException при дедлайне области и InterruptedException, если прервали владельца; при awaitAll исходы разбирают по state(), get() и exception() каждой подзадачи.
  • Области вкладываются, отмена идёт вниз по дереву, и дерево видно в дампе jcmd Thread.dump_to_file; работа, которая должна пережить запрос, областью не оформляется — там остаются CompletableFuture и свой исполнитель.

Что почитать дальше

  • Виртуальные потоки — на чём работает fork и почему «поток на подзадачу» снова дёшев.
  • CompletableFuture — асинхронные цепочки: когда результат нужен «потом», а не «здесь».
  • Пулы потоков: ExecutorService — классический слой, который scope постепенно вытесняет из прикладного кода.