ExecutorCompletionService — Executor 기반 완료 서비스 구현

ExecutorCompletionService — Executor 기반 완료 서비스 구현

ExecutorCompletionService<V>는 공급된 Executor를 사용해 태스크를 실행하는 CompletionService 구현이에요. 제출된 태스크가 완료되면 take로 접근 가능한 큐에 배치해요. 태스크 그룹을 처리할 때 일시적으로 쓰기 충분히 가벼워요.

출처: Java API Reference

본문

개념 이해하기

ExecutorCompletionService는 공급된 Executor로 태스크를 실행하는 CompletionService 구현이에요. 제출된 태스크가 완료되는 즉시 take로 접근 가능한 큐에 놓여, 완료된 순서대로 결과를 꺼낼 수 있어요.

public class ExecutorCompletionService<V>
extends Object
implements CompletionService<V>

활용 예 — 비동기 결과 처리

문제에 대한 여러 솔버(solver) 집합을 동시에 실행하고, non-null 값을 반환하는 각각의 결과를 use(Result r)로 처리한다면:

void solve(Executor e, Collection<Callable<Result>> solvers)
        throws InterruptedException, ExecutionException {
    CompletionService<Result> cs = new ExecutorCompletionService<>(e);
    solvers.forEach(cs::submit);
    for (int i = solvers.size(); i > 0; i--) {
        Result r = cs.take().get();
        if (r != null) use(r);
    }
}

활용 예 — 첫 non-null 결과 사용

태스크 집합의 첫 non-null 결과만 사용하고, 예외는 무시하며, 첫 결과가 준비되면 다른 모든 태스크를 취소한다면:

void solve(Executor e, Collection<Callable<Result>> solvers)
        throws InterruptedException {
    CompletionService<Result> cs = new ExecutorCompletionService<>(e);
    int n = solvers.size();
    List<Future<Result>> futures = new ArrayList<>(n);
    Result result = null;
    try {
        solvers.forEach(solver -> futures.add(cs.submit(solver)));
        for (int i = n; i > 0; i--) {
            try {
                Result r = cs.take().get();
                if (r != null) { result = r; break; }
            } catch (ExecutionException ignore) {}
        }
    } finally {
        futures.forEach(future -> future.cancel(true));
    }
    if (result != null) use(result);
}

생성자

public ExecutorCompletionService(Executor executor) — 기본 태스크 실행을 위한 공급 executor와 완료 큐로 LinkedBlockingQueue를 사용하는 서비스를 만들어요.

  • NullPointerExceptionexecutornull일 때

public ExecutorCompletionService(Executor executor, BlockingQueue<Future<V>> completionQueue) — 공급 executor와 공급 큐를 완료 큐로 사용하는 서비스를 만들어요. 큐는 무경계로 취급돼요 — 완료된 태스크에 대한 실패한 Queue.add는 그 태스크를 검색 불가능하게 만들어요.

  • NullPointerExceptionexecutor 또는 completionQueuenull일 때

메서드

public Future<V> submit(Callable<V> task) — 값을 반환하는 태스크를 제출하고 대기 중 결과를 나타내는 Future를 반환해요.

  • RejectedExecutionException, NullPointerException

public Future<V> submit(Runnable task, V result)Runnable 태스크를 제출하고 그 태스크를 나타내는 Future를 반환해요. 완료 시 get()result를 반환해요.

  • RejectedExecutionException, NullPointerException

public Future<V> take() throws InterruptedException — 다음으로 완료된 태스크의 Future를 꺼내 제거해요. 아직 없으면 기다려요.

  • InterruptedException

public Future<V> poll() — 다음으로 완료된 태스크의 Future를 꺼내거나, 없으면 null.

public Future<V> poll(long timeout, TimeUnit unit) throws InterruptedException — 최대 timeout까지 기다렸다가 다음 완료 태스크의 Future를 꺼내요. 시간이 지나면 null.

더 알아보기 (Learn more)