ExecutorCompletionService — Executor 기반 완료 서비스 구현
ExecutorCompletionService — Executor 기반 완료 서비스 구현
ExecutorCompletionService<V>는 공급된 Executor를 사용해 태스크를 실행하는 CompletionService 구현이에요. 제출된 태스크가 완료되면 take로 접근 가능한 큐에 배치해요. 태스크 그룹을 처리할 때 일시적으로 쓰기 충분히 가벼워요.
본문
개념 이해하기
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를 사용하는 서비스를 만들어요.
NullPointerException—executor가null일 때
public ExecutorCompletionService(Executor executor, BlockingQueue<Future<V>> completionQueue) — 공급 executor와 공급 큐를 완료 큐로 사용하는 서비스를 만들어요. 큐는 무경계로 취급돼요 — 완료된 태스크에 대한 실패한 Queue.add는 그 태스크를 검색 불가능하게 만들어요.
NullPointerException—executor또는completionQueue가null일 때
메서드
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.