如何基于Reactive MongoDB实现文档串行保存与处理?
问题描述
我有一批文档需先执行处理逻辑,再保存到MongoDB中。已引入Reactive MongoDB依赖,希望尽量避免使用非反应式依赖,要求第二个文档的处理必须在第一个文档保存完成后启动。
当前代码如下:
Mono<List<Index>> deferredCreate = Mono.defer(() -> index .flatMapMany(Flux::fromIterable) .flatMapSequential(entity -> { repository.process(entity).subscribe(); return entity; }) ) .collectList());
其中process方法实现为:
public Mono<IndexDocument> process(Index index) { if(someCondition) { return mongoOperations.save(index); } else { return mongoOperations.findAndReplace(query, index); } }
目前代码虽逐个处理文档,但第n+1个文档的处理会在第n个开始后启动,而非保存完成后。无法用.block()替代.subscribe(),否则并行流会报错;尝试过concatMap,效果相同。
请问是否有可行实现方式?因属于更大规模反应式流程的一部分,需使用反应式编程。
解决方案
你的核心问题是在flatMapSequential里直接调用repository.process(entity).subscribe(),这会让process的Mono脱离当前反应式链——当前流不会等待MongoDB操作完成,就直接返回entity触发下一个元素的处理。
正确的做法是让process的Mono成为反应式链的一部分,确保流严格等待每个保存操作完成后再处理下一个文档。可以用concatMap(或flatMapSequential,单线程场景下二者行为一致,都能保证顺序执行),但要正确把process操作嵌入链中。
修改后的代码如下:
Mono<List<Index>> deferredCreate = Mono.defer(() -> index .flatMapMany(Flux::fromIterable) // concatMap会严格按顺序等待前一个process完成后,再处理下一个元素 .concatMap(entity -> repository.process(entity) // 处理完成后映射回原entity,保证最终能收集到原始Index列表 .thenReturn(entity) ) .collectList());
关键说明:
- 移除
.subscribe(),改用.thenReturn(entity)将process的异步操作结果映射回原entity,这样流会等待MongoDB的保存/替换操作完全结束后,才会发射下一个元素。 - 如果不需要收集原始的Index列表,而是要收集处理后的
IndexDocument,可以直接去掉.thenReturn(entity),让concatMap返回repository.process(entity)即可。 - 你之前尝试
concatMap无效,大概率是因为当时的写法同样让process的Mono脱离了反应式链(比如误用了subscribe),现在的写法会确保每个步骤都在链内等待。
内容的提问来源于stack exchange,提问作者marie
相关产品推荐
相关产品推荐

