RxJava Flowable如何根据Firebase ChildEventListener事件类型拆分发射流
Firebase ChildEventListener 按事件拆分RxJava管道的Java实现
核心实现思路:先将不同类型的子节点事件包装为带类型标识的统一对象,再通过RxJava操作符过滤拆分出独立的事件处理流,实现差异化逻辑分发。
1. 定义事件类型包装结构
首先定义事件类型枚举和统一事件包装类,用于区分不同回调事件、携带对应数据:
// 子节点事件类型枚举 public enum ChildEventType { ADDED, CHANGED, REMOVED, MOVED, CANCELLED } // 统一事件包装类 public class ChildEvent { private final ChildEventType eventType; private final User user; // 仅ADDED/CHANGED/MOVED事件携带该参数 private final String previousChildName; // 仅CANCELLED事件携带该参数 private final DatabaseError error; public ChildEvent(ChildEventType eventType, User user, String previousChildName, DatabaseError error) { this.eventType = eventType; this.user = user; this.previousChildName = previousChildName; this.error = error; } // 各事件快捷构造方法 public static ChildEvent added(User user, String previousChildName) { return new ChildEvent(ChildEventType.ADDED, user, previousChildName, null); } public static ChildEvent changed(User user, String previousChildName) { return new ChildEvent(ChildEventType.CHANGED, user, previousChildName, null); } public static ChildEvent removed(User user) { return new ChildEvent(ChildEventType.REMOVED, user, null, null); } public static ChildEvent moved(User user, String previousChildName) { return new ChildEvent(ChildEventType.MOVED, user, previousChildName, null); } public static ChildEvent cancelled(DatabaseError error) { return new ChildEvent(ChildEventType.CANCELLED, null, null, error); } // 各字段Getter方法 public ChildEventType getEventType() { return eventType; } public User getUser() { return user; } public String getPreviousChildName() { return previousChildName; } public DatabaseError getError() { return error; } }
2. 封装ChildEventListener为共享Flowable
修正原参考代码的问题(未处理监听移除、参数错误、无背压策略),封装为可共享的Flowable:
private Flowable<ChildEvent> userChildEventFlowable; // 全局仅初始化一次即可 private void initUserEventFlowable() { DatabaseReference ref = FirebaseDatabase.getInstance().getReference().child("Users"); userChildEventFlowable = Flowable.create(emitter -> { ChildEventListener listener = new ChildEventListener() { @Override public void onChildAdded(@NonNull DataSnapshot dataSnapshot, @Nullable String previousChildName) { User user = dataSnapshot.getValue(User.class); if (user != null && !emitter.isCancelled()) { emitter.onNext(ChildEvent.added(user, previousChildName)); } } @Override public void onChildChanged(@NonNull DataSnapshot dataSnapshot, @Nullable String previousChildName) { User user = dataSnapshot.getValue(User.class); if (user != null && !emitter.isCancelled()) { emitter.onNext(ChildEvent.changed(user, previousChildName)); } } @Override public void onChildRemoved(@NonNull DataSnapshot dataSnapshot) { User user = dataSnapshot.getValue(User.class); if (user != null && !emitter.isCancelled()) { emitter.onNext(ChildEvent.removed(user)); } } @Override public void onChildMoved(@NonNull DataSnapshot dataSnapshot, @Nullable String previousChildName) { User user = dataSnapshot.getValue(User.class); if (user != null && !emitter.isCancelled()) { emitter.onNext(ChildEvent.moved(user, previousChildName)); } } @Override public void onCancelled(@NonNull DatabaseError databaseError) { if (!emitter.isCancelled()) { // 若需要将取消作为错误处理,可替换为 emitter.onError(databaseError.toException()) emitter.onNext(ChildEvent.cancelled(databaseError)); } } }; // 添加Firebase监听 ref.addChildEventListener(listener); // 订阅取消时自动移除监听,避免内存泄漏 emitter.setCancellable(() -> ref.removeEventListener(listener)); }, BackpressureStrategy.BUFFER) // 多订阅者共享同一个监听,避免重复注册 .share(); }
3. 拆分独立事件处理管道
通过RxJava过滤操作符,拆分出不同事件的独立处理流:
// 子节点新增事件流 public Flowable<User> getChildAddedFlow() { return userChildEventFlowable .filter(event -> event.getEventType() == ChildEventType.ADDED) .map(ChildEvent::getUser); } // 子节点修改事件流 public Flowable<User> getChildChangedFlow() { return userChildEventFlowable .filter(event -> event.getEventType() == ChildEventType.CHANGED) .map(ChildEvent::getUser); } // 子节点删除事件流 public Flowable<User> getChildRemovedFlow() { return userChildEventFlowable .filter(event -> event.getEventType() == ChildEventType.REMOVED) .map(ChildEvent::getUser); } // 监听取消事件流 public Flowable<DatabaseError> getCancelledFlow() { return userChildEventFlowable .filter(event -> event.getEventType() == ChildEventType.CANCELLED) .map(ChildEvent::getError); }
如果需要用到previousChildName等额外参数,可不做map转换,直接返回ChildEvent对象即可
4. 差异化逻辑使用示例
分别订阅不同事件流,实现独立的业务逻辑处理:
// 初始化事件流 initUserEventFlowable(); // 处理新增用户逻辑 Disposable addedDisposable = getChildAddedFlow() .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe(user -> { // 示例:新增用户插入列表头部 userList.add(0, user); userAdapter.notifyItemInserted(0); }); // 处理修改用户逻辑 Disposable changedDisposable = getChildChangedFlow() .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe(user -> { // 示例:更新列表对应位置数据 int index = findUserIndexById(user.getUid()); if (index != -1) { userList.set(index, user); userAdapter.notifyItemChanged(index); } }); // 处理删除用户逻辑 Disposable removedDisposable = getChildRemovedFlow() .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe(user -> { // 示例:移除列表对应项 int index = findUserIndexById(user.getUid()); if (index != -1) { userList.remove(index); userAdapter.notifyItemRemoved(index); } }); // 页面销毁时取消所有订阅,避免内存泄漏 @Override protected void onDestroy() { super.onDestroy(); addedDisposable.dispose(); changedDisposable.dispose(); removedDisposable.dispose(); }
内容的提问来源于stack exchange,提问作者nirkov
相关产品推荐
相关产品推荐

