【发布时间】: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