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分钟轮询任务?
解决方案
问题根源
flatMap无界并发:默认flatMap不限制并发数,每次轮询20台机器会同时启动20个任务,boundedElastic调度器为每个阻塞IO任务(SSH、数据库操作)创建新线程,长期运行线程池持续膨胀,触达系统线程上限引发OOM。interval未等待上一轮任务完成:Flux.interval严格按5分钟间隔发射事件,不管上一轮任务是否执行完毕,导致任务堆积,线程数不断增加。- 异步操作线程泄漏:从堆栈看,数据库操作使用
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
相关产品推荐
相关产品推荐

