[JAVA] CompletableFuture : 안정적 비동기 프로그래밍

Jae-Baek Song·2023년 3월 21일

모던자바인액션

목록 보기
11/11
post-thumbnail

저자: 라울-게이브리얼 우르마 , 마리오 푸스코 , 앨런 마이크로프트
도서명: 모던 자바 인 액션
출판사: 한빛미디어


Future의 단순 활용

	public static void main(String[] args) {
        ExecutorService executorService = Executors.newCachedThreadPool();
        var future = executorService.submit(Future::doSomeLongComputation);
        doSomethingElse();
        try{
            future.get(1, TimeUnit.SECONDS);
        }catch (Exception exception){

        }
    }

    private static void doSomeLongComputation() {
        for (int i = 0; i < 10; i++) {
            System.out.println(i + "번째 실행 doSomeLongComputation");
        }
    }

    private static void doSomethingElse(){
        for (int i = 0; i < 10; i++) {
            System.out.println(i + "번째 실행 doSomethingElse");
        }
    }
    
    
0번째 실행 doSomethingElse
0번째 실행 doSomeLongComputation
1번째 실행 doSomethingElse
1번째 실행 doSomeLongComputation
2번째 실행 doSomethingElse
2번째 실행 doSomeLongComputation
3번째 실행 doSomethingElse
3번째 실행 doSomeLongComputation
4번째 실행 doSomethingElse
4번째 실행 doSomeLongComputation
5번째 실행 doSomethingElse
5번째 실행 doSomeLongComputation
6번째 실행 doSomeLongComputation
7번째 실행 doSomeLongComputation
8번째 실행 doSomeLongComputation
9번째 실행 doSomeLongComputation
6번째 실행 doSomethingElse
7번째 실행 doSomethingElse
8번째 실행 doSomethingElse
9번째 실행 doSomethingElse

get 메서드를 호출했을 때 이미 계산이 완료되어 결과가 준비되었다면 즉시 결과를 반환하지만 결과가 준비되지 않았다면 작업이 완료될 때까지 우리 스레드를 블록 시킨다.

오래 걸리는 작업이 영원히 끝나지 않는 문제를 해결하기 위해 스레드가 대기할 최대 타임아웃 시간을 설정하는 것이 좋다.


에러 처리 방법

비동기를 실행하는 동안 에러가 발생하면 어떻게 될까?
예외가 발생하면 해당 스레드에만 영향을 미친다. 즉, 에러가 발생해도 가격 계산은 계속 진행되며 일의 순서가 꼬인다.
결과적으로 클라이언트는 get 메서드가 반환될 때까지 영원히 기다리게 될 수도 있다.

블록 문제가 발생할 수 있는 상황에서는 타임아웃을 활용하는 것이 좋다.
하지만 왜 에러가 발생했는지 알 수 있는 방법이 없다.
따라서 completeExceptionally 메서드를 이용해서 CompletableFutrue 내부에서 발생한 예외를 클라이언트로 전달해야 한다.

public Future<Double> getPriceAsync(String product) {
  CompletableFuture<Double> futurePrice = new CompletableFuture<>();
  new Thread(() -> {
    try {
      double price = calculatePrice(product);
      futurePrice.complete(price);
    } catch {
      futurePrice.completeExceptionally(ex); //에러를 포함시켜 Future를 종료
    }
  }).start();
  return futurePrice;
}

팩토리 메서드 supplyAsync로 CompletableFuture 만들기

public Future<Double> getPriceAsync(String product) {
	return CompletableFuture.supplyAsync(() -> calculatePrice(product));
}

비블록 코드 만들기

public List<String> findPrices(String product) {
  return shops.stream()
    .map(shop -> String.format("%s price is %.2f", shop.getName(), shop.getPrice(product)))
    .collect(toList());
}

네 개의 상점에서 각각 가격을 검색하는 동안 블록되는 시간이 발생할 것이다.

병렬 스트림으로 요청 병렬화하기

public List<String> findPrices(String product) {
  return shops.parallelStream()
    .map(shop -> String.format("%s price is %.2f", shop.getName(), shop.getPrice(product)))
    .collect(toList());
}

이제 네 개의 상점에서 병렬로 검색이 진행되므로 시간은 하나의 상점에서 가격을 검색하는 정도만 소요될 것이다.

CompletableFuture로 비동기 호출 구현하기

List<CompletableFuture<String>> priceFutures = 
  shops.stream()
    .map(shop -> CompletableFuture.suppltAsync(
      () -> String.format("%s price is %.2f", shop.getName(), shop.getPrice(product)))
    .collect(toList());
}

위 코드로 List<CompletableFuture<String>>를 얻을수 있다.
하지만 우리가 재구현하는 findPrices 메서드의 반환 형식은 List<String>이므로 모든 CompletableFuture의 동작이 완료되고 결과를 추출한 다음에 리스트를 반환해야한다.

리스트의 모든 CompletableFutre에 join을 호출해서 모든 동작이 끝나기를 기다린다.
CompletableFuture클래스의 join 메서드는 Future 인터페이스의 get 메서드와 같은 의미를 갖는다. 다만 join은 아무 예외도 발생시키지 않는다는 점이 다르다.
따라서 두 번째 map의 람다 표현식을 try/catch로 감쌀 필요가 없다.

public List<String> findPrices(String product) {
  List<CompletableFuture<String>> priceFutures = 
    shops.stream()
      .map(shop -> CompletableFuture.suppltAsync(
        () -> shop.getName() + "price is " + shop.getPrice(product)))
      .collect(toList());
      
  return priceFutures.stream()
    .map(CompletableFuture::join) //모든 비동기 동작이 끝나길 대기
    .collect(toList());
}

스트림 연산은 게으른 특성이 있으므로 하나의 파이프라인으로 처리했다면 모든 가격 정보 요청 동작이 동기적, 순차적으로 이루어지게 된다.

CompletableFuture로 각 상점의 정보를 요청할 때 기존 요청 작업이 완료되어야 join이 결과를 반환하면서 다음 상점으로 정보를 요청할 수 있기 때문이다.


더 확장성이 좋은 해결 방법

일반적으로 스레드 풀에서 제공하는 스레드 수는 4개다.
그래서 병렬 스트림 버전에서 상점이 한개 추가되는경우 1초 이상 시간이 추가된다.
CompletableFuture는 병렬 스트림에 비해 작업에 이용할 수 있는 Executor를 지정할 수 있다는 장점이 있다.

커스텀 Executor 사용하기

실제로 필요한 작업량을 고려한 풀에서 관리하는 스레드 수에 맞게 Executor를 만들면 좋을 것 같다.

풀에서 관리하는 스레드 수는 어떻게 결정할 수 있을까?

한 상점에 하나의 스레드가 할당될 수 있도록, 상점 수만큼 Executor를 설정한다.

서버 크래시 방지를 위해 하나의 Executor에서 사용할 스레드의 최대 개수는 100 이하로 설정한다.

private final Executor executor = Executors.newFixedThreadPool(Math.min(shops.size(), 100), //상점 수만큼의 스레드를 갖는 풀 생성(0~100 사이)
    new ThreadFactory() {
  public Thread new Thread(Runnable r) {
    Thread t = new Thread(r);
    t.setDeamon(true);
    return t;
  }
});

데몬 스레드를 사용하면 자바 프로그램이 종료될 때 강제로 스레드 실행이 종료될 수 있다.

스트림 병렬화와 CompletableFuture 병렬화

  • I/O가 포함되는 않은 계산 중심의 동작을 실행할 때는 스트림 인터페이스가 가장 구현하기 간단하며 효율적일 수 있다.
  • I/O를 기다리는 작업을 병렬로 실행할 때는 CompletableFuture가 더 많은 유연성을 제공하며, 대기/계산의 비율에 적합한 스레드 수를 설정할 수 있다. 스트림의 게으른 특성 때문에 스트림에서 I/O를 실제로 언제 처리할지 예측하기 어려운 문제도 있다.

비동기 작업 파이프라인 만들기

public class Discount {
  public enum Code {
    NONE(0), SILVER(5), GOLD(10), PLATINUM(15), DIAMOND(20);
    
    private final int percentage;
    
    Code(int percentage) {
      this.percentage = percentage;
    }
  }
  ...
}

enum으로 할인율을 제공하는 코드를 정의하였다.

그리고 getPrice 메서드는 ShopName:price:DiscountCode 형식의 문자열을 반환하도록 수정했다.

public String getPrice(String product) {
  double price = calcuatePrice(product);
  Discount.Code code = Discount.Code.values()[
    random.nextInt(Discount.Code.values().length)];
  return String.format("%s:%.f:%s", name, price, code);
}
public class Quote {
    private final String shopName;
    private final double price;
    private final Discount.Code discountCode;

    public Quote(String shopName, double price, Discount.Code code) {
        this.shopName = shopName;
        this.price = price;
        this.discountCode = code;
    }

    public static Quote parse(String s) {
        String[] split = s.split(":");
        String shopName = split[0];
        double price = Double.parseDouble(split[1]);
        Discount.Code discountCode = Discount.Code.valueOf(split[2]);
        return new Quote(shopName, price, discountCode);
    }

    public String getShopName() {
        return shopName;
    }

    public Double getPrice() {
        return price;
    }

    public Discount.Code getDiscountCode() {
        return discountCode;
    }
}

상점에서 제공한 문자열 파싱은 다음처럼 Quote 클래스로 캡슐화할 수 있다.

상점에서 얻은 문자열을 정적 팩토리 메서드 parse로 넘겨주면 상점 이름, 할인전 가격, 할인된 가격 정보를 포함하는 Quote 클래스 인스턴스가 생성된다.

다음으로 Discount 서비스에서는 Quote 객체를 인수로 받아 할인된 가격 문자열을 반환하는 applyDiscount 메서드도 제공한다.

public class Discount {
  public enum Code {
    ...
  }
  
  public static String applyDiscount(Quote quote) {
    return quote.getShopName() + " price is " + Discount.apply(
      quote.getPrice(), quote.getDiscountCode());
  }
  
  pivate static double apply(double price, Code code) {
    delay();
    return format(price * (100 - code.percentage) / 100);
  }
}

할인 서비스 이용

public List<String> findPrices(String product) {
  return shops.stream()
    .map(shop -> sho.getPrice(product)) //각 상점에서 할인전 가격 얻기
    .map(Quote::parse) //반환된 문자열을 Quote 객체로 변환
    .map(Discount::applyDiscount)  //Quote에 할인 적용
    .collect(toList());
}

코드를 수행해보면 순차적으로 다섯 상점에 가격을 요청하면서 5초가 소요되고, 할인코드를 적용하면서 5초가 소요된다.

앞서 확인한 것처럼 병렬 스트림으로 변환하면 성능을 개선할 수 있다. 하지만 스트림이 사용하는 스레드 풀의 크기가 고정되어 있으므로, 상점 수가 늘어나게되면 유연하게 대응할 수 없다.

따라서 CompletableFutuer에서 수행하는 태스크를 설정할 수 있는 커스텀 Executoer를 정의해서 CPU 사용을 극대화해야한다.

동기 작업과 비동기 작업 조합하기

public List<String> findPrices(String product) {
  List<CompletableFuture<String>> priceFutures = 
    shops.stream() //Stream<Shop>
      .map(shop -> CompletableFuture.supplyAsync(
        () -> shop.getPrice(product), executor)) //Stream<CompletableFuture<Double>>
      .map(future -> future.thenApply(Quote::parse)) //Stream<CompletableFuture<Quote>>
      .map(future -> future.thenCompose(quote ->
        CompletableFuture.supplyAsync(
          () -> Discount.applyDiscount(quote), executor))
      .collect(toList());
      
  return priceFutures.stream()
    .map(CompletableFuture::join)
    .collect(toList());
}
  1. 가격정보 얻기
    팩토리메서드 suuplyAsync에 람다 표현식을 전달해서 비동기적으로 상점에서 정보를 조회했다.
    반환 결과는 Stream<CompletableFuture<String>>이다.
  2. Quote 파싱하기
    CompletableFuture의 thenApply 메서드를 호출해서 Quote 인스턴스로 변환하는 Function으로 전달한다.
    thenApply 메서드는 CompletableFutur가 끝날 때까지 블록하지 않는다. 즉, 이전에 등록한 CompletableFuture가 완료된 후에 실행한다.

Funtion<Double> futurePriceInUSD = CompletableFuture.supplyAsync(() -> shop.getPrice(product)) //1번째 태스크 - 가격정보 요청
  .thenCombine(CompletableFuture.suuplyAsync(
      () -> exchangeService.getRate(Money.EUR, Money.USD)), //2번째 태스크 - 환율정보 요청
    (price, rate) -> price * rate)); //두 결과 합침

독립적으로 실행된 두 개의 CompletableFuture 결과를 합쳐야할 때 thenCombine 메서드를 사용한다.

thenCombine 메서드의 BiFunction 인수는 결과를 어떻게 합질지 정의한다.

독립적인 두 개의 비동기 태스크는 각각 수행되고, 마지막에 합쳐진다.


타임아웃 효과적으로 사용하기

Funtion<Double> futurePriceInUSD = CompletableFuture.supplyAsync(() -> shop.getPrice(product))
  .thenCombine(CompletableFuture.suuplyAsync(
      () -> exchangeService.getRate(Money.EUR, Money.USD)),
    (price, rate) -> price * rate))
  .orTimeout(3, TimeUnit.SECONDS);

Future가 작업을 끝내지 못할 경우 TimeoutException을 발생시켜 문제를 해결할 수 있다.

Funtion<Double> futurePriceInUSD = CompletableFuture.supplyAsync(() -> shop.getPrice(product))
  .thenCombine(CompletableFuture.suuplyAsync(
      () -> exchangeService.getRate(Money.EUR, Money.USD)),
      .completOnTimeout(DEFAULT_RATE, 1, TimeUnit.SECONDS),
    (price, rate) -> price * rate))
  .orTimeout(3, TimeUnit.SECONDS);

compleOnTimeout메서드를 통해 예외를 발생시키는 대신 미리 지정된 값을 사용하도록 할 수도 있다.


CompletableFuture의 종료에 대응하는 방법

get이나 join으로 CompletableFutrue가 완료될 때까지 블록하지 않고 다른 방식으로 CompletableFuture의 종료에 대응하는 방법을 소개한다.

팩토리 메서드 allOf는 전달된 모든 CompletableFuture가 완료된 후에 CompletableFuture<Void>를 반환한다.

이를 통해 모든 결과가 반환되었음을 확인할 수 있다.

CompletableFuture[] futures = findPriceStream("myPhone")
  .map(f -> f.thenAccept(System.out::println))
  .toArray(size -> new CompletableFuture[size]);
CompletableFuture.allOf(futues).join();

만약 CompletableFuture 중 하나만 완료되기를 기다리는 상황이라면 팩토리메서드 anyOf를 사용할 수 있다.


List<CompletableFuture<String>> futures = new ArrayList<>();
futures.add(CompletableFuture.supplyAsync(() -> {
    try {
        Thread.sleep(10000);
        System.out.println("11");
        return "1";
    } catch (InterruptedException e) {
        throw new RuntimeException(e);
    }
}));
futures.add(CompletableFuture.supplyAsync(() -> {
    try {
        Thread.sleep(5000);
        System.out.println("22");
        return "2";
    } catch (InterruptedException e) {
        throw new RuntimeException(e);
    }
}));

CompletableFuture.supplyAsync(() -> {
    try {
        Thread.sleep(5000);
        System.out.println("33");
        return 3;
    } catch (InterruptedException e) {
        throw new RuntimeException(e);
    }
});

futures.stream()
        .peek(i -> System.out.println("start"))
        .map(CompletableFuture::join)
        .peek(System.out::println)
        .collect(toList());
        
start
22
33
11
1
start
2

join() 호출시 ASYNC_POOL에 등록된 모든 클래스를 작업을 수행


thenApply와 thenCompose의 차이점을 예시를 통해 설명해 보겠습니다.

두 가지 메서드가 있다고 가정해 봅시다. getUserInfo(int userId)와 getUserRating(UserInfo userInfo)입니다:

public CompletableFuture<UserInfo> getUserInfo(userId)

public CompletableFuture<UserRating> getUserRating(UserInfo)

두 메서드의 반환 유형은 모두 CompletableFuture입니다.

먼저 getUserInfo()를 호출하고, 완료되면 결과 UserInfo를 가지고 getUserRating()을 호출하고 싶습니다.

getUserInfo() 메서드가 완료되면 thenApplythenCompose를 모두 시도해 봅시다. 차이점은 반환 유형에 있습니다

CompletableFuture<CompletableFuture<UserRating>> f =
    userInfo.thenApply(this::getUserRating);

CompletableFuture<UserRating> relevanceFuture =
    userInfo.thenCompose(this::getUserRating);

thenCompose()는 중첩된 퓨처를 평탄화하는 Scala의 flatMap처럼 작동합니다.

thenApply()는 중첩된 퓨처를 그대로 반환하지만, thenCompose()는 중첩된 CompletableFutures를 평탄화하여 더 많은 메서드 호출을 쉽게 연결할 수 있도록 합니다.

https://stackoverflow.com/questions/43019126/completablefuture-thenapply-vs-thencompose

2개의 댓글

comment-user-thumbnail
2023년 3월 24일

너무 멋져요

답글 달기
comment-user-thumbnail
2023년 3월 24일

너무 멋져요

답글 달기