如何基于Spring Data Reactive MongoDB无阻塞推送Flux元素至MongoDB数组
解决ReactiveMongoOperations插入Flux元素到Mongo数组的非阻塞方案
我明白你的问题了——直接把Flux传给Update.push().each()时,Spring Data会把Flux对象本身序列化存进数组,而不是提取里面的元素。这是因为each()方法并不支持直接接收反应式流(比如Flux)作为参数,它只会把传入的对象当成单个元素处理。
要在不阻塞的情况下实现需求,我们可以利用反应式操作的特性,先异步收集Flux中的元素,再执行upsert操作:
val mongo: ReactiveMongoOperations = ... data class Something(val data: String) val flux = Flux.just(Something("A"), Something("B")) // 实际为动态生成的Flux // 异步收集Flux元素为List,再执行upsert flux.collectList() .flatMap { items -> mongo.upsert( query(where("_id").isEqualTo("myId")), Update().push("myArray").each(items), "collection" ) } // 根据你的业务场景处理订阅,比如绑定到WebFlux的响应链 .subscribe()
为什么这个方案可行?
collectList()是反应式操作,它会异步收集Flux中的所有元素到一个List里,全程不会阻塞线程,符合反应式编程的非阻塞要求。- 在
flatMap中拿到收集好的List后,再传给each(),此时each()会把List中的每个Something对象逐个添加到Mongo的数组字段中,得到你期望的结果。
原代码问题的根源
Update.push().each()的重载参数是Object... values或者Collection<?> values,并没有专门处理Publisher(Flux/Mono)的重载。所以当你直接传入Flux时,它会把Flux作为一个普通Java对象进行序列化,最终存进Mongo的是Flux的内部结构(比如示例中的FluxArray),而不是流中的元素。
注意事项
如果你的Flux是无限流,collectList()会一直等待元素结束,这种情况下你需要用take(n)之类的操作符限制收集的元素数量,确保操作能正常完成。
内容的提问来源于stack exchange,提问作者Johan
相关产品推荐
相关产品推荐

