RxJava:包含异步调用的observable

Jaa*_*tum 6 java asynchronous reactive-programming rx-java

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

请考虑以下方法,该方法返回一个可观察的调用NsdManager.registerService.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?
        }
    });
}
Run Code Online (Sandbox Code Playgroud)

只有在调用侦听器之后,观察者才能提供通知,但是异步调用侦听器.

我怎么能用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) {}
};
Run Code Online (Sandbox Code Playgroud)

tom*_*ozb 2

在监听器内部实现调用:

subscriber.onNext(result) 
subscriber.onComplete()
Run Code Online (Sandbox Code Playgroud)

这result是boolean传递给侦听器的。