【问题标题】:RXJava 2 Polling with multiple subscribersRXJava 2 轮询多个订阅者
【发布时间】:2017-04-06 12:02:53
【问题描述】:

我想开始替换/使用 RXJava2 来代替 Observer 和 Listeners 进行轮询。现在只有一个问题。我有一个 Polling Observable,只有在连接了至少一个 Subscriber 时才应该启动它。如果连接了多个订阅者,则间隔应该相同。表示:一个 observable 重复轮询过程 n 秒。如果 observable 有 1..* 个订阅者,它应该继续轮询 n 秒并通知所有订阅者结果。

这就是我使用 Listeners 和/或我的 RXJava 解决方案的方式。

我的第一次尝试是创建一个单例类,它只创建一个 PublishSubject。如果有人订阅,它将在 onNext() 中获取数据。现在我的 Polling Observer 在某处启动并将数据推送到主题。这不起作用,因为它是

  • 糟糕的图案设计
  • 仅在连接了订阅者时才启动,在没有订阅者可用时停止
  • 不成功共享数据,需要两个类(用于主题和重复可观察)

    public class SingleTonClass { 
    private PublishSubject<List<Data>> subject = PublishSubject.create();
    
    
    public PublishSubject getSubject() {
        return this.subject;
    }
    
    public void setData(List<Data> data) { 
       subject.onNext(data);
    }
    }
    

我很乐意避免使用侦听器/接口来共享信息并让 rxjava2 完成它的工作。

经过研究,我发现有 refcount() 和 share() 但我不确定这是否是解决此问题的正确方法。在我的情况下,它是一个 REST 服务,如果至少有一个订阅者连接,它会轮询服务器,否则它应该停止轮询,因为在这种情况下获取数据没有意义。

我试图解决它,但它不能正常工作:

Polling using RXJava2 / RXAndroid 2 and Retrofit

【问题讨论】:

  • 如果你想拥有一个有多个订阅者的 Observable,那么 share() 可能是在 replay(1) 之后的方式,以确保新订阅者获得最新的发射。

标签: java android rx-java rx-android rx-java2


【解决方案1】:

我会这样做:

Observable<Data> dataSource = Observable.interval(INTERVAL, TIME_UNIT)
    .observeOn(Schedulers.io()) // make REST requests on IO threads
    .map(n -> {
            return requestData();
        })
    .replay(1);

replay() 运算符包括share(),而后者又包括publish()refCount() 功能。这使您的 observable hot,即所有订阅者共享一个订阅。它会自动与第一个订阅者订阅(开始新的interval 序列),并在最后一个订阅者离开时取消订阅(停止interval)。

replay(1) 还缓存最后发出的值,即新订阅者不必等待新数据到达。

【讨论】:

    猜你喜欢
    • 2018-06-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多