Введение
Эта статья, представленная как улучшение API параллелизма Java 8, представляет собой введение в возможности и варианты использования класса CompletableFuture. В то же время в CompletableFuture в Java 9 есть некоторые улучшения, которые будут объяснены позже.
Будущий расчет
Асинхронными вычислениями на фьючерсах трудно манипулировать, и обычно мы хотим думать о любой логике вычислений как о последовательности шагов. Но в случае асинхронных вычислений методы, представленные в виде обратных вызовов, обычно разбросаны по коду или глубоко вложены друг в друга. Но все может усложниться, когда нам нужно разобраться с возможной ошибкой на одном из шагов.
Интерфейс Futrue был добавлен в Java 5 как асинхронные вычисления, но у него нет возможности объединять вычисления или обрабатывать возможные ошибки.
В Java 8 был представлен класс CompletableFuture. Наряду с интерфейсом Future он также реализует интерфейс CompletionStage. Этот интерфейс определяет контракты, которые можно комбинировать с другими фьючерсами для формирования контрактов асинхронных вычислений.
CompletableFuture — это одновременно композиция и фреймворк с примерно 50 различными композициями, комбинациями, выполнением асинхронных вычислительных шагов и обработкой ошибок.
Такой большой API может быть ошеломляющим, и некоторые важные из них будут выделены ниже.
Используйте CompletableFuture как будущую реализацию
Во-первых, класс CompletableFuture реализует интерфейс Future, поэтому его можно использовать как реализацию Future, но для этого требуется дополнительная логика реализации завершения.
Например, вы можете использовать конструктор без параметров для создания экземпляра этого класса, а затем использоватьcompleteметод завершен. Потребители могут использовать метод get для блокировки текущего потока до тех пор, покаget()результат.
В приведенном ниже примере у нас есть метод, который создает экземпляр CompletableFuture, затем выполняет вычисления в другом потоке и немедленно возвращает Future.
После завершения вычисления метод завершает Future, предоставляя результат методу complete:
public Future<String> calculateAsync() throws InterruptedException {
CompletableFuture<String> completableFuture
= new CompletableFuture<>();
Executors.newCachedThreadPool().submit(() -> {
Thread.sleep(500);
completableFuture.complete("Hello");
return null;
});
return completableFuture;
}
Для разделения вычислений используемExecutorAPI, это создание и доработкаМетоды CompletableFutureМожет использоваться с любым параллельным пакетом (включая необработанные потоки).
Пожалуйста, обрати внимание,ДолженcalculateAsyncметод возвращаетFutureпример.
Мы просто вызываем метод, получаемFutureэкземпляр и вызвать его, когда мы будем готовы заблокировать результатgetметод.
Также обратите внимание, чтоgetметод выдает некоторое проверенное исключение, т.е.ExecutionException(для инкапсуляции исключений, возникших во время вычислений) иInterruptedException(исключение, указывающее, что поток, выполняющий метод, был прерван):
Future<String> completableFuture = calculateAsync();
// ...
String result = completableFuture.get();
assertEquals("Hello", result);
Если вы уже знаете результат вычисления, вы также можете вернуть результат синхронным образом.
Future<String> completableFuture =
CompletableFuture.completedFuture("Hello");
// ...
String result = completableFuture.get();
assertEquals("Hello", result);
Как и в некоторых сценариях, вы можете отменить выполнение будущей задачи.
Допустим, мы не находим результата и решаем полностью отменить асинхронное выполнение. Это можно сделать с помощью метода отмены Future. Этот методmayInterruptIfRunning, но в случае CompletableFuture это не имеет никакого эффекта, поскольку прерывания не используются для управления обработкой CompletableFuture.
Вот модифицированная версия асинхронного метода:
public Future<String> calculateAsyncWithCancellation() throws InterruptedException {
CompletableFuture<String> completableFuture = new CompletableFuture<>();
Executors.newCachedThreadPool().submit(() -> {
Thread.sleep(500);
completableFuture.cancel(false);
return null;
});
return completableFuture;
}
Когда мы используем метод Future.get() для блокировки результата,cancel()означает отменить выполнение, оно вызовет исключение CancellationException:
Future<String> future = calculateAsyncWithCancellation();
future.get(); // CancellationException
Введение в API
описание статического метода
Приведенный выше код очень прост, вот несколькоstaticметоды, использующие задачи для создания экземпляра CompletableFuture.
CompletableFuture.runAsync(Runnable runnable);
CompletableFuture.runAsync(Runnable runnable, Executor executor);
CompletableFuture.supplyAsync(Supplier<U> supplier);
CompletableFuture.supplyAsync(Supplier<U> supplier, Executor executor)
- Метод runAsync получает экземпляр Runnable, но не имеет возвращаемого значения.
- Метод SupplyAsync представляет собой функциональный интерфейс JDK8 без параметров и возвращает результат.
- Эти два метода являются обновлением исполнителя, что означает, что задача выполняется в указанном пуле потоков.Если не указано, задача обычно выполняется в пуле потоков ForkJoinPool.commonPool().
SupplyAsync() использовать
статический методrunAsyncиsupplyAsyncПозволяет нам создавать экземпляры CompletableFuture из типов функций Runnable и Supplier соответственно.
Интерфейс Runnable — это старый интерфейс, используемый потоками, который не позволяет возвращать значения.
Интерфейс Supplier — это универсальный функциональный интерфейс, который не принимает параметров и возвращает один метод параметризованного типа.
Это позволяет предоставить экземпляр Supplier в виде лямбда-выражения, которое выполняет вычисления и возвращает результат:
CompletableFuture<String> future
= CompletableFuture.supplyAsync(() -> "Hello");
// ...
assertEquals("Hello", future.get());
thenRun() использовать
В двух задачах задача A, задача B, если вам не нужно ни значение задачи A, ни ссылка в задаче B, вы можете передать Runnable lambda вthenRun()метод. В приведенном ниже примере после вызова метода future.get() мы просто выводим в консоль строку:
шаблон
CompletableFuture.runAsync(() -> {}).thenRun(() -> {});
CompletableFuture.supplyAsync(() -> "resultA").thenRun(() -> {});
- В первой строке используется
thenRun(Runnable runnable), задача A заканчивает выполнение B, и B не нужен результат A. - Во второй строке используется
thenRun(Runnable runnable), после выполнения задачи A и выполнения B она вернетresultA, но B не нужен результат A .
настоящий бой
CompletableFuture<String> completableFuture
= CompletableFuture.supplyAsync(() -> "Hello");
CompletableFuture<Void> future = completableFuture
.thenRun(() -> System.out.println("Computation finished."));
future.get();
thenAccept() использовать
В двух задачах задача A, задача B, если вам не нужно иметь возвращаемое значение в будущем, вы можете использоватьthenAcceptПолучатель метода передает ему результат вычисления. Последний вызов future.get() возвращает экземпляр типа Void.
шаблон
CompletableFuture.runAsync(() -> {}).thenAccept(resultA -> {});
CompletableFuture.supplyAsync(() -> "resultA").thenAccept(resultA -> {});
- В первой строке
runAsyncВозвращаемого значения не будет, второй способthenAccept, полученное значение resultA равно null, и задача B не вернет результат - Во второй строке
supplyAsyncЕсть возвращаемое значение, и задача B не вернет результат.
настоящий бой
CompletableFuture<String> completableFuture
= CompletableFuture.supplyAsync(() -> "Hello");
CompletableFuture<Void> future = completableFuture
.thenAccept(s -> System.out.println("Computation returned: " + s));
future.get();
тогдаПрименить () использовать
В двух задачах задача A, задача B, задача B хотят, чтобы результат был рассчитан задачей A, вы можете использоватьthenApplyметод для приема экземпляра функции, использования его для обработки результата и возврата возвращаемого значения функции Future:
шаблон
CompletableFuture.runAsync(() -> {}).thenApply(resultA -> "resultB");
CompletableFuture.supplyAsync(() -> "resultA").thenApply(resultA -> resultA + " resultB");
- Во второй строке используется thenApply(Function fn), задача A выполняет B, B нужен результат A, а задача B имеет возвращаемое значение.
настоящий бой
CompletableFuture<String> completableFuture
= CompletableFuture.supplyAsync(() -> "Hello");
CompletableFuture<String> future = completableFuture
.thenApply(s -> s + " World");
assertEquals("Hello World", future.get());
Конечно, в случае нескольких задач, если за задачей B следует задача C, продолжайте вызывать .thenXxx() вниз.
thenCompose() использовать
Далее будет очень интересный шаблон оформления;
Наилучший вариант использования API CompletableFuture — это возможность создавать экземпляры CompletableFuture в виде серии вычислительных шагов.
Результатом этой композиции является само CompletableFuture, допускающее дальнейшее продолжение композиции. Этот подход повсеместно распространен в функциональных языках и часто упоминается какmonadic设计模式.
Проще говоря, Monad — это шаблон проектирования, который разбивает рабочий процесс на несколько взаимосвязанных шагов с помощью функций. Пока вы предоставляете функцию, необходимую для следующей операции, вся операция будет выполняться автоматически.
В приведенном ниже примере мы используем метод thenCompose для последовательного объединения двух фьючерсов.
Обратите внимание, что этот метод принимает функцию, которая возвращает экземпляр CompletableFuture. Аргументы этой функции являются результатами предыдущих шагов вычисления. Это позволяет нам использовать это значение в следующей лямбде CompletableFuture:
CompletableFuture<String> completableFuture
= CompletableFuture.supplyAsync(() -> "Hello")
.thenCompose(s -> CompletableFuture.supplyAsync(() -> s + " World"));
assertEquals("Hello World", completableFuture.get());
Метод thenCompose вместе с thenApply реализует комбинированное вычисление результатов. Но их внутренние формы отличаются, у них схожая идея дизайна с методами map и flatMap классов Stream и Optional, доступных в Java 8.
Оба метода получают CompletableFuture и применяют его к результату вычисления, но метод thenCompose(flatMap) получает функцию, которая возвращает другой объект CompletableFuture того же типа. Эта функциональная структура позволяет экземплярам этих классов продолжать комбинироваться для вычислений.
thenCombine()
Возьмите результат двух задач
Если вы хотите выполнить две отдельные задачи и что-то сделать с их результатами, вы можете использовать метод thenCombine из Future:
шаблон
CompletableFuture<String> cfA = CompletableFuture.supplyAsync(() -> "resultA");
CompletableFuture<String> cfB = CompletableFuture.supplyAsync(() -> "resultB");
cfA.thenAcceptBoth(cfB, (resultA, resultB) -> {});
cfA.thenCombine(cfB, (resultA, resultB) -> "result A + B");
настоящий бой
CompletableFuture<String> completableFuture
= CompletableFuture.supplyAsync(() -> "Hello")
.thenCombine(CompletableFuture.supplyAsync(
() -> " World"), (s1, s2) -> s1 + s2));
assertEquals("Hello World", completableFuture.get());
В более простом случае, когда вы хотите использовать два будущих результата, но вам не нужно возвращать какое-либо значение результата, вы можете использоватьthenAcceptBoth, что указывает на то, что последующая обработка не требует возвращаемого значения, тогда как thenCombine указывает, что возвращаемое значение требуется:
CompletableFuture future = CompletableFuture.supplyAsync(() -> "Hello")
.thenAcceptBoth(CompletableFuture.supplyAsync(() -> " World"),
(s1, s2) -> System.out.println(s1 + s2));
Разница между thenApply() и thenCompose()
В предыдущих разделах мы показали примеры функций thenApply() и thenCompose(). Оба API используют вызовы CompletableFuture, но использование этих двух API отличается.
thenApply()
Этот метод используется для обработки ранее вызванногорезультат. Однако ключевой момент, о котором следует помнить, заключается в том, что возвращаемый тип — это тип универсального преобразования, тот же самый CompletableFuture.
Итак, когда мы хотим преобразовать результат вызова CompletableFuture, эффект будет следующим:
CompletableFuture<Integer> finalResult = compute().thenApply(s-> s + 1);
thenCompose()
Метод thenCompose() похож на thenApply() в том смысле, что оба возвращают новый вычисленный результат. Однако thenCompose() принимает предыдущее Future в качестве параметра. Это напрямую сделает результат новым Future, а не вложенным Future, который мы получили в thenApply(), но используемым для соединения двух CompletableFuture для создания нового CompletableFuture:
CompletableFuture<Integer> computeAnother(Integer i){
return CompletableFuture.supplyAsync(() -> 10 + i);
}
CompletableFuture<Integer> finalResult = compute().thenCompose(this::computeAnother);
Итак, если вы хотите продолжить вложенную цепочкуCompletableFutureметод, то лучше использоватьthenCompose().
Запускайте несколько задач параллельно
Когда нам нужно выполнить несколько задач параллельно, мы обычно хотим дождаться их выполнения, а затем обработать их объединенный результат.
ДолженCompletableFuture.allOfСтатические методы позволяют дождаться завершения всех задач:
API
public static CompletableFuture<Void> allOf(CompletableFuture<?>... cfs){...}
настоящий бой
CompletableFuture<String> future1
= CompletableFuture.supplyAsync(() -> "Hello");
CompletableFuture<String> future2
= CompletableFuture.supplyAsync(() -> "Beautiful");
CompletableFuture<String> future3
= CompletableFuture.supplyAsync(() -> "World");
CompletableFuture<Void> combinedFuture
= CompletableFuture.allOf(future1, future2, future3);
// ...
combinedFuture.get();
assertTrue(future1.isDone());
assertTrue(future2.isDone());
assertTrue(future3.isDone());
Обратите внимание, что CompletableFuture.allOf() возвращает тип CompletableFuture. Ограничение этого подхода заключается в том, что он не дает исчерпывающих результатов для всех задач. Вместо этого вам придется вручную получать результаты из фьючерсов. К счастью, метод CompletableFuture.join() и API Java 8 Streams могут решить:
String combined = Stream.of(future1, future2, future3)
.map(CompletableFuture::join)
.collect(Collectors.joining(" "));
assertEquals("Hello Beautiful World", combined);
CompletableFuture предоставляет метод join(). Его функция такая же, как и у метода get(). Оба блокируют получение значений. Разница между ними в том, что join() генерирует непроверенное исключение. Это делает его доступным в качестве ссылки на метод в методе Stream.map().
Обработка исключений
Сказав это, давайте, кстати, поговорим об обработке исключений CompletableFuture. Здесь мы вводим два метода:
public CompletableFuture<T> exceptionally(Function<Throwable, ? extends T> fn);
public <U> CompletionStage<U> handle(BiFunction<? super T, Throwable, ? extends U> fn);
посмотри на код
CompletableFuture.supplyAsync(() -> "resultA")
.thenApply(resultA -> resultA + " resultB")
.thenApply(resultB -> resultB + " resultC")
.thenApply(resultC -> resultC + " resultD");
В приведенном выше коде последовательно выполняются задачи A, B, C и D. Если задача A выдает исключение (конечно, приведенный выше код не выдает исключение), то последующие задачи выполняться не будут. Если задача C выдает исключение, то задача D не может быть выполнена.
Итак, как мы обрабатываем исключения? Глядя на код ниже, мы генерируем исключение в задаче A и обрабатываем его:
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
throw new RuntimeException();
})
.exceptionally(ex -> "errorResultA")
.thenApply(resultA -> resultA + " resultB")
.thenApply(resultB -> resultB + " resultC")
.thenApply(resultC -> resultC + " resultD");
System.out.println(future.join());
В приведенном выше коде задача A выдает исключение, а затем проходит.exceptionally()Метод обрабатывает исключение и возвращает новый результат, который передается задаче Б. Итак, окончательный вывод:
errorResultA resultB resultC resultD
String name = null;
// ...
CompletableFuture<String> completableFuture
= CompletableFuture.supplyAsync(() -> {
if (name == null) {
throw new RuntimeException("Computation error!");
}
return "Hello, " + name;
})}).handle((s, t) -> s != null ? s : "Hello, Stranger!");
assertEquals("Hello, Stranger!", completableFuture.get());
Конечно, они оба могут быть нулевыми, потому что s имеет значение null, если экземпляр CompletableFuture, на который он действует, не имеет возвращаемого значения.
Асинхронный постфиксный метод
CompletableFutureБольшинство методов API в классе имеют два сAsyncДополнительные модификации суффикса. Эти методы представлены для асинхронных потоков.
нетAsyncМетоды Postfix используют вызывающий поток для запуска следующего этапа потока выполнения. БезAsyncиспользование методаForkJoinPool.commonPool()пул потоковfork / joinвыполнять вычислительные задачи. с участиемAsyncметод с использованием переходногоExecutorзадача для запуска.
Кейс прикреплен ниже, вы можете видеть, что естьthenApplyAsyncметод. Внутри программы поток заворачивается вForkJoinTaskв экземпляре. Это дополнительно распараллеливает ваши вычисления и более эффективно использует системные ресурсы.
CompletableFuture<String> completableFuture
= CompletableFuture.supplyAsync(() -> "Hello");
CompletableFuture<String> future = completableFuture
.thenApplyAsync(s -> s + " World");
assertEquals("Hello World", future.get());
JDK 9 CompletableFuture API
В Java 9 API CompletableFuture был дополнительно улучшен за счет следующих изменений:
- Добавлен новый фабричный метод
- Поддержка задержек и тайм-аутов
- Улучшена поддержка подклассов.
Представлен новый API экземпляра:
- Executor defaultExecutor()
- CompletableFuture newIncompleteFuture()
- CompletableFuture copy()
- CompletionStage minimalCompletionStage()
- CompletableFuture completeAsync(Supplier<? extends T> supplier, Executor executor)
- CompletableFuture completeAsync(Supplier<? extends T> supplier)
- CompletableFuture orTimeout(long timeout, TimeUnit unit)
- CompletableFuture completeOnTimeout(T value, long timeout, TimeUnit unit)
Есть также несколько статических служебных методов:
- Executor delayedExecutor(long delay, TimeUnit unit, Executor executor)
- Executor delayedExecutor(long delay, TimeUnit unit)
- CompletionStage completedStage(U value)
- CompletionStage failedStage(Throwable ex)
- CompletableFuture failedFuture(Throwable ex)
Наконец, для решения проблемы тайм-аута в Java 9 представлены две новые функции:
- orTimeout()
- completeOnTimeout()
в заключении
В этой статье мы описываем методы и типичные варианты использования класса CompletableFuture.