【发布时间】:2020-04-04 13:33:33
【问题描述】:
我有一个可观察到的流式响应所有请求。我想在发出请求时创建该 observable 的过滤器,以便我可以与多个订阅者共享输出。以下是一些示例代码。
PublishSubject<String> publishSubject = PublishSubject.create();
Observable<String> fooObservable = publishSubject.filter(value -> value.startsWith("foo"))
.doOnSubscribe(disposable -> {
publishSubject.onNext("foobar");
})
.replay(1)
.refCount();
fooObservable.subscribe(val -> log.info("A val : <{}>", val));
我使用 PublishSubject 作为我的模拟服务,因为有时该服务会立即返回响应。
我的发现是,因为当有即时结果时没有当前订阅,所以我的 fooObservable 没有被填充。即我想查看时没有日志输出:
A val : <foobar>
请注意,我使用此代码得到相同的结果:
PublishSubject<String> publishSubject = PublishSubject.create();
Observable<String> fooObservable = publishSubject.filter(value -> value.startsWith("foo"))
.replay(1)
.refCount();
publishSubject.onNext("foobar");
fooObservable.subscribe(val -> log.info("A val : <{}>", val));
所以问题是 fooObservable 直到订阅后才订阅 PublishSubject,
有没有办法在第一次订阅 fooObservable 后立即运行代码?
编辑: 我想到了类似的东西:
PublishSubject<String> publishSubject = PublishSubject.create();
BehaviorSubject<String> fooObservable = BehaviorSubject.create();
publishSubject.filter(value -> value.startsWith("foo")).subscribe(fooObservable);
publishSubject.onNext("foobar");
fooObservable.subscribe(val -> log.info("A val : <{}>", val));
但是我有 2 个订阅,我不确定如何清理,因为过滤器不返回一次性订阅后的订阅。
编辑 2:后台任务的描述。
我的代码需要订阅第三方服务。该服务在我的代码中调用一个 onResponse 方法,其参数包含我的原始请求和响应。可以随时通过对 onResponse 的新调用来更新响应。
我想为这个提供方法的服务创建一个包装器:
public Observable<Response> getObservable(Request req);
如果请求匹配一个已经订阅的请求,那么 observable 应该在订阅时立即提供最新的匹配值。
当没有订阅者时,我需要取消订阅我正在打包的服务。
【问题讨论】:
-
好奇,你怎么想把值的发射绑定到订阅事件上呢?为什么不使用
BehaviorSubject之类的东西作为您的模拟服务?你可以直接调用BehaviorSubject.onNext(...),它会将该值发送给后续订阅者。 -
@Trogdor :你是在建议我用 BehaviourSubject 替换 fooObservable 吗?如果是这样,问题编辑是我对这种方法的问题。