INNER CODE UNIT · Java
stockPriceUpdateObservable
righettod/poc-graphql · src/main/java/eu/righettod/graphqlpoc/publishers/NewAssociationPublisher.java:42
Observable<String> stockPriceUpdateObservable = Observable.create(emitter -> {
ScheduledExecutorService executorService = Executors.newScheduledThreadPool(1);
executorService.scheduleAtFixedRate(verifyPresenceOfNewAssociation(emitter), 0, 20, TimeUnit.SECONDS);
});
ConnectableObservable<String> connectableObservable = stockPriceUpdateObservable.share().publish();
connectableObservable.connect();
publisher = connectableObservable.toFlowable(BackpressureStrategy.BUFFER);
}
/**
* Verify if new association has been created
*
* @param emitter Event emitter
* @return A runnable instance
*/
private Runnable verifyPresenceOfNewAssociation(ObservableEmitter<String> emitter) {
return () -> {
try {