RxJava操作符实践:10_连接操作之2_publish

一、描述

将普通的Observable转换为可连接的Observable。

可连接的Observable (connectable Observable)与普通的Observable差不多,不过它并不会在被订阅时开始发射数据,而是直到使用了Connect操作符时才会开始。用这种方法,你可以在任何时候让一个Observable开始发射数据。

有一个变体接受一个函数作为参数。这个函数用原始Observable发射的数据作为参数,产生一个新的数据作为ConnectableObservable给发射,替换原位置的数据项。实质是在签名的基础上添加一个Map操作。

二、示意图

publish

三、示例代码

1. 未使用publish操作符时

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
Observable observable = Observable.range(1, 1000000)
.sample(10, TimeUnit.MILLISECONDS);

observable.subscribe(new Subscriber<Integer>() {
@Override
public void onCompleted() {
System.out.println("onCompleted1.");
}

@Override
public void onError(Throwable e) {
System.out.println("onError1: " + e.getMessage());
}

@Override
public void onNext(Integer integer) {
System.out.println("onNext1: " + integer);
}
});

observable.subscribe(new Subscriber<Integer>() {
@Override
public void onCompleted() {
System.out.println("onCompleted2.");
}

@Override
public void onError(Throwable e) {
System.out.println("onError2: " + e.getMessage());
}

@Override
public void onNext(Integer integer) {
System.out.println("onNext2: " + integer);
}
});

2. 运行结果

1
2
3
4
5
6
7
8
9
10
onNext1: 15281
onNext1: 40401
...
onNext1: 983620
onCompleted1.
onNext2: 25895
onNext2: 48456
...
onNext2: 983356
onCompleted2.

可见Observable在订阅的时候就开始发射数据,导致两个观察者收到的数据是不一样的。

3. 使用publish操作符后

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
ConnectableObservable observable = Observable.range(1, 1000000).sample(10, TimeUnit.MILLISECONDS).publish();

observable.subscribe(new Subscriber<Integer>() {
@Override
public void onCompleted() {
System.out.println("onCompleted1.");
}

@Override
public void onError(Throwable e) {
System.out.println("onError1: " + e.getMessage());
}

@Override
public void onNext(Integer integer) {
System.out.println("onNext1: " + integer);
}
});

observable.subscribe(new Subscriber<Integer>() {
@Override
public void onCompleted() {
System.out.println("onCompleted2.");
}

@Override
public void onError(Throwable e) {
System.out.println("onError2: " + e.getMessage());
}

@Override
public void onNext(Integer integer) {
System.out.println("onNext2: " + integer);
}
});

observable.connect();

四、运行结果

1
2
3
4
5
6
7
8
9
onNext1: 20491
onNext2: 20491
onNext1: 39191
onNext2: 39191
...
onNext1: 997372
onNext2: 997372
onCompleted1.
onCompleted2.

可见在订阅的时候,并不会开始发射数据,只有等到connect连接后,才开始发射数据,所以两个观察者接收到的数据是一样的。

五、参考资料

ReactiveX官方文档

ReactiveX文档中文翻译

PS:欢迎关注 SherlockShi 个人博客

感谢你的支持,让我继续努力分享有用的技术和知识点!