【问题标题】:Making a Subscriber subscribe to a Subject in RxJava让订阅者订阅 RxJava 中的主题
【发布时间】:2020-07-04 23:13:15
【问题描述】:

在 RxJava 3 中,有没有办法将 Subscriber 订阅到 Subject? Subject 也是一个 Observable,这意味着它应该能够提供一个上游,当 Subject 的 onNext 发出某些东西时,其他人(Subscriber 实现,如 DisposableSubscriberDefaultSubscriber)应该能够获取数据。但是下面的代码给出了编译时错误,我找不到任何其他机制来实现这一点。

public class NewsPublisher {
    static Subject<Integer>  newsSubject = PublishSubject.create();

public class NewsSubscriber
    extends DisposableSubscriber<Integer> {

    void bind() {
        NewsPublisher.newsSubject.subscribe(this);
    }

错误消息:无法解析方法“订阅(NewsSubscriber)”

Subject 的subscribe() 方法只允许ConsumerObserver 作为输入。但从逻辑上讲,我想要的只是当 Subject(也是一个 Observable)发出一些东西时,我可以让一个 Subscriber 监听,因此它的 onNext() 被调用。

【问题讨论】:

    标签: rx-java system.reactive


    【解决方案1】:

    更新:

    决定自己尝试代码并找到一种解决方法。使用 lambda 可以链接 onNext() 方法。 lambda istelf 是一个匿名 Consumer 对象,因此它可能会破坏整个类的目的(或不是?取决于你还在做什么),但这很有效。

    import io.reactivex.rxjava3.disposables.Disposable;
    import io.reactivex.rxjava3.subscribers.DisposableSubscriber;
    
    public class NewsSubscriber extends DisposableSubscriber<Integer> {
    
        Disposable disposable;
    
        void bind() {
            disposable = NewsPublisher.newsSubject.subscribe(this::onNext);
        }
    
        public void onNext(Integer integer) {   
            // implement
        }
    
        public void onError(Throwable throwable) {
            if (disposable != null) disposable.dispose();
        }
    
        public void onComplete() {
            disposable.dispose();
        }
    }
    

    【讨论】:

    • Martin,更改为 PublishSubject 没有帮助。发布的错误是完整的,即; Cannot resolve method 'subscribe(NewsSubscriber)'。消息表明 Subject 的 subscribe 方法不接受 DisposableSubscriber 作为输入。这就是我的问题。我们如何做到这一点?谢谢
    • 更新了我的答案。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-10-29
    • 1970-01-01
    • 2018-08-14
    • 2016-06-17
    相关资源
    最近更新 更多