You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.24 08:45:02