RxJava并行查询本地DB与API:需先展示本地结果并去重
解决方案:拆分流并结合RxJava算子满足需求
你之前用CombineLatest的问题在于,它必须等两个Observable都发射数据才会输出结果,没法单独先展示本地查询结果。咱们可以通过共享本地Observable+拆分两个订阅流的方式来实现所有需求,具体思路如下:
- 共享本地查询的Observable,避免重复执行本地数据库查询
- 单独订阅本地流,一有结果就立即展示给用户
- 处理API流时,先捕获异常保证不中断流程,再和本地流合并完成去重,确保API结果必须等本地结果就绪后才处理
完整实现代码
public void performSearchAsync(String query) { // 1. 创建本地查询Observable,用publish+autoConnect共享订阅,避免多次执行本地查询 Observable<List<LocalSearchResult>> localObservable = Observable.just(performLocalSearch(query)) .subscribeOn(Schedulers.io()) .publish() .autoConnect(2); // 当有2个订阅者时才触发执行 // 2. 第一个订阅:本地结果就绪后立即展示 localObservable .observeOn(AndroidSchedulers.mainThread()) .subscribe(localResults -> { if (localResults != null && !localResults.isEmpty()) { applySearchResult(localResults); } }, throwable -> { // 可选:处理本地查询异常,比如Toast提示 }); // 3. 处理API流:捕获异常+和本地结果合并去重 Observable<List<ApiSearchResult>> apiProcessedObservable = SearchApi.INSTANCE.getSearchResults(query) .subscribeOn(Schedulers.io()) // 捕获API异常,返回空的成功响应,避免流中断 .onErrorReturn(throwable -> Response.success(Collections.emptyList())) // 合并本地结果,确保API结果必须等本地结果就绪后才处理 .combineLatest(localObservable, (apiResponse, localResults) -> { // 优化去重逻辑:用HashSet提高查询效率 Set<String> localResultNames = localResults.stream() .map(LocalSearchResult::getName) .collect(Collectors.toSet()); if (!apiResponse.isSuccessful()) { return Collections.emptyList(); } // 去重:过滤掉本地已存在的结果 return apiResponse.body().stream() .filter(apiResult -> !localResultNames.contains(apiResult.getName())) .collect(Collectors.toList()); }) .filter(results -> !results.isEmpty()) // 过滤空结果,避免无意义的UI更新 .observeOn(AndroidSchedulers.mainThread()); // 4. 订阅处理后的API流,更新UI apiProcessedObservable.subscribe(apiResults -> { apiSearchResultsAdapter.updateResults(apiResults); }, throwable -> { // 这里不会触发,因为前面已经用onErrorReturn处理了API异常 }); } static class SearchResultWrapper { List<LocalSearchResult> localResults; List<ApiSearchResult> apiResults; }
各需求的满足说明
- 需求1:本地结果立即展示:本地Observable单独订阅,一发射数据就切换到主线程更新UI,无需等待API结果
- 需求2:API结果去重后展示:通过
combineLatest拿到本地和API的所有结果,用HashSet高效去重后再更新UI - 需求3:API先返回则等本地结果:
combineLatest必须等两个Observable都发射数据才会执行合并逻辑,所以API结果会一直等待本地数据就绪后再处理 - 需求4:API异常不影响本地结果:用
onErrorReturn捕获API的异常,返回空的成功响应,保证API流不会中断,同时本地流的订阅不受任何影响
额外优化点
- 用
HashSet替代ArrayList存储本地结果名称,将去重的时间复杂度从O(n)降到O(1) - 用
filter过滤空的API结果,避免不必要的UI更新 - 用
publish().autoConnect(2)确保本地查询只执行一次,不会因为两个订阅流重复触发数据库操作
内容的提问来源于stack exchange,提问作者hari_shankar
相关产品推荐
相关产品推荐

