【发布时间】:2019-10-18 14:14:42
【问题描述】:
我有以下测试代码:
FlowableOnSubscribe<SomeObj> fos;
private void init() {
fos = emitter -> {
try {
while(true) {
SomeObj someObj = readFromDataInputStream();
emitter.onNext(someObj);
System.out.println("Emitted object");
}
}
catch(Exception e) {
emitter.onError(e);
}
};
}
public Single<String> doWork() {
Flowable<String> myFlow = Flowable.defer(() ->
Flowable.create(fos, BackpressureStrategy.BUFFER)
.cache()
.subscribeOn(Schedulers.io(), false)
.doOnSubscribe(x -> doSomethingToTriggerDataInputStream())
.map(x -> convertMyCustomObjectToString(x))
);
);
// lastOrError/singleOrError end up blocking :(
return myFlow.firstOrError();
}
public VertxDataFetcher<CompletionStage<String>> vertxDataFetcherTest() {
return new VertxDataFetcher<>((env, future) -> {
try {
future.complete(doWork().to(SingleInterop.get()));
}
catch(Exception e) {
future.fail();
}
});
}
public DataFetcher<CompletionStage<String>> dataFetcherTest() {
return env -> doWork().to(SingleInterop.get());
}
如果我运行示例代码,它会在第一次成功使用后挂起。换句话说,它会在初始网页加载时运行一次,但如果我在浏览器中进行刷新(Ctrl F5),它会挂起并且不再完成调用。
FWIW,经过一点调试,它看起来像 在 RxCachedThreadScheduler1 上发生 defer/map 调用,在网页刷新后它挂在与 RxCachedThreadScheduler2 的延迟调用中。
如果我切换到为每个呼叫使用单独的输入流,它不会挂断。我遵循的示例在这里(RxJava: Feed one stream (Observable) as the input of another stream...)。
但是,这不适用于我的设计,因为它需要 1 个始终保持打开状态的共享连接。这是因为我希望 GraphQL 订阅能够捕获数据输入流上发出的任何内容。如果我有多个/单独的套接字连接,GraphQL 订阅输入流将错过在非订阅输入流上发出的任何内容。 (除非最好的选择是让所有其他输入流也发送到 GraphQL 订阅流,我不知道该怎么做......)
作为旁注,在这种情况下我使用 VertxDataFetcher 还是 DataFetcher 是否重要?我目前正在使用 DataFetcher 来运行我的示例。如果我需要切换到 VertxDataFetcher,我不确定如何正确转换类型以使该方法正常工作。
【问题讨论】: