

使用ExecutorCompletionService出现OOM的场景
使用java.util.concurrent.ExecutorCompletionService异步处理任务:
java.util.concurrent.ExecutorCompletionService#submit(java.util.concurrent.Callable<V>)java.util.concurrent.ExecutorCompletionService#submit(java.lang.Runnable, V)而没有使用方法:
java.util.concurrent.ExecutorCompletionService#take或
java.util.concurrent.ExecutorCompletionService#poll()对提交的所有任务获取结果,导致内存泄露,发生OOM。
使用ExecutorCompletionService为什么会出现OOM
ExecutorCompletionService 使用我们自定义的线程池去异步执行任务,任务执行完,会把任务执行的结果java.util.concurrent.Future 缓存到队列 BlockingQueue 中:
/**
* Creates an ExecutorCompletionService using the supplied
* executor for base task execution and a
* {@link LinkedBlockingQueue} as a completion queue.
*
* @param executor the executor to use
* @throws NullPointerException if executor is {@code null}
*/
public ExecutorCompletionService(Executor executor) {
if (executor == null)
throw new NullPointerException();
this.executor = executor;
this.aes = (executor instanceof AbstractExecutorService) ?
(AbstractExecutorService) executor : null;
this.completionQueue = new LinkedBlockingQueue<Future<V>>();
}默认缓存任务执行结果的队列为队列不受限的 LinkedBlockingQueue:
this.completionQueue = new LinkedBlockingQueue<Future<V>>()当我们提交任务的时候,任务会封装为java.util.concurrent.ExecutorCompletionService.QueueingFuture:

java.util.concurrent.ExecutorCompletionService.QueueingFuture覆写了方法:
java.util.concurrent.FutureTask#done /**
* Protected method invoked when this task transitions to state
* {@code isDone} (whether normally or via cancellation). The
* default implementation does nothing. Subclasses may override
* this method to invoke completion callbacks or perform
* bookkeeping. Note that you can query status inside the
* implementation of this method to determine whether this task
* has been cancelled.
*/
protected void done() { }当任务执行完,会把结果缓存到队列中:

既然任务执行结果缓存到队列中,为了不让队列出现内存泄露,我们必须在任务执行结束后,从队列中移除任务执行结果,所以ExecutorCompletionService 为我们提供了两对方法完成此操作:
/**
* @throws RejectedExecutionException {@inheritDoc}
* @throws NullPointerException {@inheritDoc}
*/
public Future<V> submit(Runnable task, V result) {
if (task == null) throw new NullPointerException();
RunnableFuture<V> f = newTaskFor(task, result);
executor.execute(new QueueingFuture<V>(f, completionQueue));
return f;
}
public Future<V> take() throws InterruptedException {
return completionQueue.take();
}
public Future<V> poll() {
return completionQueue.poll();
}
public Future<V> poll(long timeout, TimeUnit unit)
throws InterruptedException {
return completionQueue.poll(timeout, unit);
}如果我们不调用上述两对方法,任务执行的结果一值缓存在队列中,发送内存泄露,终有一刻会发生OOM异常。
使用ExecutorCompletionService的正确姿势
案例:对批量job即solvers异步处理后,一定要获取执行结果,做其它业务处理,
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) {
//do something
};
}
}案例来自javadoc
来自javadoc的另一个案例:获取批量job第一个先执行完的结果,其余job执行cancel处理:
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);
}但我感觉这个可能会发生内存泄露风险,因为第一个job执行完,从结果队列里移除,此时其他job在执行cance之前,也可能会执行完job,会把结果缓存到队列中,而QueueingFuture没有复写cancel方法,判断此时任务是否执行完,是否已经把结果缓存到队列中,是否需要从队列中删除。
小结
使用ExecutorCompletionService处理任务,一定记得执行:
java.util.concurrent.ExecutorCompletionService#take或
java.util.concurrent.ExecutorCompletionService#poll()方法,对提交的所有任务获取结果,防止任务结果缓存队列内存泄漏!限制在本地局部变量使用!也可预防!。
建议:不要使用ExecutorCompletionService,从javadoc上,这个类的实现并不是Doug Lea的作品。