美文网首页
RxJava操作符--->组合/合并

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

作者: 谢尔顿 | 来源:发表于2018-06-26 17:43 被阅读11次

    引言

    该篇文章主要是关于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
    

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

    相关文章

      网友评论

          本文标题:RxJava操作符--->组合/合并

          本文链接:https://www.haomeiwen.com/subject/ldxmeftx.html