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

RxJava竞态条件规避:如何让查询等待数据库插入完成?

问题场景

我在ViewModel的构造函数中,使用Completable在后台线程将数据插入本地数据库:

public class MainViewModel extends ViewModel {
    public MainViewModel(){
        localRepository.insertValueIntoDatabase().subscribeOn(Schedulers.io())
                    .subscribe(() -> {
                        sharedPrefManager.setAnotherValue(true);
                    }, throwable -> {
                        Timber.e(throwable, "Failed to insert into DB");
                    });
    }
}

在MainActivity中,创建MainViewModel后会立即调用performQueryWithValue()方法,该方法需要基于构造函数中插入的数据执行查询:

viewModel = new ViewModelProvider(this).get(MainViewModel.class);
viewModel.performQueryWithValue();

performQueryWithValue()的实现如下:

public class MainViewModel extends ViewModel {
    public MainViewModel(){
        localRepository.insertValueIntoDatabase().subscribeOn(Schedulers.io())
                    .subscribe(() -> {
                        sharedPrefManager.setAnotherValue(true);
                    }, throwable -> {
                        Timber.e(throwable, "Failed to insert into DB");
                    });
    }

    public void performQueryWithValue(){
        localRepository.getValueFromDatabase().flatMapSingle(value -> {
             if(value == 0){
                return remoteRepository.performQueryOne();
             }else{
                return remoteRepository.performQueryTwo();
             }
        })
        .subscribeOn(Schedulers.io())
        .observeOn(AndroidSchedulers.mainThread())
        .subscribe(result-> {
                    // 处理结果
        }, err -> {
                    // 处理错误
        });
    }
}

当localRepository.insertValueIntoDatabase()插入操作耗时10秒时,查询会先于插入完成执行,导致获取不到正确数据。请问如何让performQueryWithValue()等待插入完成后再执行?

解决方案

方法1:让查询操作依赖插入流

将构造函数中的插入Completable保存为成员变量,在查询方法中通过andThen操作符,确保插入完成后再执行查询逻辑。

修改后的ViewModel代码:

public class MainViewModel extends ViewModel {
    // 保存插入操作的Completable,不直接在构造函数中subscribe
    private final Completable insertCompletable;

    public MainViewModel() {
        insertCompletable = localRepository.insertValueIntoDatabase()
                .subscribeOn(Schedulers.io())
                // 插入完成后的回调逻辑移到doOnComplete
                .doOnComplete(() -> sharedPrefManager.setAnotherValue(true))
                .doOnError(throwable -> Timber.e(throwable, "Failed to insert into DB"));
    }

    public void performQueryWithValue() {
        insertCompletable
                // 等插入完成后,再执行查询操作
                .andThen(localRepository.getValueFromDatabase())
                .flatMapSingle(value -> {
                    if (value == 0) {
                        return remoteRepository.performQueryOne();
                    } else {
                        return remoteRepository.performQueryTwo();
                    }
                })
                .subscribeOn(Schedulers.io())
                .observeOn(AndroidSchedulers.mainThread())
                .subscribe(result -> {
                    // 处理查询结果
                }, err -> {
                    // 处理错误(包括插入失败和查询失败)
                });
    }
}

这种方式的核心是利用RxJava的操作符串联异步流,天然保证执行顺序,无需额外的状态标记。

方法2:用LiveData标记插入状态,延迟触发查询

如果需要更灵活的控制(比如插入失败时决定是否执行查询),可以用LiveData传递插入完成的状态,在Activity中监听状态变化后再调用查询方法。

ViewModel代码

public class MainViewModel extends ViewModel {
    private final MutableLiveData<Boolean> insertCompleted = new MutableLiveData<>();

    public MainViewModel() {
        localRepository.insertValueIntoDatabase()
                .subscribeOn(Schedulers.io())
                .subscribe(() -> {
                    sharedPrefManager.setAnotherValue(true);
                    // 插入成功,通知状态
                    insertCompleted.postValue(true);
                }, throwable -> {
                    Timber.e(throwable, "Failed to insert into DB");
                    // 插入失败,也通知状态
                    insertCompleted.postValue(false);
                });
    }

    // 对外暴露观察插入状态的LiveData
    public LiveData<Boolean> getInsertCompleted() {
        return insertCompleted;
    }

    public void performQueryWithValue() {
        localRepository.getValueFromDatabase()
                .flatMapSingle(value -> {
                    if (value == 0) {
                        return remoteRepository.performQueryOne();
                    } else {
                        return remoteRepository.performQueryTwo();
                    }
                })
                .subscribeOn(Schedulers.io())
                .observeOn(AndroidSchedulers.mainThread())
                .subscribe(result -> {
                    // 处理查询结果
                }, err -> {
                    // 处理查询错误
                });
    }
}

Activity代码修改

viewModel = new ViewModelProvider(this).get(MainViewModel.class);
// 监听插入完成状态,只有成功后才执行查询
viewModel.getInsertCompleted().observe(this, isCompleted -> {
    if (isCompleted) {
        viewModel.performQueryWithValue();
    } else {
        // 这里可以处理插入失败的逻辑,比如提示用户
    }
});

注意事项

  • 避免在ViewModel构造函数中直接订阅异步流:这样会导致流的执行不受控制,后续操作无法依赖其完成状态。
  • 错误处理:两种方案都需要考虑插入失败的情况,决定是否继续执行查询或进行错误提示。

内容的提问来源于stack exchange,提问作者WHOATEMYNOODLES

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 18:32:08