【问题标题】:RXjava continuous streams of dataRXjava 连续数据流
【发布时间】:2018-10-19 04:22:25
【问题描述】:

在 RxJava 中订阅 observable 的正确方法是什么,该 observable 在未知时间连续接收来自不同来源的事件。

例如:假设我们有从服务器接收的任务,并且可以从多个不同的区域(例如,由推送通知、轮询事件、用户交互等触发)启动对服务器的调用。

我们在 UI 上唯一关心的是我们会收到已收到任务的通知。我们不在乎他们来自哪里。我实际上想开始观察 Activity 的生命周期,并且模型会根据需要更新观察者

我已经实现了下面的类,它可以满足我的需求,但我不确定它是否正确,或者 RxJava 是否已经解决了这样的问题。

这个类在 ConnectableObservable 上有效创建,可以让许多观察者订阅一个 Observable(确保所有观察者获得相同的流)。我注意到的一件事是,在以这种方式订阅 ConnectableObservable 时,调用 observeOn 和 subscribeOn 可能会导致意外结果,这可能是一个问题,因为该类无法控制谁在使用 ConnectableObservable。

public class ApiService {

private Emitter<String> myEmitter;
private ConnectableObservable myObservable;

public ApiService() {
    //Create an observable that is simply used to get the emitter.
    myObservable = Observable.create(new ObservableOnSubscribe<String>() {
        @Override
        public void subscribe(ObservableEmitter<String> e) throws Exception {
            myEmitter = e;
        }
    }).publish();
    //connect must be called here to ensure we have an instance of the emitter
    // before we have any subscribers
    myObservable.connect();

}

/**
 * This method returns the observable that all observers will subscribe to
 *
 * @return
 */
public Observable<String> getObservable() {
    return myObservable;
}

/**
 * This method is used to simulate a value that has been received from
 * an unknown source
 *
 * @param value
 */
public void run(final String value) {
    Observable.create(new ObservableOnSubscribe<String>() {
        @Override
        public void subscribe(ObservableEmitter<String> e) throws Exception {
            myEmitter.onNext(value);
            //api call
        }
    }).subscribe();
}

}

我也很好奇,假设每个订阅一个观察者的观察者在适当的时间被处理掉,这样做是否存在内存泄漏问题。

这是看到 Bob 的回答后重构的类

public class ApiService {

private PublishProcessor<String> myProcessor = PublishProcessor.create();

public void subscribe(Subscriber<String> subscriber) {
    myProcessor.subscribe(subscriber);
}

/**
 * This method is used to simulate a value that has been received from
 * an unknown source
 *
 * @param value
 */
public void run(final String value) {
    myProcessor.onNext(value);
}

}

【问题讨论】:

  • 如果有从observable 发出的high no of elements 使用flowable ...coz 这可能会导致 back-pressure 并且您的应用程序可能会耗尽内存(pub在这些情况下,-sub 模式比观察者模式更好)

标签: android rx-java rx-java2


【解决方案1】:

使用Subject (RxJava) 或Processor (RxJava 2) 进行订阅。然后,您将为主题订阅每个可观察的源。最终,您将订阅该主题并获得合并的排放流。

或者,您可以使用Relay 将下游观察者与可能来自上游的任何onComplete() 或onError() 隔离开来。当任何 observables 可能在其他 observables 之前完成时,这是一个更好的选择。

【讨论】:

  • 更好,谢谢。现在我已经知道要搜索什么,我找到了这篇简洁的文章。我还用重构类 medium.com/@ar.xa.vasquez/… 更新了我的问题
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-06-13
  • 2017-10-16
  • 1970-01-01
  • 1970-01-01
  • 2021-04-04
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多