Observable.blockingForEach触发无法创建本地线程OutOfMemoryError求助
问题根本原因
核心问题是代码中存在线程池重复创建的资源泄露:
- 每次进入
flatMap的处理逻辑时,都会执行Schedulers.from(Executors.newFixedThreadPool(4)),意味着每处理一个key就会新建一个独立的4线程固定线程池 Executors.newFixedThreadPool默认配置核心线程永不超时销毁,也没有任何代码主动关闭这些临时创建的线程池,线程会一直存活占用系统资源- 每调用一次该接口,就会新增20~200个永久存活线程,随着请求量累计,线程数很快超过操作系统允许的进程最大线程数,最终抛出
unable to create new native thread的OOM错误,和异常栈、线程转储中存在大量高编号pool-xxx闲置线程的特征完全吻合 - 额外的叠加问题:从线程转储可以看出,部分线程卡在HTTP请求的等待、读取阶段,说明外部API调用没有配置合理超时,进一步加剧了线程资源的占用。
解决方案
核心修复:全局复用调度器
不要每次请求/每次处理key都新建线程池,在类的成员变量位置仅初始化一次调度器,所有请求共用:
// 类成员位置定义,全局唯一,复用所有请求的线程资源 private val externalApiScheduler = run { val executor = Executors.newFixedThreadPool(4) as ThreadPoolExecutor // 可选配置:允许核心线程超时闲置回收,进一步降低闲时资源占用 executor.allowCoreThreadTimeOut(true) executor.setKeepAliveTime(60, java.util.concurrent.TimeUnit.SECONDS) Schedulers.from(executor) } // 业务逻辑修改后 Observable.fromIterable(listOfKeys) .flatMap({ key -> Maybe.fromCallable { Pair(key, methodWhichCallsExternalAPI()) } .toObservable() .subscribeOn(externalApiScheduler) }, 4) // 第二个参数限制flatMap最大并发数,避免同时请求量过高打垮下游 .blockingForEach { // 保留原有结果赋值逻辑 }
补充修复:添加HTTP请求超时
给methodWhichCallsExternalAPI使用的HTTP客户端配置合理的超时规则,比如连接超时5秒、读取超时10秒,避免请求长时间挂起占住线程资源。
可选优化
如果该业务没有强依赖ReactiveX的需求,这么小的数据量可以直接用普通线程池+循环处理,逻辑更简单也更容易排查问题。
内容的提问来源于stack exchange,提问作者yingxuan
相关产品推荐
相关产品推荐

