Java使用Callable实现异步调用出现OutOfMemoryError内存溢出问题咨询
问题根因
你遇到的OutOfMemoryError属于系统层面无法创建新线程的错误,错误码11对应Linux系统的EAGAIN错误,说明当前Java进程的线程数已经达到操作系统限制,无法再分配资源创建新线程。
代码存在的缺陷
- 线程数完全和业务列表
pccsSurvList的长度绑定,没有上限,当列表数据量突增时,会一次性创建远超系统承受能力的线程,直接触发创建线程失败的OOM - 线程池未复用,每次业务调用都新建独立的线程池,若该逻辑被并发调用,总线程数会成倍增长,进一步加大OOM风险
- 空指针风险:当
pccsSurvList为空时,executorService为null,后续直接调用invokeAll、shutdown方法会直接抛出空指针异常 - 线程池关闭逻辑不可靠:未做异常捕获和
finally兜底,若invokeAll执行过程中抛出中断等异常,shutdown方法不会执行,会导致线程永久泄漏,占用系统资源
修复方案
- 线程池参数固定上限,不要和业务数据量绑定:调用第三方webservice属于IO密集型任务,线程数建议设置为
2*CPU核心数,最大不要超过系统允许的单进程线程数上限(可通过ulimit -u查看),也可根据第三方接口的吞吐量调整。 - 线程池全局复用:将线程池设置为单例,避免每次业务调用反复创建销毁线程,也方便全局管控总线程数上限。
- 补全空判断逻辑:业务列表为空时直接返回,不执行后续线程池相关逻辑。
- 线程池操作加异常兜底:保证无论是否发生异常,线程池都能正常关闭,避免线程泄漏,可新增
awaitTermination设置超时等待时间,防止任务永久阻塞。 - 自定义线程池参数:不要用
Executors.newFixedThreadPool,手动构造ThreadPoolExecutor,明确指定队列长度、拒绝策略、线程工厂,避免隐式的无界队列等风险。 - 大数据量分批处理:若
pccsSurvList单次数据量超过1000条,建议拆分成多批提交,每批执行完成后再提交下一批,避免同时提交太多任务占满资源。
优化后代码示例
import com.google.common.util.concurrent.ThreadFactoryBuilder; import java.util.*; import java.util.concurrent.*; // 线程池单例,全局复用,建议放在统一的配置类中初始化 private static final ExecutorService WEBSERVICE_SYNC_POOL; static { int corePoolSize = Math.min(2 * Runtime.getRuntime().availableProcessors(), 50); int maximumPoolSize = 50; long keepAliveTime = 60L; // 队列长度根据业务可接受的最大等待任务数调整 BlockingQueue<Runnable> workQueue = new ArrayBlockingQueue<>(200); // 自定义线程名,方便问题排查,无guava依赖可自行实现ThreadFactory接口 ThreadFactory threadFactory = new ThreadFactoryBuilder().setNameFormat("webservice-sync-thread-%d").build(); // 拒绝策略可根据业务调整,CallerRunsPolicy表示队列满了由提交线程执行任务,避免丢失 RejectedExecutionHandler handler = new ThreadPoolExecutor.CallerRunsPolicy(); WEBSERVICE_SYNC_POOL = new ThreadPoolExecutor(corePoolSize, maximumPoolSize, keepAliveTime, TimeUnit.SECONDS, workQueue, threadFactory, handler); } public List<SyncFlagEntity> syncWebServiceData(List<PCCSSurveyInfoEntity> pccsSurvList) throws InterruptedException { // 空列表直接返回 if (pccsSurvList == null || pccsSurvList.isEmpty()) { return Collections.emptyList(); } List<Callable<SyncFlagEntity>> taskList = new ArrayList<>(pccsSurvList.size()); for (PCCSSurveyInfoEntity user : pccsSurvList) { taskList.add(new NotifyEmailTransactionImpl(user)); } List<Future<SyncFlagEntity>> futureList = WEBSERVICE_SYNC_POOL.invokeAll(taskList); // 处理结果,捕获执行异常 List<SyncFlagEntity> resultList = new ArrayList<>(); for (Future<SyncFlagEntity> future : futureList) { try { resultList.add(future.get()); } catch (ExecutionException e) { // 单条任务执行异常处理,根据业务需要调整 throw new RuntimeException("同步任务执行失败", e); } } return resultList; } // 应用关闭时调用该方法销毁线程池 public void shutdownPool() { WEBSERVICE_SYNC_POOL.shutdown(); try { if (!WEBSERVICE_SYNC_POOL.awaitTermination(1, TimeUnit.MINUTES)) { WEBSERVICE_SYNC_POOL.shutdownNow(); } } catch (InterruptedException e) { WEBSERVICE_SYNC_POOL.shutdownNow(); Thread.currentThread().interrupt(); } }
内容的提问来源于stack exchange,提问作者absolute
相关产品推荐
相关产品推荐

