안동민 개발노트

본문 시작

일괄 실행과 비동기 결합

여러 외부 작업을 직렬 호출하는 주문 흐름을 분해하고 일괄 Future 수집과 CompletableFuture 조합에서 부분 실패·실행기·결과 순서를 통제합니다.

배송 알림, 회계 반영, 재고 예약이 서로 독립이라면 순서대로 기다릴 이유가 없습니다.

fan-out은 독립 작업을 여러 분기로 제출하고 fan-in에서 결과를 합칩니다.

invokeAll은 작업 컬렉션을 제출하고 Future 목록을 입력 순서로 돌려주며, CompletableFuture는 변환과 복구 단계를 선언적으로 연결합니다.

비동기 조합에서도 실행기 선택과 실패 정책은 자동으로 정해지지 않습니다.

기본 공용 풀에 블로킹 I/O를 몰아넣지 말고 전용 실행기를 전달합니다.

한 하위 작업 실패가 전체 주문 실패인지 부분 성공인지 결과 타입으로 정합니다.


독립 시스템을 순차 호출한 지연

bad/SequentialOrderFanout.java
public final class SequentialOrderFanout {
    static String call(String name) throws InterruptedException {
        Thread.sleep(100);
        return name + "-ok";
    }

    public static void main(String[] args) throws Exception {
        long start = System.nanoTime();
        System.out.println(call("shipping"));
        System.out.println(call("accounting"));
        System.out.println(call("inventory"));
        System.out.println("millis=" + (System.nanoTime() - start) / 1_000_000);
    }
}

이번 실행은 shipping-ok, accounting-ok, inventory-ok 순서와 millis=305를 출력했습니다. 원문은 각각 sleep(100)을 호출한 모의 작업이며, millis에는 출력 비용과 스케줄링도 포함됩니다. 세 설정값의 합이 정확한 실행 시간을 보장하지는 않습니다.

병렬화 전에는 세 호출이 정말 독립인지와 외부 시스템 동시 요청 한도를 확인합니다.


팬아웃 설계 규칙

  • 입력 작업은 서로 결과 의존성이 없는 단위로 나눈다.
  • 외부 시스템별 동시성 한도를 실행기 크기에 반영한다.
  • 전체 시간 예산을 한 번 정하고 하위 작업에 배분한다.
  • 부분 실패를 허용하면 성공과 실패를 같은 결과 모델에 담는다.
  • CompletableFuture의 예외 단계에서 원인을 숨기는 기본값을 만들지 않는다.
  • 모든 분기가 끝난 뒤 전용 실행기를 종료한다.

invokeAll 결과 순서

src/InvokeAllOrderTasks.java
import java.util.List;
import java.util.concurrent.Callable;
import java.util.concurrent.Executors;

public final class InvokeAllOrderTasks {
    public static void main(String[] args) throws Exception {
        var executor = Executors.newFixedThreadPool(3);
        try {
            List<Callable<String>> tasks = List.of(
                    () -> "shipping-ok",
                    () -> "accounting-ok",
                    () -> "inventory-ok");
            var futures = executor.invokeAll(tasks);
            for (var future : futures) {
                System.out.println(future.get());
            }
        } finally {
            executor.shutdown();
        }
    }
}

이 main의 invokeAll은 세 작업의 완료를 기다린 뒤 입력 목록 순서의 Future를 반환합니다. 순회 출력은 shipping·accounting·inventory 순서이며 작업 완료 순서를 측정한 것은 아닙니다.

시간 제한 오버로드에서는 돌아올 때 끝나지 않은 작업을 취소합니다. 반환된 각 Future의 get은 정상값, 실행 예외 또는 취소를 구분해 처리해야 합니다. 이 main은 시간 제한 오버로드를 실행하지 않습니다.


CompletableFuture 결과 조합

src/CompletableOrderFanout.java
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.function.Supplier;

public final class CompletableOrderFanout {
    record Result(String system, boolean success, String detail) {}

    static CompletableFuture<Result> call(String system, Supplier<String> action,
                                          ExecutorService executor) {
        return CompletableFuture.supplyAsync(action, executor)
                .handle((value, error) -> error == null
                        ? new Result(system, true, value)
                        : new Result(system, false, error.getClass().getSimpleName()));
    }

    public static void main(String[] args) {
        ExecutorService executor = Executors.newFixedThreadPool(3);
        try {
            var futures = List.of(
                    call("shipping", () -> "notified", executor),
                    call("accounting", () -> {
                        throw new IllegalStateException("closed");
                    }, executor),
                    call("inventory", () -> "reserved", executor));
            CompletableFuture.allOf(futures.toArray(CompletableFuture[]::new)).join();
            System.out.println(futures.stream().map(CompletableFuture::join).toList());
        } finally {
            executor.shutdown();
        }
    }
}
분기 결과가 handle을 거쳐 정상 완료한 Result가 된다

CompletableOrderFanout의 세 공급 작업 결과와 handle이 보관하는 Result를 비교하며, 실패 원인에서 단순 클래스 이름만 남는 변환을 드러냅니다.

분기 결과가 handle을 거쳐 정상 완료한 Result가 된다
시스템supplyAsync 작업의 원래 결과handle이 만드는 Result
shipping문자열 notified 반환success=true · detail=notified
accountingIllegalStateException("closed") 발생success=false · detail=CompletionException
inventory문자열 reserved 반환success=true · detail=reserved
shipping
supplyAsync 작업의 원래 결과: 문자열 notified 반환
handle이 만드는 Result: success=true · detail=notified
accounting
supplyAsync 작업의 원래 결과: IllegalStateException("closed") 발생
handle이 만드는 Result: success=false · detail=CompletionException
inventory
supplyAsync 작업의 원래 결과: 문자열 reserved 반환
handle이 만드는 Result: success=true · detail=reserved

이번 출력의 success·detail은 표의 마지막 열과 같았습니다. handle은 작업 예외도 Result 값으로 바꾸므로 세 후속 Future가 정상 완료합니다. allOf는 그 완료를 기다리고 main은 futures 목록 순서로 결과를 수집하므로, 이 순서가 작업 완료 순서는 아닙니다.

이 JDK에서 detail에 남는 이름은 CompletionException이며, 원래 예외의 IllegalStateException과 메시지 closed는 Result에 보존하지 않습니다. 따라서 현재 출력만으로 원래 실패 원인을 복원할 수는 없습니다.

Future의 정상 완료가 업무 성공을 뜻하지는 않습니다. 현재 main은 Result 목록을 출력할 뿐이며, 후속 업무를 진행할 호출자는 success를 확인해야 합니다.

전체 실패가 요구라면 예외를 보존하고 allOf의 CompletionException 원인을 처리합니다. allOf 자체는 다른 분기를 취소하지 않습니다.

supplyAsync에는 전용 실행기를 넘기지만 handle과 thenCombine은 non-async 단계이므로 완료시키는 스레드나 결합을 등록하는 스레드에서 실행될 수 있습니다.


팬아웃 조합 방식

요구도구특징
작업 목록 전체 대기invokeAll입력 순서 Future
완료된 순서로 개별 수집CompletionService완료된 Future를 take·poll로 수집
성공한 결과 하나invokeAny반환할 때 미완료 작업 취소
단계별 변환CompletableFuturethenApply·thenCompose
여러 독립 분기allOf결과는 별도 수집

부분 실패를 합성 결과에 남기는 방법

독립 호출의 병렬화가 줄이는 시간은 실행기 대기와 외부 시스템의 제한에 따라 달라집니다. 이 장의 비동기 예제는 즉시 반환하거나 예외를 던지므로 앞의 sleep(100) 세 번과 같은 부하의 속도 비교가 아닙니다.

한 공급자가 실패해도 다른 결과로 응답할 수 있는지, 둘 다 필요해서 전체를 실패시켜야 하는지를 조합 전에 결정합니다.

성공 값만 남기는 코드는 장애를 “데이터 없음”과 섞으므로 공급자 이름, 소요 시간, 실패 종류를 가진 결과 타입이 더 낫습니다.

전체 제한 시간은 각 호출 제한 시간의 합이 아닙니다.

요청 시작 시각에서 하나의 마감 시각을 계산하고 후속 단계에 남은 시간을 전달해 예산을 관리합니다. 그래도 스케줄링과 정리 작업을 포함한 전체 반환 시간의 상한이 자동으로 보장되지는 않습니다. 이 장의 비동기 원문에는 그 마감 계산과 시간 제한이 없습니다.

취소된 작업이 외부 요청을 계속하지 않도록 클라이언트의 취소 규칙도 정해야 합니다. CompletableFuture.cancel의 mayInterruptIfRunning 인자는 작업자 인터럽트에 사용되지 않으므로, 이를 FutureTask의 cancel(true)와 같은 중단 방식으로 가정하지 않습니다. 이 원문은 외부 요청이나 취소를 실행하지 않습니다.


연습 문제

두 공급자의 작업을 전용 실행기에 모두 제출하고 둘 다 성공하면 최저가를 반환하세요.

한쪽이 실패하면 성공한 가격을 쓰고 둘 다 실패하면 예외를 던집니다.

정답과 해설
exercise/BestPriceSolution.java
import java.util.OptionalInt;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executors;

public final class BestPriceSolution {
    static OptionalInt safe(java.util.function.IntSupplier supplier) {
        try { return OptionalInt.of(supplier.getAsInt()); }
        catch (RuntimeException e) { return OptionalInt.empty(); }
    }

    public static void main(String[] args) {
        var executor = Executors.newFixedThreadPool(2);
        try {
            var a = CompletableFuture.supplyAsync(() -> safe(() -> 120), executor);
            var b = CompletableFuture.supplyAsync(() -> safe(() -> {
                throw new IllegalStateException();
            }), executor);
            int best = a.thenCombine(b, (left, right) -> {
                if (left.isPresent() && right.isPresent()) return Math.min(left.getAsInt(), right.getAsInt());
                if (left.isPresent()) return left.getAsInt();
                if (right.isPresent()) return right.getAsInt();
                throw new IllegalStateException("all providers failed");
            }).join();
            System.out.println("best=" + best);
        } finally {
            executor.shutdown();
        }
    }
}

이번 출력은 best=120이었습니다. 이 main은 한 공급자의 120과 다른 공급자의 RuntimeException만 다룹니다. 둘 다 성공해 최저가를 고르는 경로와 둘 다 실패해 예외를 던지는 경로는 코드에 있지만 이 입력에서는 실행하지 않았습니다.

safe는 RuntimeException을 빈 가격으로 바꾸면서 원인 정보도 버립니다. Error는 이 catch의 대상이 아닙니다. 실패를 빈 가격으로 표현하는 정책이 업무상 허용되는지 먼저 결정해야 합니다.


병렬 외부 호출의 종료 기준

fan-out은 독립성과 외부 동시성 한도가 확인된 작업에만 적용합니다.

invokeAll은 단순한 일괄 수집에, CompletableFuture는 단계 조합과 부분 복구에 적합합니다.

어느 쪽이든 전용 실행기와 실패 의미를 명시해야 합니다.