Flux#toIterable未按文档实现懒加载引发内存溢出问题咨询
问题描述
我正在使用Reactor完成一个学校项目,遇到了Flux相关的问题。
上游有一个从数据库读取数据、生成数据行供下游处理器使用的Flux,示例代码如下:
public Flux<Row> emitRow(...) { return Flux.create(cursor -> { // ... 数据库读取与校验逻辑 ... emitRow(cursor, row); cursor.complete() }); }
由于沙箱环境内存有限,将所有数据读入内存并不可行,因此我最初的设计是每读取一行就先处理该行,再获取下一行。
最初的实现如下,运行正常,但指导老师要求我们改用toIterable()方法:
public Mono<Outcome> process(Flux<Row> data) { return data .flatMap(row -> processRow(row), SINGLE_THREAD, SINGLE_THREAD) .map(results -> new Outcome(results)); }
现在我尝试适配toIterable()时遇到了问题,代码运行效果和JavaDoc描述不符:当前实现似乎涉及两个队列,第一个是Iterable迭代器的阻塞队列,第二个是原始Flux.create创建的队列,而Flux.create生成的第二个队列似乎会将所有数据缓冲到内存中,导致应用抛出内存溢出异常。我的实现代码如下:
public Mono<Outcome> process(Flux<Row> data) { Iterator<Row> itr = data.toIterable(1).iterator(); return Mono.fromCallable(() -> processRow(itr)) .map(results -> new Outcome(results)); } public Results processRow(Iterator<Row> itr) { while(itr.hasNext()) // <--- 此处预期为阻塞调用 { Row r = itr.next(); dbContentBuilder.handle(r); } return new Results(dbContentBuilder.build()); }
请问为什么会出现该问题?如何调整实现才能在使用Iterable的前提下,避免Flux.create将所有数据缓冲到同步队列中,而是通过iterator.next()主动请求下一块数据?官方文档显示该方法应当是懒加载队列,调用.next()时会阻塞。
回答
问题原因
内存溢出的核心原因是当前Flux.create的实现没有遵守Reactor的背压规则:
- 默认情况下
Flux.create使用OverflowStrategy.BUFFER作为溢出策略,会无视下游的请求量,缓冲所有上游发射的元素直到内存耗尽。 toIterable(1)的逻辑本身符合预期:每次调用next()时仅向上游请求1个元素,但你的Flux.create内部逻辑是一次性把所有数据库数据全部读出来发射,根本不响应下游的请求信号,所有多出来的元素都被缓冲在Flux.create的内部队列里,最终导致OOM。
修复方案
你只需要调整emitRow方法中Flux.create的实现,让它响应下游的拉取请求,按需读取数据库数据即可:
public Flux<Row> emitRow(...) { return Flux.create(sink -> { // 注册请求监听器,只有当下游请求数据时才读取对应行数的数据库内容 sink.onRequest(requested -> { long emitted = 0; while(emitted < requested && 还有未读的数据库数据) { Row row = 读取下一行数据库数据; sink.next(row); emitted++; } if(没有未读数据) { sink.complete(); } }); // 因为我们严格按下游请求量发射数据,不会出现溢出,所以可以用IGNORE策略避免额外的缓冲开销 }, FluxSink.OverflowStrategy.IGNORE); }
修改完成后你原有的toIterable(1)逻辑无需调整,迭代器的hasNext()/next()会正常阻塞请求下一行数据,不会出现多余的内存占用。
如果暂时不方便修改上游emitRow的实现,也可以临时在toIterable之前添加背压控制算子,避免上游无限制发射:
Iterator<Row> itr = data.onBackpressureLatest().toIterable(1).iterator();
注意该方案只是临时规避,还是建议优先修改上游Flux.create的实现遵守背压规则,才能从根本上解决内存溢出问题。
内容的提问来源于stack exchange,提问作者user3342825
相关产品推荐
相关产品推荐

