带RxJava2的Android Room中的交易

Ste*_*eve 5 android rx-java rx-java2 android-room

我的应用程序的要求是允许用户执行多个步骤,然后在完成后根据每个步骤中的条目将值写入数据库。UI中的每个步骤都可能有助于将需要写入数据库的操作。数据可能在多个表中,并且与这些表中的不同行有关。如果任何数据库操作失败,则整个操作应失败。

我最初考虑将所有数据加载到内存中,进行操作,然后在每个可能的实体中调用update方法(使用REPLACE的冲突策略),但是内存中可能有非常大量的数据。

我认为我可以组装一个List,其中显示中的每个Fragment都贡献一个或多个Completable,然后在UI流结束时使用Completable.concat()顺序执行这些。如下所示:

    Completable one = Completable.fromAction(() -> Log.w(LOG_TAG, "(1)")).delay(1, TimeUnit.SECONDS);
    Completable two = Completable.fromAction(() -> Log.w(LOG_TAG, "(2)")).delay(2, TimeUnit.SECONDS);
    Completable three = Completable.fromAction(() -> Log.w(LOG_TAG, "(3)")).delay(3, TimeUnit.SECONDS);
    Completable four = Completable.fromAction(() -> Log.w(LOG_TAG, "(4)")).delay(3, TimeUnit.SECONDS);

    Completable.concatArray(one, two, three, four)
            .doOnSubscribe(__ -> {
                mRoomDatabase.beginTransaction();
            })
            .doOnComplete(() -> {
                mRoomDatabase.setTransactionSuccessful();
            })
            .doFinally(() -> {
                mRoomDatabase.endTransaction();
            })
            .observeOn(AndroidSchedulers.mainThread())
            .subscribeOn(Schedulers.io())
            .subscribe();
Run Code Online (Sandbox Code Playgroud)

Completables实际上是Room DAO插入/更新/删除方法的包装。我可能还会在完成时执行UI操作,这就是为什么我在主线程上进行观察的原因。

当我执行此代码时,我得到以下日志:

W/MyPresenter: Begin transaction.
W/MyPresenter: (1)
W/MyPresenter: (2)
W/MyPresenter: (3)
W/MyPresenter: (4)
W/MyPresenter: Set transaction successful.
W/MyPresenter: End transaction.
W/System.err: java.lang.IllegalStateException: Cannot perform this operation because there is no current transaction.
W/System.err:     at android.database.sqlite.SQLiteSession.throwIfNoTransaction(SQLiteSession.java:915)
W/System.err:     at android.database.sqlite.SQLiteSession.endTransaction(SQLiteSession.java:398)
W/System.err:     at android.database.sqlite.SQLiteDatabase.endTransaction(SQLiteDatabase.java:524)
W/System.err:     at android.arch.persistence.db.framework.FrameworkSQLiteDatabase.endTransaction(FrameworkSQLiteDatabase.java:88)
W/System.err:     at android.arch.persistence.room.RoomDatabase.endTransaction(RoomDatabase.java:220)
W/System.err:     at ...lambda$doTest$22$MyPresenter(MyPresenter.java:490)
Run Code Online (Sandbox Code Playgroud)

为什么我到达doFinally时交易消失了?我也欢迎对这种方法的质量或可行性发表任何评论,因为我对RxJava和Room还是很陌生。

Ste*_*eve 9

通过记录当前线程并细读Android开发人员文档,我想我可能最终可以理解我做错了什么。

1)事务必须在同一线程上发生。这就是为什么它告诉我没有交易的原因。我显然在线程之间跳动。

2)doOnSubscribe,doOnComplete和doFinally方法是副作用,因此不属于实际流本身。这意味着它们不会出现在我订阅的调度程序上。它们将出现在我观察到的调度程序上。

3)因为我想在完成时在UI线程上接收结果,但是希望副作用在后台线程上发生,所以我需要更改观察的位置。

Completable.concatArray(one, two, three, four)
                .observeOn(Schedulers.single()) // OFF UI THREAD
                .doOnSubscribe(__ -> {
                    Log.w(LOG_TAG, "Begin transaction. " + Thread.currentThread().toString());
                    mRoomDatabase.beginTransaction();
                })
                .doOnComplete(() -> {
                    Log.w(LOG_TAG, "Set transaction successful."  + Thread.currentThread().toString());
                    mRoomDatabase.setTransactionSuccessful();
                })
                .doFinally(() -> {
                    Log.w(LOG_TAG, "End transaction."  + Thread.currentThread().toString());
                    mRoomDatabase.endTransaction();
                })
                .subscribeOn(Schedulers.single())
                .observeOn(AndroidSchedulers.mainThread()) // ON UI THREAD
                .subscribeWith(new CompletableObserver() {
                    @Override
                    public void onSubscribe(Disposable d) {
                        Log.w(LOG_TAG, "onSubscribe."  + Thread.currentThread().toString());
                    }

                    @Override
                    public void onComplete() {
                        Log.w(LOG_TAG, "onComplete."  + Thread.currentThread().toString());
                    }

                    @Override
                    public void onError(Throwable e) {
                        Log.e(LOG_TAG, "onError." + Thread.currentThread().toString());
                    }
                });
Run Code Online (Sandbox Code Playgroud)

现在,日志记录语句如下所示:

W/MyPresenter: onSubscribe.Thread[main,5,main]
W/MyPresenter: Begin transaction. Thread[RxSingleScheduler-1,5,main]
W/MyPresenter: (1)
W/MyPresenter: (2)
W/MyPresenter: (3)
W/MyPresenter: (4)
W/MyPresenter: Set transaction successful.Thread[RxSingleScheduler-1,5,main]
W/MyPresenter: End transaction.Thread[RxSingleScheduler-1,5,main]
W/MyPresenter: onComplete.Thread[main,5,main]
Run Code Online (Sandbox Code Playgroud)

我相信这可以实现我所追求的目标,但是基于Room的RxJava Completables的分步组装是否会成功还有待观察。我会随时留意任何评论/答案,并可能报告后代。