-
Notifications
You must be signed in to change notification settings - Fork 119
/
Copy pathSubscribeOnAndObserveOn.java
64 lines (49 loc) · 1.35 KB
/
SubscribeOnAndObserveOn.java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
package com.packtpub.reactive.chapter06;
import static com.packtpub.reactive.common.Helpers.debug;
import java.util.concurrent.CountDownLatch;
import rx.Observable;
import rx.schedulers.Schedulers;
import com.packtpub.reactive.common.Program;
/**
* Demonstrates using subscribeOn and observeOn with {@link Schedulers} and {@link Observable}s.
*
* @author meddle
*/
public class SubscribeOnAndObserveOn implements Program {
@Override
public String name() {
return "A few examples of observeOn and subscribeOn";
}
@Override
public int chapter() {
return 6;
}
@Override
public void run() {
CountDownLatch latch = new CountDownLatch(1);
Observable<Integer> range = Observable
.range(20, 5)
.flatMap(n -> Observable
.range(n, 3)
.subscribeOn(Schedulers.computation())
.doOnEach(debug("Source"))
);
Observable<Character> chars = range
.observeOn(Schedulers.newThread())
.map(n -> n + 48)
.doOnEach(debug("+48 ", " "))
.observeOn(Schedulers.computation())
.map(n -> Character.toChars(n))
.map(c -> c[0])
.doOnEach(debug("Chars ", " "))
.finallyDo(() -> latch.countDown());
chars.subscribe();
System.out.println("Hey!");
try {
latch.await();
} catch (InterruptedException e) {}
}
public static void main(String[] args) {
new SubscribeOnAndObserveOn().run();
}
}