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

Quarkus Mutiny Multi/Uni内存泄漏与Vertx线程阻塞问题求助

问题分析与解决思路

你的问题核心是两个:Java堆内存溢出(OOM)和Vert.x EventLoop线程阻塞,下面针对这两个问题给出具体的排查和解决方向:

一、解决内存溢出(OOM)问题

1. 根源:一次性加载百万级数据到内存

当前getAllMatchingItem方法一次性把所有匹配数据加载成List<Item>,百万级数据会直接占用大量堆空间,后续的分块操作只是对这个大List做切片,原List始终在内存中无法被GC回收,这是OOM的主要原因。

解决方法:改用流式查询+流式分块

不要一次性拉取全量数据,使用Panache的stream()方法实现流式查询,再通过Mutiny的Multi.group().intoLists()做流式分块,每次只加载一个批次的数据到内存:

修改getAllMatchingItem为返回流式Multi:

public Multi<Item> getAllMatchingItem(Request request) {
    return Multi.createFrom().item(() -> repo.find("...").stream())
            .onItem().transformToMulti(stream -> Multi.createFrom().iterable(stream))
            .runSubscriptionOn(Infrastructure.getDefaultWorkerPool());
}

核心处理函数去掉全量List加载:

public Uni<Void> processData(Request request) {
    return getAllMatchingItem(request)
            .group().intoLists().of(10000) // 流式分批次,仅保留当前批次数据在内存
            .onItem().transformToUni(resultBatched -> 
                // 移除AtomicReference,通过Tuple传递批次数据和关联数据
                Uni.combine().all().unis(
                        getOtherItemBasedOnItem(resultBatched),
                        Uni.createFrom().item(resultBatched)
                ).asTuple()
                .chain(tuple -> createProcessedItems(tuple.getItem2(), tuple.getItem1()))
                .chain(this::insertProcessedItems)
                .onTermination().invoke(() -> LOGGER.info("Data process success"))
                .replaceWithVoid()
            )
            .merge(1) // 串行处理批次,避免并发导致内存激增
            .collect().asList()
            .replaceWithVoid();
}

2. 移除不必要的AtomicReference

原代码中用AtomicReference传递批次数据是完全多余的,会导致批次数据被额外持有,延迟GC。通过Uni.combine().all().unis()把批次数据和关联查询结果打包成Tuple,就能在后续链式调用中直接获取,避免不必要的对象引用。

二、解决Vert.x EventLoop线程阻塞问题

1. 根源:同步阻塞操作侵入EventLoop线程

虽然你在辅助方法中用了runSubscriptionOn指定worker线程,但存在两个潜在问题:

  • 用Uni.createFrom().item(() -> repo.find(...).list)包裹Panache同步查询,虽然runSubscriptionOn会让lambda在worker线程执行,但如果后续用emitOn切回EventLoop,若有未处理的阻塞操作会卡住EventLoop;
  • 不必要的emitOn(MutinyHelper.executor(Vertx.currentContext()))会强制把结果切换回EventLoop,而后续操作如果不需要EventLoop(比如数据库写入、数据处理),会增加线程切换开销,甚至导致阻塞。

解决方法:

(1)优先使用Panache异步API

替换所有同步的repo.find(...).list为异步的listAsync(),避免手动包裹同步方法,更贴合Mutiny的异步模型:

public Uni<List<OtherItem>> getOtherItemBasedOnItem(List<Item> items) {
    return repo.find("id in ?1", items.getIDSet()).listAsync()
            .runSubscriptionOn(Infrastructure.getDefaultWorkerPool());
}

(2)移除不必要的emitOn

如果后续操作不需要在EventLoop线程执行(比如数据处理、数据库写入),直接去掉emitOn(MutinyHelper.executor(Vertx.currentContext())),减少线程切换,避免EventLoop被阻塞。

(3)确保耗时操作都在worker线程执行

对于doCreateProcessedItems这类CPU密集型或耗时的本地计算,确保createProcessedItems的runSubscriptionOn生效,且不要切回EventLoop:

public Uni<List<ProcessedItem>> createProcessedItems(List<Item> items, List<OtherItem> otherItems) {
    return Uni.createFrom().item(() -> doCreateProcessedItems(items, otherItems))
            .runSubscriptionOn(Infrastructure.getDefaultWorkerPool());
    // 移除不必要的emitOn
}

三、其他优化建议

  • 调整批次大小:如果10000条数据仍占用过多内存,可以尝试减小批次(比如5000),降低单批次内存占用;
  • 使用Panache批量写入:insertProcessedItems中的repo.persist(processedItems)如果是同步操作,改用persistAsync(processedItems),并确保在worker线程执行;
  • 监控内存使用:通过JVM参数(如-Xmx)调整堆大小,同时用JProfiler等工具排查是否存在内存泄漏点。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 14:00:56