引言
该篇文章主要是关于RxJava的组合/变换操作符使用的代码讲解。组合/变换操作符总共有四大类:
(1)组合多个被观察者
- 按发送顺序:concat()、concatArray()
- 按时间:merge()、mergeArray()
- 错误处理:concatDelayError()、mergeDelayError()
(2)合并多个事件
- 按数量:zip()
- 按时间:combineLatest()、combineLatestDelayError()
- 合并成一个事件发送:reduce()、collect()
(3)发送事件前追加发送事件
- startWith()
- startWithArray()
(4)统计发送事件数量
- count()
1. concat()/concatArray()
组合多个被观察者一起发送数据,合并后按发送顺序串行执行。
二者区别:组合被观察者的数量,即concat()组合被观察者数量≤4个,而concatArray()则可>4个。
Observable.concat(
Observable.just(1,2,3),
Observable.just(4,5,6),
Observable.just(7,8,9),
Observable.just(10,11,12))
.subscribe(new Observer<Integer>() {
@Override
public void onSubscribe(Disposable d) {
}
@Override
public void onNext(Integer value) {
Log.d(Constant.TAG,"接收到了事件"+value);
}
@Override
public void onError(Throwable e) {
Log.d(Constant.TAG,"对Error事件做出响应");
}
@Override
public void onComplete() {
Log.d(Constant.TAG,"对Complete事件做出响应");
}
});
Observable.concatArray(Observable.just(1,2),
Observable.just(3,4),
Observable.just(5,6),
Observable.just(7,8),
Observable.just(9,10))
.subscribe(new Observer<Integer>() {
@Override
public void onSubscribe(Disposable d) {
}
@Override
public void onNext(Integer value) {
Log.d(Constant.TAG,"接收到了事件"+value);
}
@Override
public void onError(Throwable e) {
Log.d(Constant.TAG,"对Error事件做出响应");
}
@Override
public void onComplete() {
Log.d(Constant.TAG,"对Complete事件做出响应");
}
});
concat()的log信息:
06-22 14:00:51.141 12967-12967/? D/RxJava: 接收到了事件1
06-22 14:00:51.141 12967-12967/? D/RxJava: 接收到了事件2
06-22 14:00:51.141 12967-12967/? D/RxJava: 接收到了事件3
06-22 14:00:51.142 12967-12967/? D/RxJava: 接收到了事件4
06-22 14:00:51.141 12967-12967/? D/RxJava: 接收到了事件5
06-22 14:00:51.141 12967-12967/? D/RxJava: 接收到了事件6
06-22 14:00:51.141 12967-12967/? D/RxJava: 接收到了事件7
06-22 14:00:51.141 12967-12967/? D/RxJava: 接收到了事件8
06-22 14:00:51.143 12967-12967/? D/RxJava: 对Complete事件做出响应
concatArray()的log信息:
06-22 14:00:51.143 12967-12967/? D/RxJava: 对Complete事件做出响应
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件1
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件2
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件3
06-22 14:00:51.144 12967-12967/? D/RxJava: 接收到了事件4
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件5
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件6
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件7
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件8
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件9
06-22 14:00:51.143 12967-12967/? D/RxJava: 接收到了事件10
06-22 14:00:51.143 12967-12967/? D/RxJava: 对Complete事件做出响应
2. merge()/mergeArray()
组合多个被观察者一起发送数据,合并后按时间线并行执行。
1.二者区别:和上述的concat和concatArray的一样;
2.区别上述concat操作符,同样是组合多个被观察者一起发送数据,但concat操作符合并后是按发送顺序串行执行。
Observable.merge(
Observable.intervalRange(0,3,1,1, TimeUnit.SECONDS),
Observable.intervalRange(2,3,1,1, TimeUnit.SECONDS))
.subscribe(new Observer<Long>() {
@Override
public void onSubscribe(Disposable d) {
}
@Override
public void onNext(Long value) {
Log.d(Constant.TAG,"接收到了事件"+value);
}
@Override
public void onError(Throwable e) {
Log.d(Constant.TAG,"对Error事件做出响应");
}
@Override
public void onComplete() {
Log.d(Constant.TAG,"对Complete事件做出响应");
}
});
log信息:
06-22 14:23:11.358 14031-14082/com.gjj.frame D/RxJava: 接收到了事件0
06-22 14:23:11.366 14031-14083/com.gjj.frame D/RxJava: 接收到了事件2
06-22 14:23:12.357 14031-14082/com.gjj.frame D/RxJava: 接收到了事件1
06-22 14:23:12.358 14031-14082/com.gjj.frame D/RxJava: 接收到了事件3
06-22 14:23:13.358 14031-14082/com.gjj.frame D/RxJava: 接收到了事件2
06-22 14:23:13.359 14031-14082/com.gjj.frame D/RxJava: 接收到了事件4
06-22 14:23:13.362 14031-14082/com.gjj.frame D/RxJava: 对Complete事件做出响应
3. concatArrayDelayError()/mergeArrayDelayError()
使用concat和merge操作符时,若其中一个被观察者发出onError事件,则会马上终止其他被观察者继续发送事件,若希望onError事件推迟到其他被观察者发送事件结束后才处罚,就需要使用对应的concatDelayError或mergeDelayError()操作符。
(1)无使用concatArrayDelayError()的情况
Observable.concat(Observable.create(new ObservableOnSubscribe<Integer>() {
@Override
public void subscribe(ObservableEmitter<Integer> e) throws Exception {
e.onNext(1);
e.onNext(2);
e.onNext(3);
//对error事件,因为无使用concatDelayError,所以第二个Observable将不会发送事件
e.onError(new NullPointerException());
e.onComplete();
}
}),Observable.just(4,5,6)).subscribe(new Observer<Integer>() {
@Override
public void onSubscribe(Disposable d) {
}
@Override
public void onNext(Integer value) {
Log.d(Constant.TAG,"接收到了事件"+value);
}
@Override
public void onError(Throwable e) {
Log.d(Constant.TAG,"对error事件做出响应");
}
@Override
public void onComplete() {
Log.d(Constant.TAG,"对Complete事件做出响应");
}
});
测试结果:第一个悲观者发送Error事件后,第2个被观察者则不会继续发送事件。
log信息:
06-25 11:03:06.905 21337-21337/com.gjj.frame D/RxJava: 接收到了事件1
06-25 11:03:06.905 21337-21337/com.gjj.frame D/RxJava: 接收到了事件2
06-25 11:03:06.905 21337-21337/com.gjj.frame D/RxJava: 接收到了事件3
06-25 11:03:06.906 21337-21337/com.gjj.frame D/RxJava: 对error事件做出响应
(2)使用concatArrayDelayError()的情况
Observable.concatArrayDelayError(Observable.create(new ObservableOnSubscribe<Integer>() {
@Override
public void subscribe(ObservableEmitter<Integer> e) throws Exception {
e.onNext(1);
e.onNext(2);
e.onNext(3);
//对error事件,因为无使用concatDelayError,所以第二个Observable将不会发送事件
e.onError(new NullPointerException());
e.onComplete();
}
}),Observable.just(4,5,6)).subscribe(new Observer<Integer>() {
@Override
public void onSubscribe(Disposable d) {
}
@Override
public void onNext(Integer value) {
Log.d(Constant.TAG,"接收到了事件"+value);
}
@Override
public void onError(Throwable e) {
Log.d(Constant.TAG,"对error事件做出响应");
}
@Override
public void onComplete() {
Log.d(Constant.TAG,"对Complete事件做出响应");
}
});
测试结果:第1个被观察者的error事件将在第2个被观察者发送完事件后再继续发送。
log信息:
06-25 11:08:55.097 21509-21509/com.gjj.frame D/RxJava: 接收到了事件1
06-25 11:08:55.097 21509-21509/com.gjj.frame D/RxJava: 接收到了事件2
06-25 11:08:55.097 21509-21509/com.gjj.frame D/RxJava: 接收到了事件3
06-25 11:08:55.097 21509-21509/com.gjj.frame D/RxJava: 接收到了事件4
06-25 11:08:55.097 21509-21509/com.gjj.frame D/RxJava: 接收到了事件5
06-25 11:08:55.097 21509-21509/com.gjj.frame D/RxJava: 接收到了事件6
06-25 11:08:55.097 21509-21509/com.gjj.frame D/RxJava: 对error事件做出响应
4. Zip()
合并多个被观察者(Observable)发送的事件,生成一个新的事件序列(即组合过后的事件序列),并最终发送。
//创建第1个观察者
Observable<Integer> observable1 = Observable.create(new ObservableOnSubscribe<Integer>() {
@Override
public void subscribe(ObservableEmitter<Integer> e) throws Exception {
e.onNext(1);
e.onNext(2);
e.onNext(3);
}
}).subscribeOn(Schedulers.io());//设置被观察者1再工作线程1中工作
//创建第2个观察者
Observable<String> observable2 = Observable.create(new ObservableOnSubscribe<String>() {
@Override
public void subscribe(ObservableEmitter<String> e) throws Exception {
e.onNext("A");
e.onNext("B");
e.onNext("C");
e.onNext("D");
e.onComplete();
}
}).subscribeOn(Schedulers.newThread());//设置被观察者2再工作线程2中工作
Observable.zip(observable1, observable2, new BiFunction<Integer, String, String>() {
@Override
public String apply(Integer integer, String s) throws Exception {
return integer+s;
}
}).subscribe(new Observer<String>() {
@Override
public void onSubscribe(Disposable d) {
}
@Override
public void onNext(String value) {
Log.d(Constant.TAG,"最终收到的事件 = "+value);
}
@Override
public void onError(Throwable e) {
Log.d(Constant.TAG,"onError");
}
@Override
public void onComplete() {
Log.d(Constant.TAG,"onComplete");
}
});
log信息:
06-26 16:30:02.147 29926-29985/com.gjj.frame D/RxJava: 最终收到的事件 = 1A
06-26 16:30:02.150 29926-29984/com.gjj.frame D/RxJava: 最终收到的事件 = 2B
06-26 16:30:02.151 29926-29984/com.gjj.frame D/RxJava: 最终收到的事件 = 3C
注意:最终合并的事件数量是多个被观察者中最少的数量,多余的事件将不会发送。
5. combineLatest()
当两个Observable中的任何一个发送了数据后,将先发送了数据的Observables的最新(最后)一个数据与另外一个Observable发送的每一个数据结合,最终基于该函数的结果发送数据。
Observable.combineLatest(Observable.just(1L, 2L, 3L), Observable.intervalRange(0, 3, 1, 1, TimeUnit.SECONDS), new BiFunction<Long, Long, Long>() {
@Override
public Long apply(Long aLong, Long aLong2) throws Exception {
Log.d(Constant.TAG,"合并的数据是:"+aLong+" "+aLong2);
return aLong+aLong2;
}
}).subscribe(new Consumer<Long>() {
@Override
public void accept(Long aLong) throws Exception {
Log.d(Constant.TAG,"合并的结果是:"+aLong);
}
});
log信息:
06-26 16:48:37.010 30604-30634/com.gjj.frame D/RxJava: 合并的数据是:3 0
06-26 16:48:37.011 30604-30634/com.gjj.frame D/RxJava: 合并的结果是:3
06-26 16:48:38.010 30604-30634/com.gjj.frame D/RxJava: 合并的数据是:3 1
06-26 16:48:38.011 30604-30634/com.gjj.frame D/RxJava: 合并的结果是:4
06-26 16:48:39.012 30604-30634/com.gjj.frame D/RxJava: 合并的数据是:3 2
06-26 16:48:39.013 30604-30634/com.gjj.frame D/RxJava: 合并的结果是:5
6. combineLatestDelayError()
作用类似于concatArrayDelayError()。
7. reduce()
把被观察者需要发送的事件聚合成一个事件&发送
Observable.just(1,2,3,4)
.reduce(new BiFunction<Integer, Integer, Integer>() {
@Override
public Integer apply(Integer integer, Integer integer2) throws Exception {
Log.d(Constant.TAG,"本次计算的数据是:"+integer+"乘"+integer2);
return integer * integer2;
}
}).subscribe(new Consumer<Integer>() {
@Override
public void accept(Integer integer) throws Exception {
Log.d(Constant.TAG,"最终计算的结果是:"+integer);
}
});
log信息:
06-26 16:59:56.401 31613-31613/com.gjj.frame D/RxJava: 本次计算的数据是:1乘2
06-26 16:59:56.402 31613-31613/com.gjj.frame D/RxJava: 本次计算的数据是:2乘3
06-26 16:59:56.402 31613-31613/com.gjj.frame D/RxJava: 本次计算的数据是:6乘4
06-26 16:59:56.402 31613-31613/com.gjj.frame D/RxJava: 最终计算的结果是:24
8. collect()
将被观察者Observable发送的数据事件收集到一个数据结构里
Observable.just(1,2,3,4,5,6)
.collect(new Callable<ArrayList<Integer>>() {
@Override
public ArrayList<Integer> call() throws Exception {
return new ArrayList<>();
}
}, new BiConsumer<ArrayList<Integer>, Integer>() {
@Override
public void accept(ArrayList<Integer> list, Integer integer) throws Exception {
list.add(integer);
}
}).subscribe(new Consumer<ArrayList<Integer>>() {
@Override
public void accept(ArrayList<Integer> integers) throws Exception {
Log.d(Constant.TAG,"本次发送的数据是:"+integers);
}
});
log信息:
06-26 17:04:40.264 31785-31785/com.gjj.frame D/RxJava: 本次发送的数据是:[1, 2, 3, 4, 5, 6]
9. startWith()/startWithArray()
在一个被观察者发送事件钱,追加发送一些数据/一个新的被观察者
Observable.just(3,4)
.startWith(0)//追加单个数据
.startWithArray(1,2)//追加多个数据
.subscribe(new Consumer<Integer>() {
@Override
public void accept(Integer integer) throws Exception {
Log.d(Constant.TAG,"接收到了事件"+integer);
}
});
log信息:
06-26 17:39:24.075 4052-4052/com.gjj.frame D/RxJava: 接收到了事件1
06-26 17:39:24.075 4052-4052/com.gjj.frame D/RxJava: 接收到了事件2
06-26 17:39:24.075 4052-4052/com.gjj.frame D/RxJava: 接收到了事件0
06-26 17:39:24.075 4052-4052/com.gjj.frame D/RxJava: 接收到了事件3
06-26 17:39:24.075 4052-4052/com.gjj.frame D/RxJava: 接收到了事件4
10.count()
统计被观察者发送事件的数量。
Observable.just(1,2,3,4)
.count()
.subscribe(new Consumer<Long>() {
@Override
public void accept(Long integer) throws Exception {
Log.d(Constant.TAG,"发送的事件数量 = "+integer);
}
});
log信息:
06-26 17:42:20.639 4750-4750/com.gjj.frame D/RxJava: 发送的事件数量 = 4