/
linto
/
opensearch-example
Обзор
Документация
Войти
/
linto
/
opensearch-example
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
src/test/java/com/example/demo/future/FutureTest.java
113 строк
4 KB
linto
init
30 июл 2025, 11:44
30 июл 2025, 11:44
dc5aaf3
Код
Авторство
О чём код?
package com.example.demo.future; import java.time.LocalDateTime; import java.time.temporal.ChronoUnit; import java.util.ArrayList; import java.util.Iterator; import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.TestInstance; import org.junit.jupiter.api.TestInstance.Lifecycle; import lombok.AllArgsConstructor; import lombok.Data; import lombok.extern.slf4j.Slf4j; @Slf4j @TestInstance(Lifecycle.PER_CLASS) public class FutureTest { private final ExecutorService executor = Executors.newFixedThreadPool(2); private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1); private List<DelayedFuture<String>> awaitedFutures = new ArrayList<>(); private final long TIMEOUT_MILLIS = 5000L; @BeforeAll void runScheduler() { scheduler.scheduleAtFixedRate(this::checkResponseReceived, 1L, 1L, TimeUnit.SECONDS); } @Test void realLifeTest() { CompletableFuture<Void> responseFuture = awaitResponse(3000, "success") .thenAcceptAsync(this::processResponse, executor); CompletableFuture<Void> requestFuture = sendRequestAsync(10000).exceptionally(t -> { responseFuture.completeExceptionally(new RuntimeException("Error sending request", t)); return null; }); CompletableFuture<Void> requestFuture2 = sendRequestAsync(10000); // CompletableFuture<Void> requestFuture3 = sendRequestAsync(3000); // CompletableFuture<Void> requestFuture4 = sendRequestAsync(3000); // CompletableFuture<Void> requestFuture5 = sendRequestAsync(3000); responseFuture.join(); } CompletableFuture<String> awaitResponse(long millis, String response) { LocalDateTime now = LocalDateTime.now(); CompletableFuture<String> future = new CompletableFuture<>(); DelayedFuture<String> delayedFuture = new DelayedFuture<String>(future, response, now.plus(millis, ChronoUnit.MILLIS), now.plus(TIMEOUT_MILLIS, ChronoUnit.MILLIS)); awaitedFutures.add(delayedFuture); return future; } private CompletableFuture<Void> sendRequestAsync(long millis) { return CompletableFuture.runAsync(() -> { log.info("sendRequestAsync()"); try { // отправляем запрос Thread.sleep(millis); // throw new RuntimeException("Can't send request!"); } catch (InterruptedException e) { throw new RuntimeException(e); } }, executor); } private void checkResponseReceived() { log.info("checkResponseReceived(): {}", awaitedFutures.size()); LocalDateTime now = LocalDateTime.now(); Iterator<DelayedFuture<String>> iterator = awaitedFutures.iterator(); while (iterator.hasNext()) { DelayedFuture<String> delayedFuture = iterator.next(); if (delayedFuture.getTimeout().isBefore(now)) { delayedFuture.getFuture().completeExceptionally(new TimeoutException("Response timed out")); } if (delayedFuture.getEnd().isBefore(now)) { delayedFuture.getFuture().complete(delayedFuture.getResult()); iterator.remove(); } } } private void processResponse(String response) { log.info("processResponse()"); log.info("Response received: {}", response); } @Data @AllArgsConstructor public static class DelayedFuture<T> { private CompletableFuture<T> future; private T result; private LocalDateTime end; private LocalDateTime timeout; } }