안동민 개발노트

본문 시작

Future.get과 병렬 합산

제출 직후 get이 후속 제출을 늦추는 흐름을 분석하고, 작업을 먼저 제출한 뒤 수집 시작 시점의 공통 예산으로 결과를 모읍니다.

Future를 사용한다고 자동으로 병렬 실행이 되는 것은 아닙니다.

작업 하나를 제출하고 즉시 get한 뒤 다음 작업을 제출하면 호출자가 첫 완료를 기다리는 동안 두 번째 계산은 시작조차 하지 않습니다.

독립 작업을 먼저 제출하고 나중에 결과를 수집하면 실행 구간이 겹칠 수 있습니다.

get이 블로킹하는 위치도 응답 설계의 일부입니다.

메인 요청 스레드에서 무제한 get을 호출하면 느린 작업 하나가 전체 응답을 붙잡습니다.

공통 마감 시간을 전달하거나 비동기로 결과를 연결하고, 시간 초과 뒤의 취소와 종료 정책도 정합니다.


submit·get 반복에 따른 직렬화

bad/SequentialFutureWait.java
import java.util.concurrent.Executors;

public final class SequentialFutureWait {
    public static void main(String[] args) throws Exception {
        var executor = Executors.newFixedThreadPool(2);
        long start = System.nanoTime();
        try {
            int first = executor.submit(() -> {
                Thread.sleep(200);
                return 20;
            }).get();
            int second = executor.submit(() -> {
                Thread.sleep(200);
                return 22;
            }).get();
            System.out.println("sum=" + (first + second));
            System.out.println("millis=" + (System.nanoTime() - start) / 1_000_000);
        } finally {
            executor.shutdown();
        }
    }
}

한 번의 실행에서는 sum=42, millis=409를 출력했습니다. 두 번째 제출이 첫 get의 반환 뒤에 있으므로 두 sleep 작업은 차례로 실행됩니다. millis에는 제출·수집·출력 비용과 스케줄링이 포함됩니다.

아래 ParallelRangeSum은 sleep 없이 작은 정수 범위를 합산합니다. 이 절에는 같은 sleep 작업을 병렬로 실행한 400ms 대 200ms 비교 실험이 없습니다.


병렬 분할의 검토 항목

  • 서로 독립적인 작업만 동시에 제출한다.
  • 제출 단계에서 모든 Future 손잡이를 확보한다.
  • 결과 수집 순서와 업무 결과 순서를 별도로 결정한다.
  • 분할 비용이 계산 이득보다 작은지 측정한다.
  • 한 작업 실패가 전체 실패인지 부분 결과 허용인지 정한다.
  • 정한 수집 마감 시간에서 남은 시간을 각 get에 전달한다.

병렬 범위 합산

src/ParallelRangeSum.java
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;

public final class ParallelRangeSum {
    static long sum(long start, long endInclusive) {
        long result = 0;
        for (long value = start; value <= endInclusive; value++) {
            result += value;
        }
        return result;
    }

    public static void main(String[] args) throws Exception {
        ExecutorService executor = Executors.newFixedThreadPool(2);
        try {
            List<Future<Long>> futures = List.of(
                    executor.submit(() -> sum(1, 50)),
                    executor.submit(() -> sum(51, 100)));
            long total = 0;
            for (Future<Long> future : futures) {
                total += future.get();
            }
            System.out.println("total=" + total);
        } finally {
            executor.shutdown();
        }
    }
}
첫 get 이전에 제출된 작업 수를 비교한다

SequentialFutureWait의 두 sleep 작업과 ParallelRangeSum의 두 합산 작업을 코드 순서 관점에서 비교합니다. 같은 부하의 시간 비교가 아닙니다.

첫 get 이전에 제출된 작업 수를 비교한다
원문 프로그램첫 get 전에 제출한 작업다음 작업이 시작할 수 있는 경계
SequentialFutureWait첫 sleep 작업 하나첫 결과를 받은 다음에야 두 번째 sleep 작업을 제출
ParallelRangeSum1..50 합산과 51..100 합산 모두첫 결과 수집 전에 둘 다 실행할 수 있음
SequentialFutureWait
첫 get 전에 제출한 작업: 첫 sleep 작업 하나
다음 작업이 시작할 수 있는 경계: 첫 결과를 받은 다음에야 두 번째 sleep 작업을 제출
ParallelRangeSum
첫 get 전에 제출한 작업: 1..50 합산과 51..100 합산 모두
다음 작업이 시작할 수 있는 경계: 첫 결과 수집 전에 둘 다 실행할 수 있음

둘 다 미리 제출하면 실행이 겹칠 수 있지만, 작은 작업의 실제 겹침을 측정한 것은 아닙니다. 두 프로그램의 작업 내용이 다르므로 출력 시간을 같은 부하의 속도 비교로 쓰지 않습니다.

결과는 5,050입니다.

범위가 겹치거나 빠지지 않는지 분할 경계도 테스트해야 합니다.

이 작은 합산은 구조 설명용입니다. 병렬 준비 비용과 계산 이득은 실제 입력 크기와 부하에서 비교해야 합니다.


하나의 기한으로 여러 결과 기다리기

src/DeadlineFutureCollector.java
import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;

public final class DeadlineFutureCollector {
    static <T> List<T> collect(List<Future<T>> futures, Duration budget) throws Exception {
        long deadline = System.nanoTime() + budget.toNanos();
        List<T> results = new ArrayList<>();
        try {
            for (Future<T> future : futures) {
                long remaining = deadline - System.nanoTime();
                if (remaining <= 0) throw new TimeoutException("batch deadline");
                results.add(future.get(remaining, TimeUnit.NANOSECONDS));
            }
            return List.copyOf(results);
        } catch (Exception e) {
            futures.forEach(future -> future.cancel(true));
            throw e;
        }
    }

    public static void main(String[] args) throws Exception {
        var executor = java.util.concurrent.Executors.newFixedThreadPool(2);
        try {
            var futures = List.of(executor.submit(() -> 7), executor.submit(() -> 8));
            System.out.println(collect(futures, Duration.ofSeconds(1)));
        } finally {
            executor.shutdownNow();
        }
    }
}
잔여 예산 검사와 취소 요청의 경계를 구분한다

DeadlineFutureCollector.collect의 기한 검사, 결과 추가, 목록 반환, 예외 후 취소 시도와 재전파를 보여줍니다.

잔여 예산 검사와 취소 요청의 경계를 구분한다
코드 지점원문이 하는 일반환·취소의 경계
remaining ≤ 0TimeoutException("batch deadline") 발생이미 완료한 Future라도 이 검사에서 중단 가능
future.get(remaining, NANOSECONDS)정상 결과를 results에 추가아직 호출자에게 부분 목록을 반환하지 않음
모든 get 뒤 List.copyOf(results)결과 목록의 복사본 반환원소 객체를 깊게 복사하지 않으며 null 원소는 거부
Exception을 잡은 catch모든 Future에 cancel(true) 시도 후 예외 재전파취소 성공 여부를 검사하지 않고 부분 결과도 반환하지 않음
remaining ≤ 0
원문이 하는 일: TimeoutException("batch deadline") 발생
반환·취소의 경계: 이미 완료한 Future라도 이 검사에서 중단 가능
future.get(remaining, NANOSECONDS)
원문이 하는 일: 정상 결과를 results에 추가
반환·취소의 경계: 아직 호출자에게 부분 목록을 반환하지 않음
모든 get 뒤 List.copyOf(results)
원문이 하는 일: 결과 목록의 복사본 반환
반환·취소의 경계: 원소 객체를 깊게 복사하지 않으며 null 원소는 거부
Exception을 잡은 catch
원문이 하는 일: 모든 Future에 cancel(true) 시도 후 예외 재전파
반환·취소의 경계: 취소 성공 여부를 검사하지 않고 부분 결과도 반환하지 않음

한 번의 실행에서는 [7, 8]을 반환했으며 시간 초과·취소 경로는 관측하지 않았습니다. 예산은 collect 진입 때 시작해 앞선 제출 시간은 제외합니다. 잔여 get 예산을 재사용해도 스케줄링·취소·실행기 종료를 포함한 전체 반환 시각의 상한은 보장하지 않습니다.


결과 수집 방식 비교

결과 요구수집 방식특성
제출 순서 결과Future 목록 순회느린 앞 작업이 뒤 결과 지연
완료 순 결과CompletionService빠른 결과 먼저 처리
단계 조합CompletableFuture비동기 변환 연결
수집 대기 예산기한 get각 호출에 잔여 예산 전달

분할 수를 정하는 근거

작업 조각이 너무 크면 가장 느린 조각이 전체 완료를 늦추고, 너무 작으면 제출·큐잉·결과 결합 비용이 실제 계산보다 커집니다.

CPU 계산은 사용 가능한 프로세서 수를 출발점으로 삼고, 각 조각의 입력 분포가 비슷한지 측정합니다.

외부 I/O가 섞인다면 프로세서 수만으로 병렬도를 정하지 말고 대상 시스템의 연결 상한과 허용 요청률을 먼저 봅니다.

병렬 합산의 정확성도 성능과 함께 확인합니다.

int 범위를 넘을 수 있으면 부분합과 최종합에 long을 검토하되, long의 범위도 넘는 입력은 별도로 처리해야 합니다. 부동소수점은 결합 순서가 바뀌면 마지막 비트가 달라질 수 있음을 허용 오차에 반영합니다.

기준 순차 구현과 무작위 입력을 비교해 값이 맞는지 확인한 뒤에만 처리 시간을 비교하세요.

입력이 작을 때는 순차 경로를 선택하는 임계값도 둡니다.

측정에서 병렬 준비 비용이 계산 절감보다 큰 구간을 찾아 기준을 정하면 간단한 요청이 큐와 스레드를 불필요하게 거치지 않습니다.


연습 문제

1부터 1,000까지를 두 범위로 나눠 홀수와 짝수 개수를 각각 Callable로 계산하고, 두 작업을 모두 제출한 뒤 결과를 합치세요.

정답과 해설
exercise/ParallelParitySolution.java
import java.util.concurrent.Executors;

public final class ParallelParitySolution {
    record Counts(int odd, int even) {
        Counts plus(Counts other) { return new Counts(odd + other.odd, even + other.even); }
    }

    static Counts count(int start, int end) {
        int odd = 0, even = 0;
        for (int n = start; n <= end; n++) {
            if ((n & 1) == 0) {
                even++;
            } else {
                odd++;
            }
        }
        return new Counts(odd, even);
    }

    public static void main(String[] args) throws Exception {
        var executor = Executors.newFixedThreadPool(2);
        try {
            var left = executor.submit(() -> count(1, 500));
            var right = executor.submit(() -> count(501, 1_000));
            System.out.println(left.get().plus(right.get()));
        } finally {
            executor.shutdown();
        }
    }
}

결과는 odd 500, even 500입니다.

두 submit이 두 get보다 앞에 있어 두 계산이 겹칠 수 있습니다. 이 main은 각 계산의 시작·종료 시각을 기록하지 않습니다.


병렬 합산을 채택하는 기준

Future는 작업을 대표할 뿐 병렬 구조를 자동 구성하지 않습니다.

독립 작업을 먼저 모두 제출하고 수집 지점의 블로킹과 시간 상한을 명시해야 실제 병렬성과 응답 규칙이 함께 살아납니다.