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
相关产品推荐
相关产品推荐

