I have rx.Singles (or rx.Observables) which are to be executed in a dependency order as shown in this DAG (Directed-Acyclic-Graph) :
oneshould be executed firsttwoandthreeshould be executed only afteroneis completetwoandthreecan run in parallelfourshould be executed only aftertwois completefiveshould be executed only aftertwoandthreeare completefourandfivecan be run in parallelsixshould be executed at last, afterfourandfive.
How best can we achieve this ?
The following solution (Solution#1) is working as expected :
package rxtest;
import java.util.concurrent.Executor;
import java.util.concurrent.Executors;
import org.junit.Test;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import rx.Observable;
import rx.Single;
import rx.schedulers.Schedulers;
public class ObservablesDagTest {
private static final Logger logger = LoggerFactory.getLogger(ObservablesDagTest.class);
private static Executor customExecutor = Executors.newFixedThreadPool(20);
@Test
public void dagTest() {
Single<Integer> one = createSingle(1, 100);
Single<Integer> two = createSingle(2, 200);
Single<Integer> three = createSingle(3, 500);
Single<Integer> four = createSingle(4, 400);
Single<Integer> five = createSingle(5, 200);
Single<Integer> six = createSingle(6, 150);
executeDag(one, two, three, four, five, six);
}
private void executeDag(Single<Integer> one, Single<Integer> two, Single<Integer> three, Single<Integer> four,
Single<Integer> five, Single<Integer> six) {
logger.info("BEGIN");
Single<Integer> twoCached = two.toObservable().cache().toSingle();
Observable.concat(
one.toObservable(),
Observable.merge(
Single.concat(twoCached, four),
Observable.concat(Single.merge(twoCached, three), five.toObservable())),
six.toObservable())
.toBlocking()
.subscribe(i -> logger.info("Received : " + i));
logger.info("END");
}
private Single<Integer> createSingle(int j, int sleepMs) {
Single<Integer> single = Single.just(j)
.flatMap(i -> Single.<Integer> create(s -> {
logger.info("onSubscribe : {}", i);
sleep(sleepMs);
s.onSuccess(i);
}).subscribeOn(Schedulers.from(customExecutor)));
return single;
}
private void sleep(int ms) {
try {
Thread.sleep(ms);
}
catch (InterruptedException e) {
}
}
}
Output
21:16:23.879 [main] INFO rxtest.ObservablesDagTest BEGIN
21:16:23.974 [pool-1-thread-1] INFO rxtest.ObservablesDagTest onSubscribe : 1
21:16:24.079 [main] INFO rxtest.ObservablesDagTest Received : 1
21:16:24.101 [pool-1-thread-3] INFO rxtest.ObservablesDagTest onSubscribe : 3
21:16:24.101 [pool-1-thread-2] INFO rxtest.ObservablesDagTest onSubscribe : 2
21:16:24.303 [main] INFO rxtest.ObservablesDagTest Received : 2
21:16:24.303 [main] INFO rxtest.ObservablesDagTest Received : 2
21:16:24.304 [pool-1-thread-4] INFO rxtest.ObservablesDagTest onSubscribe : 4
21:16:24.602 [main] INFO rxtest.ObservablesDagTest Received : 3
21:16:24.603 [pool-1-thread-5] INFO rxtest.ObservablesDagTest onSubscribe : 5
21:16:24.704 [main] INFO rxtest.ObservablesDagTest Received : 4
21:16:24.804 [main] INFO rxtest.ObservablesDagTest Received : 5
21:16:24.805 [pool-1-thread-6] INFO rxtest.ObservablesDagTest onSubscribe : 6
21:16:24.956 [main] INFO rxtest.ObservablesDagTest Received : 6
21:16:24.956 [main] INFO rxtest.ObservablesDagTest END
Is there any more efficient/fluent solution than this ?
