【发布时间】:2021-01-19 12:49:30
【问题描述】:
如何使用 OkHttp3 (v4.4.1) 实现长轮询,以获取响应的每一行的 RxJava (v2.2.11) Observable?可以在不阻塞线程保持读取行的情况下完成吗?如果我需要阻塞某个线程,那么我应该阻塞哪个线程?关于使用 OkHttp3 实现长轮询的任何一般示例? Google 在这个话题上对我很害羞...
TL;DR
我使用 OKHttp3 作为 HTTP 客户端,并将其包装在 makeGetObservable 方法调用中,该方法调用返回 Observable 响应,使用 newCall 回调向 Observable 发出事件。现在我正在尝试添加对长轮询服务的支持,我担心线程。
下面的代码演示了我正在尝试做的事情(并且似乎可以工作),但我很确定它是不行的。
// return Observable<Response>
makeGetObservable("http://my.service.com/api/events")
// check for error and map to Observable<ResponseBody>
.map(this::mapRespBodyOrError)
// flat map to Observable<String> representing line of long polling response
.flatMap(respBody -> Observable.create(emitter -> {
// open reader on response body stream
try (BufferedReader reader = new BufferedReader(respBody.charStream())) {
String line;
// block and wait to read a line from input
while((line = reader.readLine()) !=null) {
// once line was read from response body input stream emit it as observable event
emitter.onNext(line);
}
}
}));
【问题讨论】:
-
您在寻找
observeOn()吗? -
可能:) 我不知道更改调度程序是否是最合适的答案。如果我不必创建线程来阻止它,我会更喜欢
-
花了一些时间试验线程,发现被阻塞的线程是来自 okhttp3 客户端线程池的线程。默认值为 5,因此这是停止工作的限制。使用@Progman 建议的
observeOn和 IO 调度程序可以解决问题,因为它只会为每次阻塞读取生成新线程
标签: java rx-java rx-java2 okhttp long-polling