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

Reactor定时任务引发Native Thread OOM问题求解

问题描述

我刚接触Reactor,若表述不清还请见谅。我编写了一个类,通过Flux.interval(Duration.ofMinutes(5))实现每5分钟执行一次重复任务:SSH连接本地Linux机器执行bash脚本,记录输出后存入数据库。初期运行正常,但一段时间后出现java.lang.OutOfMemoryError: unable to create native thread: possibly out of memory or process/resource limits reached错误。

该Java程序运行在Docker容器中,查看日志可见io-executor-thread-<some-thread-number>,当线程编号达到约250时触发OOM。相关代码如下:

@PostConstruct
private void pollMachines() {
    Flux.interval(Duration.ofMinutes(5))
            .map(this::getAllMachines)
            .flatMap(name -> process1(name)
                    .map(this::process2)                                                
                    .doOnError(throwable -> LOGGER.info("Some error happened with {}", name))
                    .onErrorResume(throwable -> setOfflineStatus(name))
            )
            .subscribeOn(Schedulers.boundedElastic())
            .subscribe();
}

异常堆栈信息如下:

07:03:49.854 [io-executor-thread-256] ERROR reactor.core.publisher.Operators - Operator called default onErrorDropped
backend     | reactor.core.Exceptions$ErrorCallbackNotImplemented: java.lang.OutOfMemoryError: unable to create native thread: possibly out of memory or process/resource limits reached
backend     | Caused by: java.lang.OutOfMemoryError: unable to create native thread: possibly out of memory or process/resource limits reached
backend     |   at java.base/java.lang.Thread.start0(Native Method)
backend     |   at java.base/java.lang.Thread.start(Thread.java:798)
backend     |   at java.base/java.util.concurrent.ThreadPoolExecutor.addWorker(ThreadPoolExecutor.java:937)
backend     |   at java.base/java.util.concurrent.ThreadPoolExecutor.execute(ThreadPoolExecutor.java:1354)
backend     |   at io.micronaut.scheduling.instrument.InstrumentedExecutor.execute(InstrumentedExecutor.java:42)
backend     |   at java.base/java.util.concurrent.CompletableFuture.asyncSupplyStage(CompletableFuture.java:1714)
backend     |   at java.base/java.util.concurrent.CompletableFuture.supplyAsync(CompletableFuture.java:1931)
backend     |   at io.micronaut.data.runtime.operations.ExecutorAsyncOperations.update(ExecutorAsyncOperations.java:152)
backend     |   at io.micronaut.data.runtime.operations.ExecutorReactiveOperations.lambda$update$12(ExecutorReactiveOperations.java:194)
backend     |   at io.micronaut.core.async.publisher.CompletableFuturePublisher$CompletableFutureSubscription.request(CompletableFuturePublisher.java:78)
backend     |   at reactor.core.publisher.MonoNext$NextSubscriber.request(MonoNext.java:108)
backend     |   at reactor.core.publisher.MonoNext$NextSubscriber.request(MonoNext.java:108)
backend     |   at reactor.core.publisher.MonoFlatMap$FlatMapInner.onSubscribe(MonoFlatMap.java:238)
backend     |   at reactor.core.publisher.MonoNext$NextSubscriber.onSubscribe(MonoNext.java:70)
backend     |   at reactor.core.publisher.MonoNext$NextSubscriber.onSubscribe(MonoNext.java:70)
backend     |   at io.micronaut.core.async.publisher.CompletableFuturePublisher.subscribe(CompletableFuturePublisher.java:49)
backend     |   at reactor.core.publisher.MonoFromPublisher.subscribe(MonoFromPublisher.java:63)
backend     |   at reactor.core.publisher.InternalMonoOperator.subscribe(InternalMonoOperator.java:64)
backend     |   at reactor.core.publisher.MonoFlatMap$FlatMapMain.onNext(MonoFlatMap.java:157)
backend     |   at reactor.core.publisher.FluxMap$MapSubscriber.onNext(FluxMap.java:122)
backend     |   at reactor.core.publisher.MonoNext$NextSubscriber.onNext(MonoNext.java:82)
backend     |   at reactor.core.publisher.MonoNext$NextSubscriber.onNext(MonoNext.java:82)
backend     |   at io.micronaut.core.async.publisher.Publishers$1.doOnNext(Publishers.java:248)
backend     |   at io.micronaut.core.async.subscriber.CompletionAwareSubscriber.onNext(CompletionAwareSubscriber.java:56)
backend     |   at io.micronaut.core.async.publisher.CompletableFuturePublisher$CompletableFutureSubscription.lambda$request$0(CompletableFuturePublisher.java:89)
backend     |   at java.base/java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:859)
backend     |   at java.base/java.util.concurrent.CompletableFuture$UniWhenComplete.tryFire(CompletableFuture.java:837)
backend     |   at java.base/java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:506)
backend     |   at java.base/java.util.concurrent.CompletableFuture$AsyncSupply.run(CompletableFuture.java:1705)
backend     |   at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
backend     |   at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
backend     |   at java.base/java.lang.Thread.run(Thread.java:829)

补充说明:每次轮询完20台左右机器后,下一轮轮询会启动新线程。请问如何避免该OOM异常,实现稳定的每5分钟轮询任务?

解决方案

问题根源

  1. flatMap无界并发:默认flatMap不限制并发数,每次轮询20台机器会同时启动20个任务,boundedElastic调度器为每个阻塞IO任务(SSH、数据库操作)创建新线程,长期运行线程池持续膨胀,触达系统线程上限引发OOM。
  2. interval未等待上一轮任务完成:Flux.interval严格按5分钟间隔发射事件,不管上一轮任务是否执行完毕,导致任务堆积,线程数不断增加。
  3. 异步操作线程泄漏:从堆栈看,数据库操作使用CompletableFuture.supplyAsync,默认依赖Micronaut的io-executor线程池,若任务未正确收尾,线程无法回收。

修复步骤

1. 限制flatMap并发数

给flatMap设置合理的并发上限,避免一次性创建过多线程:

.flatMap(name -> process1(name)
        .map(this::process2)
        .doOnError(throwable -> LOGGER.info("Some error happened with {}", name))
        .onErrorResume(throwable -> setOfflineStatus(name)),
    5) // 限制并发数为5,可根据系统资源调整

2. 确保上一轮任务完成后再触发下一轮

改用repeat+delayElement组合,替代interval,保证一轮任务全部执行完毕后再等待5分钟启动下一轮:

@PostConstruct
private void pollMachines() {
    Flux.defer(() -> Flux.fromIterable(getAllMachines()))
            .flatMap(name -> process1(name)
                    .map(this::process2)
                    .doOnError(throwable -> LOGGER.info("Some error happened with {}", name))
                    .onErrorResume(throwable -> setOfflineStatus(name)),
                5)
            .then() // 等待当前轮所有任务完成
            .delayElement(Duration.ofMinutes(5)) // 等待5分钟
            .repeat() // 重复执行
            .subscribeOn(Schedulers.boundedElastic())
            .subscribe(
                null,
                error -> LOGGER.error("Polling task failed", error) // 添加全局错误处理
            );
}

3. 优化阻塞操作的线程管理

  • 显式将阻塞IO任务绑定到boundedElastic调度器,避免Micronaut默认线程池膨胀:
private Mono<String> process1(String name) {
    return Mono.fromCallable(() -> {
        // SSH连接执行脚本的阻塞逻辑
        return executeSshScript(name);
    }).subscribeOn(Schedulers.boundedElastic());
}
  • 配置Micronaut的io-executor线程池参数(application.yml):
micronaut:
  executors:
    io:
      type: scheduled
      core-pool-size: 10
      max-pool-size: 20
      keep-alive-minutes: 5

4. 避免资源泄漏

  • 检查process1、process2、setOfflineStatus中的资源(SSH连接、数据库连接)是否正确关闭,未关闭的资源会导致线程无法回收。

内容的提问来源于stack exchange,提问作者MKRabbit

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 07:45:32