I am new to RxJava and cannot figure out how to implement repeatable polling with a ConnectableObservable, with 2 subscribers processing the events on different threads.
I have a pipeline that roughly looks like this:
I would like to repeat the whole pipeline after a delay in a similar fashion to the solution from https://github.com/ReactiveX/RxJava/issues/448
Observable.fromCallable(() -> pollValue())
.repeatWhen(o -> o.concatMap(v -> Observable.timer(20, TimeUnit.SECONDS)));
or Dynamic delay value with repeatWhen().
This works fine with a plain (non-connectable) Observable, but not with multicasting.
Code example:
Works
final int[] i = {0};
Observable<Integer> integerObservable =
Observable.defer(() -> Observable.fromArray(i[0]++, i[0]++, i[0]++));
integerObservable
.observeOn(Schedulers.newThread())
.collect(StringBuilder::new, (sb, x) -> sb.append(x).append(","))
.map(StringBuilder::toString)
.toObservable().repeatWhen(o -> o.concatMap(v -> Observable.timer(1, TimeUnit.SECONDS)))
.subscribe(System.out::println);
Does not work:
final int[] i = {0};
ConnectableObservable<Integer> integerObservable =
Observable.defer(() -> Observable.fromArray(i[0]++, i[0]++, i[0]++))
.publish();
integerObservable.observeOn(Schedulers.newThread()).subscribe(System.out::println);
integerObservable
.observeOn(Schedulers.newThread())
.collect(StringBuilder::new, (sb, x) -> sb.append(x).append(","))
.map(StringBuilder::toString)
.toObservable().repeatWhen(o -> o.concatMap(v -> Observable.timer(1, TimeUnit.SECONDS)))
.subscribe(System.out::println);
integerObservable.connect();
