从Couchbase SDK2.7迁移至3.2 如何使用Flux实现等效查询逻辑
Couchbase SDK 3.2 迁移:Flux实现旧版Observable批量查询等效逻辑
旧版基于SDK 2.x写的Observable+JsonDocument批量ID查询逻辑,在3.2版本中可以直接按照以下规则替换为Reactor Flux实现:
- 旧版
bucket.async()异步入口替换为3.x的bucket.reactive(),返回基于Reactor的响应式操作实例 - 旧版
JsonDocument类在3.x中移除,单key get操作返回GetResult对象,可通过该对象的contentAsObject()方法获取原始JSON结构,和旧版JsonDocument携带的JSON内容完全对齐 - RxJava 1的操作符和Reactor Flux操作符基本一一对应:
Observable.from替换为Flux.fromIterable,toList().toBlocking().single()替换为collectList().block()即可
基础等效实现代码
import com.couchbase.client.java.kv.GetResult; import reactor.core.publisher.Flux; import java.util.List; // 提前完成Couchbase集群、桶初始化,拿到reactiveBucket响应式操作实例 List<GetResult> foundDocs = Flux.fromIterable(campaignIdList) .flatMap(id -> reactiveBucket.get(id)) .collectList() .block();
对齐旧版JsonDocument取JSON内容的实现
如果需要直接拿到和旧版JsonDocument中一致的JSON对象,加一层map转换即可:
import com.couchbase.client.java.json.JsonObject; import com.couchbase.client.java.kv.GetResult; import reactor.core.publisher.Flux; import java.util.List; List<JsonObject> foundJsonContents = Flux.fromIterable(campaignIdList) .flatMap(id -> reactiveBucket.get(id)) .map(GetResult::contentAsObject) .collectList() .block();
优化建议
- 如果需要控制批量请求的并发度,避免瞬时请求过多打满连接,可以给
flatMap传入第二个并发参数,例如flatMap(id -> reactiveBucket.get(id), 16)即可限制最大并发为16 - 批量ID查询场景优先使用3.x原生提供的
getAll接口,比循环单key get性能更高,示例:
List<GetResult> batchFoundDocs = reactiveBucket.getAll(campaignIdList) .map(entry -> entry.getValue()) .collectList() .block();
内容的提问来源于stack exchange,提问作者rishabh katyal
相关产品推荐
相关产品推荐

