【问题标题】:How to subclass Observable in RxJava?如何在 RxJava 中继承 Observable?
【发布时间】:2015-11-26 09:56:49
【问题描述】:

我知道你应该不惜一切代价避免这种情况,但是如果我有一个 RxJava 中的子类 Observable 的有效用例怎么办?可能吗?我该怎么做?

在这种特定情况下,我有一个当前返回请求的“存储库”类:

class Request<T> {
    public abstract Object key();
    public abstract Observable<T> asObservable();

    [...]

    public Request<T> transform(Func1<Request<T>, Observable<T>> transformation) {
        Request<T> self = this;
        return new Request<T>() {
             @Override public Object key() { return self.key; }
             @Override public Observable<T> asObservable() { return transformation.call(self); }
        }
    }
}

然后我在需要请求密钥(如缓存)的上下文中使用转换方法修改响应可观察(asObservable):

 service.getItemList() // <- returns a Request<List<Item>>
     .transform(r -> r.asObservable()
             // The activity is the current Activity in Android
             .compose(Operators.ensureThereIsAnAccount(activity))
             // The cache comes last because we don't need auth for cached responses
             .compose(cache.cacheTransformation(r.key())))
     .asObservable()
     [...  your common RxJava code ...]

现在,如果我的 Request 类是 Observable 子类会非常方便,因为这样我可以消除所有 .asObservable() 调用,客户甚至不需要知道我的 Request 类。

【问题讨论】:

  • 如果你确定你真的想把事情弄得这么乱:github.com/ReactiveX/RxJava/wiki/Creating-Observables,但上面的代码似乎混淆了问题。
  • 我在 subclassing Observables 上找不到任何参考。我错过了什么吗?
  • Observable 可能不打算被子类化。记住 Effective Java 第 16 条:优先考虑组合而不是继承。为什么你认为子类化在这里是正确的?
  • 但是 RxJava 中有 ConnectableObservable。
  • 引用第 16 条的介绍:“继承是实现代码重用的强大方法,但它并不总是最好的工具。使用不当会导致软件脆弱。它是安全的在包中使用继承,其中子类和超类的实现在同一个程序员的控制下。在扩展专门设计和记录的类时使用继承也是安全的(第 17 条)。跨包继承普通的具体类然而,边界是危险的。”

标签: rx-java


【解决方案1】:

可以继承 Observable(我们为 Subjects 和 ConnectableObservables 这样做),但需要额外考虑,因为您需要传入 OnSubscribe 回调来处理传入的 Subscribers .我不清楚如果有人订阅了你的请求应该做什么,所以我会给你两个扩展 Observable 的例子:

没有共享可变状态的可观察

如果您没有要在订阅者之间共享的可变状态,您可以扩展 Observable 并将您的操作传递给 super

public final class MyObservable extends Observable<Long> {
    public MyObservable() {
        super(new OnSubscribe<Long>() {
            @Override public void call(Subscriber<? super Long> child) {
                child.onNext(System.currentTimeMillis());
                child.onCompleted();
            }
        });
    }
}

具有共享可变状态的可观察

这通常比较棘手,因为您需要从 OnSubscribe 方法和 Observable 的方法访问共享状态,但 Java 不会让您在 super 之前接触 OnSubscribe 内部类的实例字段完全的。解决方案是从构造函数中分解出这样的共享状态和OnSubscribe,并使用静态工厂方法来设置两者:

public final class MySharedObservable extends Observable<Long> {
    public static MySharedObservable create() {
        final AtomicLong counter = new AtomicLong();
        OnSubscribe<Long> onSubscribe = new OnSubscribe<Long>() {
            @Override
            public void call(Subscriber<? super Long> t1) {
                t1.onNext(counter.incrementAndGet());
                t1.onCompleted();
            }
        };
        return new MySharedObservable(onSubscribe, counter);
    }
    private AtomicLong counter;

    private MySharedObservable(OnSubscribe<Long> onSubscribe, AtomicLong counter) {
        super(onSubscribe);
        this.counter = counter;
    }
    public long getCounter() {
        return counter.get();
    }
}

【讨论】:

  • 这看起来不错,但我如何取消上述订阅? onSubscribe 返回 void,并且没有明显的取消方法(假设您的计数器执行异步操作,例如启动线程或发出网络请求......如何停止外部事物?)
  • 糟糕,没关系。订阅者上有一个 .isUnsubscribed() ,您可以调用它。 RxJava 完全是来自 C# Rx 背景的疯狂
猜你喜欢
  • 1970-01-01
  • 2023-03-26
  • 1970-01-01
  • 2019-02-10
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-07-21
相关资源
最近更新 更多