Пять минут, чтобы понять CompletableFuture

задняя часть

предисловие

Интервьюер: Какие функции JAVA8 вы знаете?

Я: лямбда-выражения, функциональные интерфейсы, потоки

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

Я: Внезапно у меня закружилась голова, какого черта, я только слышал о Будущем, а у меня в сердце было десять тысяч слов ммп.

Не паникуйте, сегодня Xiaobian покажет вам CompletableFuture.

текст

Прежде чем официально разобраться в CompleteFuture, редактору необходимо познакомить вас с Future. Интерфейс Future был представлен в JKD5. Первоначальной целью его дизайна было моделирование результатов, которые произойдут в какой-то момент в будущем. Он моделирует асинхронный расчет и возвращает ссылку на результат выполнения. Когда расчет завершен , вызывающий Эта ссылка доступна. Таким образом, потоки, которые запускают трудоемкие операции, могут быть освобождены через Future, чтобы он мог выполнять другую ценную работу без ожидания. Редактор уже использовал его в работе, и вот пример.

List<CommissionDetailVo> bigList = new ArrayList<>();
List<CommissionDetailVo> onePage = new ArrayList<>();
List<Future> futureList = new ArrayList<>();
try {
    for ( int i = 0; i < THREADS; i++) {

        final int k = i;
        Future list = threadPool.submit(() -> {
            return esSearch(k,settlementDetailSearchVo);
        });
        futureList.add(list);
    }
}catch (Exception e) {
    e.printStackTrace();
    logger.error("Es查询失败", e.getMessage());
}
for ( Future list: futureList) {
        onePage =(List<CommissionDetailVo>) list.get();
        bigList.addAll(onePage);
}

Тогда некоторые студенты спросят, раз есть будущее, зачем мне изучать CompletableFuture? Не торопитесь, пожалуйста, послушайте редактора.

В процессе использования Future вы обнаружите, что у Future есть два способа вернуть результат.

  1. Блокирующий метод get()
  2. Метод опроса isDone()

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

Если мы ожидаем, что результат будет получен автоматически после завершения асинхронной задачи, то Future не может удовлетворить наши потребности. Поэтому JDK8 представил CompletableFuture, о котором мы говорим сегодня, который реализует интерфейс Future, восполняет недостатки режима Future, а также предоставляет очень мощную функцию расширения Future, которая помогает нам упростить сложность асинхронного программирования и предоставляет функционал программирование.Возможность обработки результатов вычислений с помощью обратных вызовов, а также предоставляет функции преобразования, комбинирования и сериализации CompletableFuture.

1. После завершения асинхронной задачи метод объекта может быть автоматически вызван обратно;

2. После сбоя асинхронной задачи метод объекта может быть автоматически вызван обратно;

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

См. следующий пример:

CompletableFuture<Integer> future = CompletableFuture.supplyAsync( () -> feeItem.getAmount());
future.thenAccept(amount -> System.out.println(amount));
future.exceptionally(throwable -> {System.out.println("发生异常"+throwable.toString());return 111;});

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

Базовый API

Runasync и Supplyasync

Создать асинхронную операцию

public static CompletableFuture<Void> runAsync(Runnable runnable)
public static CompletableFuture<Void> runAsync(Runnable runnable, Executor executor)
public static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier)
public static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier, Executor executor)

Если при вызове не указан пул потоков, ForkJoinPool.commonPool() будет использоваться в качестве пула потоков по умолчанию для выполнения асинхронного кода. Разница между двумя методами: runAsync не поддерживает возвращаемые значения, а SupplyAsync поддерживает возвращаемые значения.

Пример:

CompletableFuture<Integer> future1 = CompletableFuture.supplyAsync(() -> feeItem.getAmount());

Метод обратного вызова, когда результат расчета завершен

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

public CompletableFuture<T> whenComplete(BiConsumer<? super T,? super Throwable> action)
public CompletableFuture<T> whenCompleteAsync(BiConsumer<? super T,? super Throwable> action)
public CompletableFuture<T> whenCompleteAsync(BiConsumer<? super T,? super Throwable> action, Executor executor)
public CompletableFuture<T> exceptionally(Function<Throwable,? extends T> fn)

В чем разница между whenComplete и whenCompleteAsync?

whenComplete: поток, выполняющий текущую задачу, продолжит выполнение задачи, переданной в whenComplete.
whenCompleteAsync: отправленные задачи в whenCompleteAsync отправляются в другие пулы потоков для выполнения.

Пример:

CompletableFuture<Integer> future3 = CompletableFuture.supplyAsync(() -> feeItem.getAmount());
future3.whenComplete( (amount, throwable)-> System.out.println(amount));

handle

Хотя метод handle возвращает объект CompleatableFuture, значение объекта отличается от значения, рассчитанного исходным CompleatableFuture. При вычислении значения исходного CompletableFuture или возникновении исключения будет инициирован расчет объекта CompletableFuture, а результат будет рассчитан параметром BiFunction. Поэтому эта группа методов имеет как функции whenComplete, так и функции преобразования.

public <U> CompletableFuture<U>  handle(BiFunction<? super T,Throwable,? extends U> fn)
public <U> CompletableFuture<U>  handleAsync(BiFunction<? super T,Throwable,? extends U> fn)
public <U> CompletableFuture<U>  handleAsync(BiFunction<? super T,Throwable,? extends U> fn, Executor executor)

Пример:

CompletableFuture<Integer> future3 = CompletableFuture.supplyAsync(() -> feeItem.getAmount());
CompletableFuture<Integer> future4 = future3.handleAsync((amount,throwable) -> amount*3);
System.out.println(future4.get());

конвертировать

thenApply передаст результат исходного вычисления CompletableFuture функции fn и использует результат fn в качестве результата вычисления нового CompletableFuture, то есть CompletableFuture преобразуется в CompletableFuture

public <U> CompletableFuture<U> thenApply(Function<? super T,? extends U> fn)
public <U> CompletableFuture<U> thenApplyAsync(Function<? super T,? extends U> fn)
public <U> CompletableFuture<U> thenApplyAsync(Function<? super T,? extends U> fn, Executor executor)

Пример:

CompletableFuture<Integer> a =CompletableFuture.supplyAsync(() -> 1).thenApply(i -> i+1).thenApply(i -> i*i).whenComplete((r,throwable) -> System.out.println(r));
System.out.println(a.get());

Чистое потребление выполняет действие

thenAcpet выполняет действие только над результатом вычисления и не возвращает новое вычисленное значение.Тип возвращаемого значения — Void.

public CompletableFuture<Void>  thenAccept(Consumer<? super T> action)
public CompletableFuture<Void>  thenAcceptAsync(Consumer<? super T> action)
public CompletableFuture<Void>  thenAcceptAsync(Consumer<? super T> action, Executor executor)

thenCombine

Два результата объекта CompletableFuture могут быть объединены

public <U,V> CompletionStage<V> thenCombine(CompletionStage<? extends U> other,BiFunction<? super T,? super U,? extends V> fn);
public <U,V> CompletionStage<V> thenCombineAsync(CompletionStage<? extends U> other,BiFunction<? super T,? super U,? extends V> fn);
public <U,V> CompletionStage<V> thenCombineAsync(CompletionStage<? extends U> other,BiFunction<? super T,? super U,? extends V> fn,Executor executor);

thenCompose

Объедините две асинхронные операции, когда первая операция завершится, перваяCompletableFutureОбъект вызывает thenCompose, передавая ему функцию. когда первыйCompletableFutureПосле завершения выполнения его результат будет использоваться как параметр функции, а возвращаемое значение этой функции является первымCompletableFutureВозврат выполнения ввода вычисляет второйCompletableFutureобъект.

public <U> CompletableFuture<U> thenCompose(Function<? super T, ? extends CompletionStage<U>> fn)

Пример:

List<CompletableFuture<Integer>> list = feeItemList.stream()
        .map(feeItems ->
        CompletableFuture.supplyAsync(() -> feeItems.getAmount()))
        .map(future -> future.thenApply(i -> i * 2))
        .map(future -> future.thenCompose(amount -> CompletableFuture.supplyAsync( () -> amount*3 )))
        .collect(Collectors.toList());     System.out.println(list.stream().map(CompletableFuture::join).collect(Collectors.toList()));

Метод соединения Compfuture здесь заключается в получении вычисленного значения CompletableFuture.

Вспомогательные методы allOf и anyOf

allOf: вычисляется, когда все CompletableFuture были выполнены

Основной поток блокируется до тех пор, пока не будут выполнены все CompletableFuturers.

Anyof: выполнить расчет при выполнении любая задача

public static CompletableFuture<Void>       allOf(CompletableFuture<?>... cfs)
public static CompletableFuture<Object>     anyOf(CompletableFuture<?>... cfs)

Пример:

CompletableFuture<Integer> future1 = CompletableFuture.supplyAsync(() -> feeItem.getAmount());
CompletableFuture<Integer> future2 = CompletableFuture.supplyAsync(() -> feeItem1.getAmount());
CompletableFuture<Object> f1 = CompletableFuture.anyOf(future2, future1);
CompletableFuture<Void> f2= CompletableFuture.allOf(future2, future1);
System.out.println(f1.get()+" "+f2.get());

Есть еще много методов CompletableFuture, около 60, поэтому я не буду здесь вдаваться в подробности.

Параллельный поток или CompletableFuture

Существует два способа выполнения параллельных операций над коллекциями.Один – использовать параллельные потоки, то есть parallelStream (ожидание результата будет заблокировано), и другой – использовать CompletableFuture. CompletableFuture обеспечивает большую гибкость, и вы можете передавать пул потоков самостоятельно, поэтому вы можете настроить размер пула потоков, чтобы убедиться, что общие вычисления не блокируются из-за того, что потоки ожидают ввода-вывода. Итак, общая рекомендация выглядит следующим образом:

1. Если вы выполняете ресурсоемкие операции и не имеете операций ввода/вывода, рекомендуется использовать параллельные потоки, т.к. реализация проста и эффективность коллег может быть самой высокой. (Если все потоки требуют больших вычислительных ресурсов, нет необходимости создавать больше потоков, чем число объединенных процессоров).

2. Если ваша параллельная единица работы также включает ожидание операций ввода-вывода или сетевых операций, гибкость использования CompletableFuture будет лучше.

Используйте FullableFuture для оптимизации производительности вашего кода

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

优化前
До оптимизации

В чем проблема с этим?

1. Пропускная способность системы невысокая

2. Интерфейс вызывается часто и долго

Мы можем сделать следующие улучшения:

性能优化
оптимизация производительности

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

Получайте запросы и запрашивайте пакетами. После возврата результатов пакетного запроса результаты возвращаются в соответствующий поток в соответствии с информацией о пакете.

Ниже приведен пример мелкого письма:

public class AsyncServie {

    @Autowired
    private OrderService orderService;

    public AtomicInteger atomicInteger = new AtomicInteger(0);

    class Request {
        String orderCode;
        String serialNo;
        CompletableFuture<Map<String,Object>> future;
    }

    LinkedBlockingDeque<Request> blockingDeque = new LinkedBlockingDeque<>();

    public Map<String, Object> queryOrderInfo(String orderCode) throws InterruptedException, ExecutionException {
        Request request = new Request();
        request.orderCode = orderCode;
        request.serialNo = UUID.randomUUID().toString();
        CompletableFuture<Map<String, Object>> future = new CompletableFuture<>();
        request.future = future;
        blockingDeque.add(request);

        atomicInteger.getAndIncrement();
        System.out.println(atomicInteger);
        //监听有没有返回值 一直阻塞
         return future.get();
    }

    @PostConstruct
    public void init() {
        //每隔10ms去队列获取请求,批量发起
        ScheduledExecutorService executorService = Executors.newScheduledThreadPool(2);
        executorService.scheduleAtFixedRate( () -> {
            int size = blockingDeque.size();
            if (0 == size) {
                return;
            }
            List<Map<String, String>> params = Lists.newArrayList();
            List<Request> requestList = Lists.newArrayList();
            for (int i=0; i< size; i++){
                Request request = blockingDeque.poll();
                Map<String ,String> map = Maps.newHashMap();
                map.put("orderCode", request.orderCode);
                map.put("serialNo", request.serialNo);
                params.add(map);
                requestList.add(request);
            }
            System.out.println("批量处理的size"+size);
            System.out.println(Thread.currentThread().getName());
            List<Map<String, Object>>  orderInfo =  orderService.getOrderInfo(params);

            //匹配对应的serialNo
             Optional.ofNullable(requestList).ifPresent(requests -> requests.forEach(request -> {
                String serialNo = request.serialNo;
                Optional.ofNullable(orderInfo).ifPresent(orderInfos -> orderInfos.forEach(response -> {
                    String serial = response.get("serialNo").toString();
                    if (Objects.equals(serialNo, serial)) {
                        request.future.complete(response);
                    }
                }));
            }));


        }, 0, 10, TimeUnit.MILLISECONDS);

    }
}

Ну, это все на сегодня.