RxJava操作符--->组合/合并

2018-06-26  本文已影响11人  谢尔顿

引言

该篇文章主要是关于RxJava的组合/变换操作符使用的代码讲解。组合/变换操作符总共有四大类:

(1)组合多个被观察者

(2)合并多个事件

(3)发送事件前追加发送事件

(4)统计发送事件数量

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

参考文章:
Android RxJava:组合 / 合并操作符 详细教程

上一篇下一篇

猜你喜欢

热点阅读