Я слышал, вы еще не знаете CompletableFuture?

Java
Я слышал, вы еще не знаете CompletableFuture?

Публичный аккаунт WeChat:Программист Сяо И

предисловие

java8Это было очень распространено в повседневной разработке и кодировании.Освоение и использование может использовать несколько строк упрощенного кода в разработке для выполнения необходимых функций. Сегодня представимCompletableFutureкак использовать практику в производственной среде.CompletableFutureкласс какJava 8 Concurrency APIЗнакомые учащиеся, знакомые с улучшением, должны знать, что в CompletableFuture в Java 9 есть некоторые улучшения, и Xiaoyi введет объяснение позже. Необходимые знания, которые вам необходимо знать, чтобы прочитать эту статью:函数式编程,线程池原理等. Студенты, которые еще не знакомы с ним, могут прочитать предыдущие статьи, без лишних слов, приступим.

моделирование сцены

Чтобы лучше выразиться, мы объясним на примерах.Предположим, Сяойи получил сегодня задание TL и попросил выполнить функцию вытягивания данных в режиме реального времени.После завершения он будет уведомлен о том, что вытягивание завершено. Предположим, что данные извлечения необходимо получить от трех служб A, B и C, а службу D необходимо вызвать после завершения извлечения и отправки.

Изменение требования 1: Извлечение данных должно быть получено из службы E, но это будет зависеть от результатов, полученных из службы A. Изменение требования 2: более 10 000 данных могут быть извлечены из службы A за раз, но производительность службы E не может поддерживать большие вызовы.ProviderЗаканчивается карманами с ограниченным потоком. Изменение требования 3: В процессе извлечения данных необходимо обеспечить целостность данных, и не должно возникать статистических ошибок.

зачем использоватьCompletableFutureВозможно, вы сказали, что это может обеспечить jdk5.0.FutureДля этого мы поместили три сервисных интерфейса из A, B и C, необходимые для загрузки данных вFutureTask, выполните результат данных асинхронно, а затем вызовите службу D синхронно. Хорошо, просто реализовать эту функцию не проблема, но в чем недостатки и как ее можно улучшить? Мы можем видеть через комментарии исходного кодаFutureРезультат, возвращаемый классом, необходимо заблокировать и дождаться, пока метод get вернет результат, что обеспечиваетisDone()метод, чтобы проверить, закончились ли асинхронные вычисления,get()Метод ожидает завершения асинхронной операции и получает результат вычисления. Подождите, пока все будущие задачи будут завершены, уведомите поток о получении результатов и слиянии.

С точки зрения производительности необходимо дождаться выполнения всех задач в коллекции Future (это требование не является проблемой, а затем посмотреть вниз), с точки зрения надежности,FutrueИнтерфейс не имеет методов для выполнения вычислительной композиции или обработки возможных ошибок. Из расширения функции,FutureИнтерфейс не может выполнять несколько асинхронных вычислений независимо друг от друга, при этом второе зависит от результата первого. И сегодняшний геройCompletableFutureВсе они выполняют вышеуказанные функции с примерно 50 различными композициями, комбинациями, выполнением асинхронных вычислительных шагов и обработкой ошибок. (Нереально выучить все методы, освоить души и основные методы, их можно состряпать по закону)

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

Существует слишком много API, чтобы их можно было кратко перечислить. Читатели могут учиться сами. Цель этой статьи не в том, чтобы представить API.

//任务 A 执行完执行 B,执行 B 不需要依赖 A 的结果同时 B 不返回结果。
CompletableFuture.supplyAsync(() -> "resultA").thenRun(() -> {}); 

/** 任务 A 执行完执行 B,B 执行依赖 A 结果同时 B 不返回结果*/
CompletableFuture.supplyAsync(() -> "resultA").thenAccept(resultA -> {});

/**  任务 A 执行完执行 B,B 执行依赖 A 结果同时 B 返回结果*/
CompletableFuture.supplyAsync(() -> "resultA").thenApply(resultA -> resultA + " resultB");

CompletableFuture<String> completableFuture = CompletableFuture.supplyAsync(() -> "orange")
.thenCompose(s -> CompletableFuture.supplyAsync(() -> s + " csong"));
//trueassertEquals("orangecsong", completableFuture.get());

Ваш вопрос:thenComposeРазве метод не такой же, как thenApply для реализации комбинированного расчета результатов? Это было немного запутанно, когда я только учился, на самом деле их внутренние формы разные, они такие же, как Stream и Java, доступные в Java 8.OptionalКатегорияmapа такжеflatMapВ этом методе заложена аналогичная дизайнерская идея. оба получают CompletableFuture и применяют его к вычисляемому результату, ноthenCompose(flatMap)Метод получает функцию, которая возвращает другой объект CompletableFuture того же типа.

CompletableFuture<String> completableFuture = CompletableFuture.supplyAsync(() -> "orange")
.thenCombine(CompletableFuture.supplyAsync(() -> " chizongzi"), (s1, s2) -> s1 + s2));
// assertEquals("orange chizongzi", completableFuture.get());

thenCombineМетод предназначен, когда вы хотите использовать несколько результатов вычислений, и последующая обработка также должна полагаться на возвращаемое значение, возвращается первый результат вычисления."orange", возвращается второй результат вычисления "chizongzi", склейка результатов, то результат"orange chizongzi"Ла. Вы можете спросить, а что, если результат не нужно обрабатывать?thenAcceptBothсможет реализовать вашу функцию. тогда это иthenApplyВ чем разница?thenCompose()Метод заключается в использовании предыдущего Future в качестве параметра. это напрямую сделает результат новым Future вместо того, что мы сделали вthenApply()Вложенный Future посередине используется для соединения двух CompletableFuture, который должен генерировать новый CompletableFuture. Поэтому, если вы хотите продолжить вложение и связывание методов CompletableFuture, лучше всего использоватьthenCompose().

 public static CompletableFuture<Void> allOf(CompletableFuture<?>... cfs){...}

Когда нам нужно выполнить несколько задач параллельно, мы обычно хотим дождаться их выполнения, а затем обработать их объединенный результат. CompletableFuture предоставляетallOfСтатический метод позволяет дождаться завершения всех задач, но его возвращаемый тип — CompletableFuture. Ограничение состоит в том, что он не возвращает исчерпывающие результаты для всех задач. Вместо этого вам придется вручную получать результаты из фьючерсов. Итак, как это решить, CompletableFuture предоставляетjoin()Это можно решить, и Xiaoyi также может использовать Stream для достижения этой цели.

String multiFutures= Stream.of(future1, future2, future3).map(CompletableFuture::join)
.collect(Collectors.joining(" "));assertEquals("Today is sun", multiFutures);

Так как же CompletableFuture обрабатывает исключения?

public CompletableFuture<T> exceptionally(Function<Throwable, ? extends T> fn);
public <U> CompletionStage<U> handle(BiFunction<? super T, Throwable, ? extends U> fn);

Если результат A, результат B, результат C имеют исключения в приобретении

CompletableFuture.supplyAsync(() -> "resultA").thenApply(resultA -> resultA + " resultB")
.thenApply(resultB -> resultB + " resultC")

CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {throw new RuntimeException();})
.exceptionally(ex -> "errorResultA")
.thenApply(resultA -> resultA + " resultB")
.thenApply(resultB -> resultB + " resultC")

В приведенном выше коде задача A выдает исключение, а затем проходитexceptionally()Метод обрабатывает исключение и возвращает новый результат, который передается задаче Б. При иновке метода future.join результат выдаст "errorResultA resultB result C" Вышеупомянутый метод в основном использует базовый функциональный API.Talk is cheap , show me code. Умные друзья, давайте потренируемся!

С момента последней статьиВы все еще беспокоитесь о тайм-ауте интерфейса rpc?В конце статьи мы говорим о больших пакетных вызовах, среди которых есть последовательные вызовы Invoke, а по сути разбираем, как реализовать асинхронные вызовы с помощью CompletableFuture?

CompletableFuture в действии

public class AsyncInvokeUtil {

    private AsyncInvokeUtil() {}

    /**
     * @param paramList 源数据 (需处理数据载体)
     * @param buildParam 中转函数 (获取的结果做一层trans,来满足调用服务条件)
     * @param transParam 中转函数 (获取的结果做一层trans,来满足调用服务条件)
     * @param processFunction 中转处理函数
     * @param size 分批大小
     * @param executorService 暴露外部自定义实现线程池(demo没判空,可以做成非必传)
     * @param <R>
     * @param <T>
     * @param <P>
     * @param <k>
     * @return
     * @throws ExecutionException
     * @throws InterruptedException
     */
    public static <R, T, P, k> List<R> partitionAsyncInvokeWithRes(List<T> paramList,
                                                                      Function<List<T>, P> buildParam,
                                                                      Function<P, List<k>> transParam,
                                                                      Function<List<k>,List<R>> processFunction,
                                                                      Integer size,
                                                                      ExecutorService executorService) throws ExecutionException, InterruptedException {
        List<CompletableFuture<List<R>>> completableFutures = Lists.partition(paramList, size).stream()
                .map(buildParam)
                .map(transParam)
                .map(eachList -> CompletableFuture.supplyAsync(() ->
                        processFunction.apply(eachList), executorService))
                .collect(Collectors.toList());
        //get
        CompletableFuture<Void> finishCompletableFuture = CompletableFuture.allOf(completableFutures.toArray(new CompletableFuture[0]));
        finishCompletableFuture.get();
        return completableFutures.stream().map(CompletableFuture::join)
                .filter(Objects::nonNull).reduce(new ArrayList<>(), (resList1, resList2) -> {
                    resList1.addAll(resList2);
                    return resList1;
                });
    }
}

наконец

  • Статьи все оригинальные, а оригинальность непростая.Благодаря платформе Nuggets я чувствую, что кое-что приобрел.Помогите Sanlian,пополните
  • Все коды, диаграммы последовательности и диаграммы архитектуры, задействованные в статье, являются общими и могут быть получены, оставив сообщение на официальном аккаунте.
  • Если в статье есть ошибки, пожалуйста, прокомментируйте и оставьте сообщение, чтобы указать, и добро пожаловать на перепечатку, пожалуйста, отметьте источник