【问题标题】:RxJava: observable that contains an asynchronous callRxJava:包含异步调用的可观察对象
【发布时间】:2014-08-12 08:46:32
【问题描述】:

我正在尝试理解 RxJava 并遇到以下情况。

考虑下面的方法,它返回一个调用 NsdManager.registerService 的 observable。 registerService 方法需要一个监听器,在注册成功(或失败)时调用。

public Observable<Boolean> registerService() {
    return Observable.create(new Observable.OnSubscribe<Boolean>() {

        @Override
        public void call(Subscriber<? super Boolean> subscriber) {
            nsdManager.registerService(serviceInfo, NsdManager.PROTOCOL_DNS_SD, registrationListener);

            // how to proceed?
        }
    });
}

观察者只有在监听器被调用后才能提供通知,但监听器是异步调用的。

如何使用 RxJava 做到这一点?


我想出了以下内容,使用 BehaviorSubject。不知道这是否是最好的解决方案,但它确实有效。

private BehaviorSubject<Boolean> registrationSubject;

public Observable<Boolean> registerService() {
    registrationSubject = BehaviorSubject.create();

    Observable.create(new Observable.OnSubscribe<Boolean>() {
        @Override
        public void call(Subscriber<? super Boolean> subscriber) {
            NsdServiceInfo serviceInfo  = new NsdServiceInfo();
            serviceInfo.setServiceName(serviceName);
            serviceInfo.setServiceType(NSD_SERVICE_TYPE);
            serviceInfo.setPort(serverSocket.getLocalPort());

            nsdManager.registerService(serviceInfo, NsdManager.PROTOCOL_DNS_SD, registrationListener);
        }
    }).subscribe(registrationSubject);

    return registrationSubject;
}

private NsdManager.RegistrationListener registrationListener = new NsdManager.RegistrationListener() {
    @Override
    public void onRegistrationFailed(NsdServiceInfo serviceInfo, int errorCode) {
        registrationSubject.onNext(false);
        registrationSubject.onCompleted();
    }

    @Override
    public void onServiceRegistered(NsdServiceInfo serviceInfo) {
        registrationSubject.onNext(true);
        registrationSubject.onCompleted();
    }

    @Override
    public void onUnregistrationFailed(NsdServiceInfo serviceInfo, int errorCode) { }

    @Override
    public void onServiceUnregistered(NsdServiceInfo serviceInfo) {}
};

【问题讨论】:

    标签: java asynchronous reactive-programming rx-java


    【解决方案1】:

    我认为尽可能避免使用主题会更好。 在您的解决方案中,您只使用主题来调用onNextonCompleted。但是,在Observable.create() 方法中,您已经可以访问subscriber,您可以在其上调用这些方法。换句话说,您可以将事件处理程序的完整设置封装在 Observable.create() 方法中。

    public Observable<Boolean> registerService() {
        return Observable.create(new Observable.OnSubscribe<Boolean>() {
            @Override
            public void call(final Subscriber<? super Boolean> subscriber) {
                NsdServiceInfo serviceInfo  = new NsdServiceInfo();
                serviceInfo.setServiceName(serviceName);
                serviceInfo.setServiceType(NSD_SERVICE_TYPE);
                serviceInfo.setPort(serverSocket.getLocalPort());
    
                nsdManager.registerService(serviceInfo, NsdManager.PROTOCOL_DNS_SD, 
                    new NsdManager.RegistrationListener() {
                        @Override
                        public void onRegistrationFailed(NsdServiceInfo serviceInfo, int errorCode) {
                            if (!subscriber.isUnsubscribed()) {
                                subscriber.onNext(false);
                                subscriber.onCompleted();
                            }
                        }
    
                        @Override
                        public void onServiceRegistered(NsdServiceInfo serviceInfo) {
                            if (!subscriber.isUnsubscribed()) {
                                subscriber.onNext(true);
                                subscriber.onCompleted();
                            }
                        }
    
                        @Override
                        public void onUnregistrationFailed(NsdServiceInfo serviceInfo, int errorCode) { }
    
                        @Override
                        public void onServiceUnregistered(NsdServiceInfo serviceInfo) {}
                    }
                );
            }
        });
    }
    

    【讨论】:

      【解决方案2】:

      在监听器实现调用内部:

      subscriber.onNext(result) 
      subscriber.onComplete()
      

      result 是传递给侦听器的 boolean

      【讨论】:

      • 在调用外部服务if (!subscriber.isUnsubscribed()) { subscriber.onNext(nsdManager.register...); subscriber.onComplete(); }之前,您可能还想检查传递的订阅者是否仍然订阅
      猜你喜欢
      • 1970-01-01
      • 2015-01-12
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多