【问题标题】:Combining Firebase realtime data listener with RxJava将 Firebase 实时数据侦听器与 RxJava 相结合
【发布时间】:2016-03-31 10:27:07
【问题描述】:

我在我的应用程序中使用 Firebase 以及 RxJava。 Firebase 能够在后端数据发生更改(添加、删除、更改……)时通知您的应用。 我正在尝试将 Firebase 的功能与 RxJava 结合起来。

我正在侦听的数据称为 LeisureObservable 发出 LeisureUpdate,其中包含 Leisure 和更新类型(添加、删除、移动、更改)。

这是我允许订阅此事件的方法。

private Observable<LeisureUpdate> leisureUpdatesObservable;
private ChildEventListener leisureUpdatesListener;
private int leisureUpdatesSubscriptionsCount;

@NonNull
public Observable<LeisureUpdate> subscribeToLeisuresUpdates() {
    if (leisureUpdatesObservable == null) {
        leisureUpdatesObservable = Observable.create(new Observable.OnSubscribe<LeisureUpdate>() {

            @Override
            public void call(final Subscriber<? super LeisureUpdate> subscriber) {
                leisureUpdatesListener = firebase.child(FirebaseStructure.LEISURES).addChildEventListener(new ChildEventListener() {
                    @Override
                    public void onChildAdded(DataSnapshot dataSnapshot, String s) {
                        final Leisure leisure = convertMapToLeisure((Map<String, Object>) dataSnapshot.getValue());
                        subscriber.onNext(new LeisureUpdate(leisure, LeisureUpdate.ADDED));
                    }

                    @Override
                    public void onChildChanged(DataSnapshot dataSnapshot, String s) {
                        final Leisure leisure = convertMapToLeisure((Map<String, Object>) dataSnapshot.getValue());
                        subscriber.onNext(new LeisureUpdate(leisure, LeisureUpdate.CHANGED));
                    }

                    @Override
                    public void onChildRemoved(DataSnapshot dataSnapshot) {
                        final Leisure leisure = convertMapToLeisure((Map<String, Object>) dataSnapshot.getValue());
                        subscriber.onNext(new LeisureUpdate(leisure, LeisureUpdate.REMOVED));
                    }

                    @Override
                    public void onChildMoved(DataSnapshot dataSnapshot, String s) {
                        final Leisure leisure = convertMapToLeisure((Map<String, Object>) dataSnapshot.getValue());
                        subscriber.onNext(new LeisureUpdate(leisure, LeisureUpdate.MOVED));
                    }

                    @Override
                    public void onCancelled(FirebaseError firebaseError) {
                        subscriber.onError(new Error(firebaseError.getMessage()));
                    }
                });
            }
        });
    }
    leisureUpdatesSubscriptionsCount++;
    return leisureUpdatesObservable;
}

首先,我想使用Observable.fromCallable() 方法来创建Observable,但我想这是不可能的,因为Firebase 使用回调,对吧?

我保留Observable 的单个实例,以便始终拥有一个Observable 可以订阅多个Subscriber

当每个人都取消订阅并且我需要停止监听 Firebase 中的事件时,问题就出现了。 我无论如何都没有找到让Observable 了解是否还有订阅。所以我一直在计算我打了多少电话给subscribeToLeisuresUpdates()leisureUpdatesSubscriptionsCount

然后每次有人想退订它都必须打电话

@Override
public void unsubscribeFromLeisuresUpdates() {
    if (leisureUpdatesObservable == null) {
        return;
    }
    leisureUpdatesSubscriptionsCount--;
    if (leisureUpdatesSubscriptionsCount == 0) {
        firebase.child(FirebaseStructure.LEISURES).removeEventListener(leisureUpdatesListener);
        leisureUpdatesObservable = null;
    }
}

这是我发现让Observable 在有订阅者时发出项目的唯一方法,但我觉得必须有一种更简单的方法,特别是在没有更多订阅者收听可观察对象时理解。

有遇到过类似问题或有不同方法的人吗?

【问题讨论】:

  • 另外,你可以使用我的图书馆rxFirebase

标签: android firebase rx-java


【解决方案1】:

你可以使用 Observable.fromEmitter,类似这样的东西

    return Observable.fromEmitter(new Action1<Emitter<LeisureUpdate>>() {
        @Override
        public void call(final Emitter<LeisureUpdate> leisureUpdateEmitter) {
            final ValueEventListener listener = new ValueEventListener() {
                @Override
                public void onDataChange(DataSnapshot dataSnapshot) {
                    // process update
                    LeisureUpdate leisureUpdate = ...
                    leisureUpdateEmitter.onNext(leisureUpdate);
                }

                @Override
                public void onCancelled(DatabaseError databaseError) {
                    leisureUpdateEmitter.onError(new Throwable(databaseError.getMessage()));
                    mDatabaseReference.removeEventListener(this);
                }
            };
            mDatabaseReference.addValueEventListener(listener);
            leisureUpdateEmitter.setCancellation(new Cancellable() {
                @Override
                public void cancel() throws Exception {
                    mDatabaseReference.removeEventListener(listener);
                }
            });
        }
    }, Emitter.BackpressureMode.BUFFER);

【讨论】:

【解决方案2】:

把它放在你的 Observable.create() 最后。

subscriber.add(Subscriptions.create(new Action0() {
                    @Override public void call() {
                        ref.removeEventListener(leisureUpdatesListener);
                    }
                }));

【讨论】:

    【解决方案3】:

    我建议您检查下一个库作为参考(或仅使用它):

    RxJava : https://github.com/nmoskalenko/RxFirebase

    RxJava 2.0:https://github.com/FrangSierra/Rx2Firebase

    其中一个适用于 RxJava,另一个适用于 RxJava 2.0 的新 RC。有兴趣的可以看看here两者的区别。

    【讨论】:

      【解决方案4】:

      这里是使用 RxJava2 和 Firebase 的 CompletionListener 的示例代码:

      Completable.create(new CompletableOnSubscribe() {
              @Override
              public void subscribe(final CompletableEmitter e) throws Exception {
                  String orderKey = FirebaseDatabase.getInstance().getReference().child("orders").push().getKey();
                  FirebaseDatabase.getInstance().getReference().child("orders").child(orderKey).setValue(order,
                          new DatabaseReference.CompletionListener() {
                              @Override
                              public void onComplete(DatabaseError databaseError, DatabaseReference databaseReference) {
                                  if (e.isDisposed()) {
                                      return;
                                  }
                                  if (databaseError == null) {
                                      e.onComplete();
                                  } else {
                                      e.onError(new Throwable(databaseError.getMessage()));
                                  }
                              }
                          });
      
              }
          }).subscribeOn(Schedulers.io()).observeOn(AndroidSchedulers.mainThread());
      

      【讨论】:

      【解决方案5】:

      另一个惊人的库,它将帮助您将所有 firebase 实时数据库逻辑包装在 rx 模式下。

      https://github.com/Link184/Respiration

      您可以在这里创建您的 firebase 存储库并从 GeneralRepository 扩展它,例如:

      @RespirationRepository(dataSnapshotType = Leisure.class)
      public class LeisureRepository extends GeneralRepository<Leisure>{
          protected LeisureRepository(Configuration<Leisure> repositoryConfig) {
              super(repositoryConfig);
          }
      
          // observable will never emit any events if there are no more subscribers
          public void killRepository() {
              if (!behaviorSubject.hasObservers()) {
                  //protected fields
                  behaviorSubject.onComplete();
                  databaseReference.removeEventListener(valueListener);
              }
          }
      }
      

      您可以通过这种方式“杀死”您的存储库:

      // LeisureRepositoryBuilder is a generated class by annotation processor, will appear after a successful gradle build
      LeisureRepositoryBuilder.getInstance().killRepository();
      

      但我认为对于您的情况,最好扩展 com.link184.respiration.repository.ListRepository 以避免通过 LeisureUpdate 从 java.util.Map 到 Leisure 模型的数据映射

      【讨论】:

      • 你也应该解释一下如何使用这个库。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-08-08
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-02-22
      相关资源
      最近更新 更多