【发布时间】:2018-11-09 15:57:28
【问题描述】:
以下是我当前的代码
private final List<Disposable> subscriptions = new ArrayList<>();
for (Instrument instrument : instruments) {
// Waiting for OrderBook to generate Reliable results.
GenericBook Book =
service
.getBook(instrument.getData())
.filter(gob -> onBookUpdate(gob))
.blockingFirst();
subscriptions.add(
service
.getBook(instrument.getData())
.subscribe(
gob -> {
try {
onBookUpdate(gob);
} catch (Exception e) {
logger.error("Error on subscription:", e);
}
},
e -> logger.error("Error on subscription:", e)));
}
所以它的作用是对它首先阻塞的每个仪器等待,直到onBookUpdate(gob) 的输出变为真。 onBookUpdate(gob) 返回布尔值。
一旦我们将第一个 onBookUpdate 设置为 true,我就会将该订阅者推送到订阅变量中。
这会变慢,因为我必须等待每个乐器,然后再继续下一个乐器。
我的目标是并行运行所有这些,然后等待所有完成并将它们推送到订阅变量。
我试过 zip 但没用
List<Observable<GenericOrderBook>> obsList = null;
for (Instrument instrument : instruments) {
// This throws nullException.
obsList.add(service
.getBook(instrument.getData())
.filter(gob -> onBookUpdate(gob))
.take(1));
}
}
// Some how wait over here until all get first onBookUpdate as true.
String o = Observable.zip(obsList, (i) -> i[0]).blockingLast();
【问题讨论】:
标签: java rx-java observable rx-java2