You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

RxJava并行查询本地DB与API:需先展示本地结果并去重

解决方案:拆分流并结合RxJava算子满足需求

你之前用CombineLatest的问题在于,它必须等两个Observable都发射数据才会输出结果,没法单独先展示本地查询结果。咱们可以通过共享本地Observable+拆分两个订阅流的方式来实现所有需求,具体思路如下:

  1. 共享本地查询的Observable,避免重复执行本地数据库查询
  2. 单独订阅本地流,一有结果就立即展示给用户
  3. 处理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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 04:25:41